根据气流中 sql 查询的结果创建动态任务

问题描述

我正在尝试使用 TaskGroup 创建动态任务,并将结果保存在变量中。根据数据库查询,该变量每 N 分钟修改一次,但是当第二次修改该变量时,调度程序会崩溃

基本上我需要根据查询中收到的唯一行数创建任务。

以 TaskGroup(f"task") 作为任务:

    data_variable = Variable.get("df")
    data = data_variable

    try :
        if data != False and data !='none':
            df = pd.read_json(data)

            for field_id in df.field.unique():
             

                task1 = Pythonoperator(
                   
                )
                task2 = Pythonoperator(
                   
                )

               
                task1 >> task2

    except:
        pass

有没有办法用任务组来做到这一点?

解决方法

这是 Airflow 的反模式。

虽然您可以在顶级代码中使用 Variable.get("df"),但您不应该这样做。变量/连接/使用任何数据库创建查询的任何其他代码应仅在操作符范围内或使用 Jinja 模板完成。这样做的原因是 Airflow 会定期解析 DAG 文件(如果您没有更改 min_file_process_interval 的默认值,则每 30 秒一次),因此每 30 秒与数据库交互一次的代码将导致该数据库的负载过重。 对于其中一些情况,未来的气流版本会发出警告(请参阅 PR

气流任务应尽可能保持静态(或缓慢变化)。

相关问答

Selenium Web驱动程序和Java。元素在(x,y)点处不可单击。其...
Python-如何使用点“。” 访问字典成员?
Java 字符串是不可变的。到底是什么意思?
Java中的“ final”关键字如何工作?(我仍然可以修改对象。...
“loop:”在Java代码中。这是什么,为什么要编译?
java.lang.ClassNotFoundException:sun.jdbc.odbc.JdbcOdbc...