Airflow构建DAG时>>运算符、set_downstream及chain方法报错如何解决
Airflow DAG结构搭建问题排查
预期要实现的DAG结构如下:
报错复现过程
第一次编写:位运算符定义依赖报错
代码:
with DAG( dag_id="dag_example", schedule_interval="@once", default_args={ "owner": "airflow", "start_date": datetime(2020, 11, 1), "retries": 1, "retry_delay": timedelta(seconds=10), }, catchup=False, ) as dag: tA = PythonOperator( task_id="A", python_callable=A, provide_context=True, ), tB = PythonOperator( task_id="B", python_callable=B, provide_context=True, ), taS = PythonOperator( task_id="aS", python_callable=aS, provide_context=True, ), taF = PythonOperator( task_id="aF", python_callable=aF, provide_context=True, ), tbS = PythonOperator( task_id="bS", python_callable=A, provide_context=True, ), tbF = PythonOperator( task_id="bF", python_callable=A, provide_context=True, ) tA >> tB tA >> taS tA >> taF tB >> tbF tB >> tbS
运行报错:TypeError: unsupported operand type(s) for >>: 'tuple' and 'tuple'
改用set_downstream方法调整依赖定义:
tA.set_downstream(tB) tA.set_downstream(taS) tA.set_downstream(taF) tB.set_downstream(tbS) tB.set_downstream(tbF)
新报错:AttributeError: 'tuple' object has no attribute 'set_downstream'
第二次编写:chain方法嵌套列表报错
代码:
with DAG( dag_id="dag_example", schedule_interval="@once", default_args={ "owner": "airflow", "start_date": datetime(2020, 11, 1), "retries": 1, "retry_delay": timedelta(seconds=10), # "on_failure_callback": snow_fail, # "on_success_callback": snow_succ, }, catchup=False, ) as dag: chain( PythonOperator( task_id="A", python_callable=A, provide_context=True, ), [ PythonOperator( task_id="aF", python_callable=aF, provide_context=True, trigger_rule=TriggerRule.ONE_FAILED, ), PythonOperator( task_id="B", python_callable=B, provide_context=True, ), [ PythonOperator( task_id="bF", python_callable=bF, provide_context=True, trigger_rule=TriggerRule.ONE_FAILED, ), PythonOperator( task_id="bS", python_callable=bS, provide_context=True, trigger_rule=TriggerRule.ONE_SUCCESS, ), ], PythonOperator( task_id="aS", python_callable=aS, provide_context=True, trigger_rule=TriggerRule.ONE_SUCCESS, ), ], )
报错:AttributeError: 'list' object has no attribute 'update_relative'
第三次编写:调整chain参数后结构不符预期
调整后代码:
with DAG( dag_id="dag_example", schedule_interval="@once", default_args={ "owner": "airflow", "start_date": datetime(2020, 11, 1), "retries": 1, "retry_delay": timedelta(seconds=10), # "on_failure_callback": snow_fail, # "on_success_callback": snow_succ, }, catchup=False, ) as dag: chain( PythonOperator( task_id="A", python_callable=A, provide_context=True, ), [ PythonOperator( task_id="aF", python_callable=aF, provide_context=True, trigger_rule=TriggerRule.ONE_FAILED, ), PythonOperator( task_id="aS", python_callable=aS, provide_context=True, trigger_rule=TriggerRule.ONE_SUCCESS, ), ], PythonOperator( task_id="B", python_callable=B, provide_context=True, ), [ PythonOperator( task_id="bF", python_callable=bF, provide_context=True, trigger_rule=TriggerRule.ONE_FAILED, ), PythonOperator( task_id="bS", python_callable=bS, provide_context=True, trigger_rule=TriggerRule.ONE_SUCCESS, ), ], )
生成的DAG结构和预期不符,实际结构如下:
最终解决方案
核心错误原因
第一次编写时,每个PythonOperator变量定义末尾多写了逗号,导致tA、tB等变量被识别为单元素元组,而非Operator对象,因此无法调用set_downstream方法,也无法使用>>运算符定义依赖。
另外chain方法不支持嵌套列表传参,仅支持同层级的任务列表,因此嵌套写法会报错。
修正后的完整代码
from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.utils.trigger_rule import TriggerRule # 提前定义对应业务逻辑函数 def A(**context): pass def B(**context): pass def aS(**context): pass def aF(**context): pass def bS(**context): pass def bF(**context): pass with DAG( dag_id="dag_example", schedule_interval="@once", default_args={ "owner": "airflow", "start_date": datetime(2020, 11, 1), "retries": 1, "retry_delay": timedelta(seconds=10), }, catchup=False, ) as dag: tA = PythonOperator( task_id="A", python_callable=A, provide_context=True, ) tB = PythonOperator( task_id="B", python_callable=B, provide_context=True, ) taS = PythonOperator( task_id="aS", python_callable=aS, provide_context=True, trigger_rule=TriggerRule.ONE_SUCCESS ) taF = PythonOperator( task_id="aF", python_callable=aF, provide_context=True, trigger_rule=TriggerRule.ONE_FAILED ) tbS = PythonOperator( task_id="bS", python_callable=bS, provide_context=True, trigger_rule=TriggerRule.ONE_SUCCESS ) tbF = PythonOperator( task_id="bF", python_callable=bF, provide_context=True, trigger_rule=TriggerRule.ONE_FAILED ) # 按预期定义依赖 tA >> [tB, taS, taF] tB >> [tbS, tbF]
内容的提问来源于stack exchange,提问作者robert_553z8
相关产品推荐
相关产品推荐

