如何在Airflow中用Python实现任务串行执行且不依赖上游结果
问题描述
- 如何串行运行多个ExternalPythonOperator(不同DAG任务需要不同的包/版本),且不依赖上游任务的执行结果(即使上游失败也继续执行)。
- 要求任务依次执行,无论之前的任务成功或失败。
- 不创建单独DAG文件的原因:希望在与其他任务完全错开的时间段内依次运行几个资源密集型任务,避免干扰;同时这些任务需要彼此隔离,防止服务器资源限制和外部因素导致互相干扰。
我的代码
import logging import os import shutil import sys import tempfile import time from pprint import pprint import pendulum from airflow import DAG from airflow.decorators import task log = logging.getLogger(__name__) PYTHON = sys.executable BASE_DIR = tempfile.gettempdir() my_default_args = { 'owner': 'me', #'email': ['myemail@myemail.com'], 'email_on_failure': True, #'email_on_retry': True, #'retries': 1, # 'retry_delay': timedelta(minutes=1) } with DAG( dag_id='some_dag_id_comes_here', schedule='1 * * * *', start_date=pendulum.datetime(2021, 1, 1, tz="UTC"), # this is from whre it starts counting time to run taks, NOT like cron catchup=False, default_args=my_default_args, tags=['xyz1'], ) as dag: @task.external_python(task_id="task1", python='/opt/airflow/my_env/bin/python3') def func1(): print('elements of task 1') time.sleep(10) @task.external_python(task_id="task2", python='/opt/airflow/my_env/bin/python3') def func2(): print('elements of task 2') time.sleep(10) task1 >> task2
已尝试方案
- 尝试过跨DAG触发依赖的相关方案,但需求完全相反,不适用。
- 尝试过Airflow官方文档中的触发规则,遇到两个错误:
代码1:
错误1:@task.external_python(task_id="task1", python='/opt/airflow/my_env/bin/python3', trigger_rule=TriggerRule.ALL_DONE)
代码2:Broken DAG: [/opt/airflow/dags/test_file.py] Traceback (most recent call last): File "<frozen importlib._bootstrap>", line 219, in _call_with_frames_removed File "/opt/airflow/dags/name.py", line 224, in <module> task1 >> task2 >> task3 >> task4 >> task5 TypeError: unsupported operand type(s) for >>: '_TaskDecorator' and '_TaskDecorator'
错误2:@task.external_python(task_id="task1", python='/opt/airflow/my_env/bin/python3', trigger_rule=TriggerRule.none_skipped)Broken DAG: [/opt/airflow/dags/test_file.py] Traceback (most recent call last): File "/opt/airflow/dags/test_file.py", line 51, in <module> ,trigger_rule=TriggerRule.none_skipped File "/usr/local/lib/python3.8/enum.py", line 384, in __getattr__ raise AttributeError(name) from None AttributeError: none_skipped
解决方案
错误原因分析
- 错误1:
@task.external_python装饰器不支持直接传入trigger_rule参数,强行传递会导致装饰器返回的对象无法使用>>操作符构建任务依赖。 - 错误2:
TriggerRule是枚举类,成员名称需大写,正确值应为TriggerRule.NONE_SKIPPED而非小写的none_skipped。
修改后的代码
改用传统的ExternalPythonOperator类定义任务,直接设置触发规则,确保下游任务在上游完成(无论成功失败)后执行:
import logging import os import sys import tempfile import time import pendulum from airflow import DAG from airflow.operators.python import ExternalPythonOperator from airflow.utils.trigger_rule import TriggerRule log = logging.getLogger(__name__) PYTHON = sys.executable BASE_DIR = tempfile.gettempdir() my_default_args = { 'owner': 'me', 'email_on_failure': True, } def func1(): print('elements of task 1') time.sleep(10) # 可故意抛出异常测试:raise Exception("Task1 failed") def func2(): print('elements of task 2') time.sleep(10) with DAG( dag_id='some_dag_id_comes_here', schedule='1 * * * *', start_date=pendulum.datetime(2021, 1, 1, tz="UTC"), catchup=False, default_args=my_default_args, tags=['xyz1'], ) as dag: task1 = ExternalPythonOperator( task_id="task1", python='/opt/airflow/my_env/bin/python3', python_callable=func1 ) task2 = ExternalPythonOperator( task_id="task2", python='/opt/airflow/my_env/bin/python3', python_callable=func2, trigger_rule=TriggerRule.ALL_DONE # 关键设置:上游完成即执行,不管成功失败 ) task1 >> task2
扩展说明
TriggerRule.ALL_DONE是满足需求的触发规则:只要上游任务完成(无论成功、失败或跳过),下游任务就会执行。- 若需要多个串行任务,只需给每个后续任务都添加
trigger_rule=TriggerRule.ALL_DONE,例如:def func3(): print('elements of task 3') time.sleep(10) task3 = ExternalPythonOperator( task_id="task3", python='/opt/airflow/another_env/bin/python3', # 不同虚拟环境 python_callable=func3, trigger_rule=TriggerRule.ALL_DONE ) task2 >> task3
内容的提问来源于stack exchange,提问作者sogu
相关产品推荐
相关产品推荐

