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

如何将函数返回的变量传递给DAG及依赖DAG的任务

实现日期参数跨任务、跨DAG传递方案

核心思路

利用Airflow的XCom实现同DAG内任务间的参数传递,再通过TriggerDagRunOperator的conf参数实现跨DAG的参数传递。下面是具体实现步骤:

1. 补全基础依赖与日期生成任务

先完善日期生成函数,并用PythonOperator创建任务生成日期——该算子默认会把函数返回值推送到XCom(默认key为return_value):

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from datetime import datetime, timedelta

# 根据实际需求定义要减去的天数
days_to_subtract = 1

def date_fn():
    d = datetime.today() - timedelta(days=days_to_subtract)
    # 转成字符串避免datetime对象序列化问题
    return d.strftime("%Y-%m-%d")

2. 修改processthe_files DAG,实现同DAG内参数传递

在processthe_files DAG中添加日期生成任务,让内部任务拉取XCom的日期参数,同时配置触发跨DAG的参数传递:

SCHEDULE = "0 8 * * 4"

with DAG(
    dag_id="processthe_files",
    start_date=datetime(2024, 10, 8),
    schedule_interval=SCHEDULE,
    catchup=False
) as dag:
    # 1. 生成日期并推送到XCom
    generate_date = PythonOperator(
        task_id="generate_date",
        python_callable=date_fn
    )

    # 2. file_processing任务拉取XCom的日期参数
    # 假设你的Job类支持通过parameters传递参数,根据实际情况调整
    file_processing = Job(
        parameters={"target_date": "{{ ti.xcom_pull(task_ids='generate_date') }}"}
    ).to_task

    # 3. 触发processtables DAG时,通过conf传递日期
    trigger_processtables = TriggerDagRunOperator(
        task_id='trigger_processtables',
        trigger_dag_id='processtables',  # 和目标DAG的dag_id保持一致
        wait_for_completion=True,
        # 将XCom的日期放入conf,传递给目标DAG
        conf={"target_date": "{{ ti.xcom_pull(task_ids='generate_date') }}"},
        dag=dag
    )

    # 设置任务依赖链
    generate_date >> file_processing >> trigger_processtables

3. 修改processtables DAG,接收跨DAG传递的参数

在processtables DAG的任务中,从dag_run.conf中取出传递过来的日期参数:

with DAG(
    dag_id="processtables",  # 和TriggerDagRunOperator中的trigger_dag_id一致
    start_date=datetime(2024, 10, 8),
    schedule_interval=None,
    catchup=False
) as dag:
    processtables = Job(
        # 从dag_run.conf中获取日期参数
        parameters={"target_date": "{{ dag_run.conf.get('target_date') }}"}
    ).to_task

    processtables

关键细节说明

  • XCom自动推送:PythonOperator默认会把函数返回值推送到XCom,无需手动调用ti.xcom_push(),通过task_ids即可指定拉取的任务。
  • 日期序列化:将datetime对象转为字符串传递,避免Airflow序列化datetime对象时出现兼容性问题。
  • conf参数传递:TriggerDagRunOperator的conf会作为触发目标DAG的配置,目标DAG可通过dag_run.conf访问这些参数。
  • 原代码修正:你原代码中processthefiles >> trigger_processtables是变量名错误,已修正为file_processing >> trigger_processtables。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 19:46:18