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

Airflow新手求助:如何从XCom获取列表并循环执行DAG任务?

解决方案:Airflow动态循环执行任务(基于XCom动态列表)

核心思路

Airflow在DAG定义阶段是静态解析的,无法直接获取运行时生成的XCom数据,因此不能用硬编码循环的方式生成任务。推荐使用Dynamic Task Mapping(Airflow 2.2+支持),它能在任务运行阶段解析XCom数据并动态生成任务实例,完美适配你的需求。

方案1:并行执行动态任务(默认)

如果不需要严格顺序执行,直接用expand方法生成并行任务:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.dummy import DummyOperator
from datetime import datetime

# 前置任务:生成要循环的列表并推送到XCom
def generate_target_list(**context):
    # 替换为你的实际数据来源(如数据库查询、API返回等)
    target_list = ["first", "second", "third"]
    context['ti'].xcom_push(key='some_list', value=target_list)

# 单个循环任务的逻辑
def process_single_item(item, **context):
    # 替换为你的实际任务逻辑,如调用业务接口、处理文件等
    print(f"Processing item: {item}")
    # 可选:将处理结果推送到XCom
    context['ti'].xcom_push(key=f'result_{item}', value=f'completed_{item}')

with DAG(
    dag_id='dynamic_parallel_tasks',
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    start = DummyOperator(task_id='start')

    # 生成列表的前置任务
    generate_list_task = PythonOperator(
        task_id='generate_target_list',
        python_callable=generate_target_list,
        provide_context=True
    )

    # 动态映射生成任务实例
    dynamic_processing_task = PythonOperator.partial(
        task_id='process_item',
        python_callable=process_single_item,
        provide_context=True
    ).expand(
        # 从XCom拉取动态列表
        item="{{ ti.xcom_pull(task_ids='generate_target_list', key='some_list') }}"
    )

    end = DummyOperator(task_id='end')

    # 设置依赖链
    start >> generate_list_task >> dynamic_processing_task >> end

方案2:顺序执行动态任务(匹配你的链式需求)

如果需要像你之前代码那样按顺序逐个执行任务,只需在动态任务中添加depends_on_past=True,让每个任务实例依赖前一个实例的成功:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.dummy import DummyOperator
from datetime import datetime

def generate_target_list(**context):
    target_list = ["first", "second", "third"]
    context['ti'].xcom_push(key='some_list', value=target_list)

def process_single_item(item, **context):
    print(f"Processing item: {item}")
    context['ti'].xcom_push(key=f'result_{item}', value=f'completed_{item}')

with DAG(
    dag_id='dynamic_sequential_tasks',
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    start = DummyOperator(task_id='start')

    generate_list_task = PythonOperator(
        task_id='generate_target_list',
        python_callable=generate_target_list,
        provide_context=True
    )

    dynamic_processing_task = PythonOperator.partial(
        task_id='process_item',
        python_callable=process_single_item,
        provide_context=True,
        depends_on_past=True  # 强制顺序执行,前一个实例完成才执行下一个
    ).expand(
        item="{{ ti.xcom_pull(task_ids='generate_target_list', key='some_list') }}"
    )

    end = DummyOperator(task_id='end')

    start >> generate_list_task >> dynamic_processing_task >> end

替代方案:触发外部DAG循环执行

如果你的Airflow版本较低(低于2.2),无法使用动态映射,可以用TriggerDagRunOperator触发子DAG处理单个项,结合PythonOperator控制循环:

子DAG(处理单个项)

# child_dag.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime

def process_item(item, **context):
    print(f"Processing item in child DAG: {item}")

with DAG(
    dag_id='child_processing_dag',
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    process_task = PythonOperator(
        task_id='process_single_item',
        python_callable=process_item,
        op_kwargs={'item': "{{ dag_run.conf['item'] }}"},
        provide_context=True
    )

主DAG(生成列表并触发子DAG)

# parent_dag.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.dummy import DummyOperator
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from datetime import datetime

def generate_target_list(**context):
    target_list = ["first", "second", "third"]
    context['ti'].xcom_push(key='some_list', value=target_list)

def trigger_child_dags(**context):
    target_list = context['ti'].xcom_pull(task_ids='generate_target_list', key='some_list')
    prev_task = None
    for idx, item in enumerate(target_list):
        trigger_task = TriggerDagRunOperator(
            task_id=f'trigger_child_{idx}',
            trigger_dag_id='child_processing_dag',
            conf={'item': item},
            wait_for_completion=True  # 等待子DAG完成再执行下一个
        )
        if prev_task:
            prev_task >> trigger_task
        else:
            context['task_instance'].set_upstream(context['ti'].task)
        prev_task = trigger_task
    # 将最后一个触发任务的下游设为end任务
    prev_task.set_downstream(context['dag'].get_task('end'))

with DAG(
    dag_id='parent_trigger_dag',
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    start = DummyOperator(task_id='start')

    generate_list_task = PythonOperator(
        task_id='generate_target_list',
        python_callable=generate_target_list,
        provide_context=True
    )

    trigger_task = PythonOperator(
        task_id='trigger_child_dags',
        python_callable=trigger_child_dags,
        provide_context=True
    )

    end = DummyOperator(task_id='end')

    start >> generate_list_task >> trigger_task

关键说明

  • 动态映射是Airflow 2.2+的官方推荐方案,比手动循环更简洁、更稳定,完全避免了DAG定义阶段获取XCom的问题。
  • depends_on_past=True仅适用于顺序执行的动态任务,确保任务按列表顺序逐个完成。
  • 替代方案适用于老版本Airflow,但维护成本较高,优先推荐动态映射。

内容的提问来源于stack exchange,提问作者Prashant singh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 01:35:04