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

如何通过Dataset的extra参数结合数据感知调度传递run_date给下游DAG

用Airflow Dataset传递上游run_date到下游DAG的可行方案

方案一:把run_date嵌进Dataset URI里(通用所有Airflow版本)

Dataset本身没提供直接传参数的通道,但我们可以自己把run_date塞进Dataset的URI里,下游DAG触发后解析URI就能拿到这个日期,完全靠数据感知调度触发,不用额外任务或Operator。

上游DAG实现

把run_date格式化成URI的一部分,比如日期路径(可根据你的存储系统调整格式):

from airflow import DAG
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
from airflow.datasets import Dataset
from datetime import datetime

def make_upstream_dataset(run_date):
    # 将run_date转为字符串嵌入URI,这里以S3路径为例
    return Dataset(f"s3://your-bucket/processed-data/{run_date.strftime('%Y-%m-%d')}")

with DAG(
    dag_id="upstream_process",
    schedule=None,
    params={"run_date": datetime(2024, 5, 20)},
    catchup=False
) as dag:
    process_task = SQLExecuteQueryOperator(
        task_id="process_raw_data",
        sql="INSERT INTO processed_table SELECT * FROM raw_table WHERE dt = '{{ params.run_date }}'",
        # 用模板语法传入参数生成对应Dataset
        outlets=[make_upstream_dataset("{{ params.run_date }}")]
    )

下游DAG实现

下游DAG用通配符+匹配所有带日期的Dataset,触发后从dag_run配置中提取URI并解析run_date:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.datasets import Dataset
from datetime import datetime
import re

def pull_run_date(**context):
    # 获取触发当前DAG的Dataset信息
    triggered_dataset = context["dag_run"].conf["dataset_trigger"]["datasets"][0]
    # 用正则从URI中提取日期,规则对应上游URI格式
    date_match = re.search(r'/(\d{4}-\d{2}-\d{2})/', triggered_dataset["uri"])
    if date_match:
        run_date = date_match.group(1)
        # 将日期存入XCom供后续任务使用
        context["ti"].xcom_push(key="upstream_run_date", value=run_date)
        return run_date
    raise ValueError("无法从Dataset URI中解析run_date,请检查上游URI格式")

with DAG(
    dag_id="downstream_analysis",
    # 用通配符匹配所有符合格式的Dataset
    schedule=[Dataset("s3://your-bucket/processed-data/+")],
    catchup=False
) as dag:
    extract_date = PythonOperator(
        task_id="extract_upstream_run_date",
        python_callable=pull_run_date,
        provide_context=True
    )

    # 后续任务直接调用XCom中的日期
    analysis_task = PythonOperator(
        task_id="run_analysis",
        python_callable=lambda date: print(f"基于上游{date}的数据执行分析"),
        op_kwargs={"date": "{{ ti.xcom_pull(task_ids='extract_upstream_run_date') }}"}
    )

    extract_date >> analysis_task

方案二:用Dataset Event Metadata(Airflow 2.7+可用)

如果你的Airflow版本是2.7及以上,官方新增了Dataset事件的metadata功能,无需修改URI,直接将run_date附加到事件的extra字段即可。

上游DAG修改

在任务中生成Dataset事件时携带extra参数:

from airflow.datasets import DatasetEvent

def attach_run_date_to_dataset(run_date, **context):
    target_dataset = Dataset("s3://your-bucket/processed-data")
    # 将run_date放入extra字段
    event = DatasetEvent(
        dataset=target_dataset,
        extra={"run_date": run_date.strftime('%Y-%m-%d')}
    )
    # 将事件绑定到当前任务实例
    context["task_instance"].dataset_events.append(event)

# 把这个函数绑定到上游任务的post_execute阶段
process_task.post_execute = attach_run_date_to_dataset

下游DAG获取参数

直接从dag_run的配置中提取extra里的run_date:

def get_run_date_from_metadata(**context):
    run_date = context["dag_run"].conf["dataset_trigger"]["datasets"][0]["extra"]["run_date"]
    context["ti"].xcom_push(key="upstream_run_date", value=run_date)

你之前尝试Dataset的extra参数未成功,大概率是因为Airflow版本低于2.7,或者没有通过DatasetEvent的方式传递——直接在Dataset定义中添加extra是无效的,必须通过事件对象传递。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 10:13:12