You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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:
    @task.external_python(task_id="task1", python='/opt/airflow/my_env/bin/python3',
    trigger_rule=TriggerRule.ALL_DONE)
    
    错误1:
    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)
    
    错误2:
    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. 错误1:@task.external_python装饰器不支持直接传入trigger_rule参数,强行传递会导致装饰器返回的对象无法使用>>操作符构建任务依赖。
  2. 错误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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 00:45:58