【发布时间】:2021-06-12 16:32:21
【问题描述】:
我正在尝试使用 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 pipeline orchestration airflow-2.x