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

如何查询更新Google Cloud Composer Airflow Dataset的最后一个DAG

实现方案

1. 在更新Dataset的DAG中添加触发源标记

当两个DAG更新//my_Dataset时,需要在更新操作里带上当前DAG的ID作为标记,这样后续才能追踪到是谁触发的更新。

比如在第一个更新DAG(dag_a)中:

from airflow.models import Dataset
from airflow.operators.python import PythonOperator

dataset = Dataset("//my_Dataset")

def update_dataset_with_tag():
    # 更新Dataset时,在extra字段中传入当前DAG ID
    dataset.update(extra={"triggering_dag": "dag_a"})

update_task = PythonOperator(
    task_id="update_my_dataset",
    python_callable=update_dataset_with_tag,
    dag=dag_a
)

第二个更新DAG(dag_b)同理,只需把triggering_dag的值改为"dag_b"即可。

2. 在my_dag中获取最新触发源并设置参数

在my_dag里,通过查询Airflow元数据库中的DatasetEvent表,获取该Dataset的最新更新记录,从中提取触发源DAG ID,再根据不同的DAG设置对应参数。

完整的my_dag代码示例:

from airflow.models import Dataset, DAG, PythonOperator, DatasetEvent
from airflow.utils.db import create_session
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1)
}

dataset = Dataset("//my_Dataset")

def fetch_trigger_source_and_set_params(**context):
    with create_session() as session:
        # 查询最新的Dataset更新事件
        latest_update = session.query(DatasetEvent)\
            .filter(DatasetEvent.dataset_uri == dataset.uri)\
            .order_by(DatasetEvent.timestamp.desc())\
            .first()
        
        # 根据触发源设置不同参数
        if latest_update and "triggering_dag" in latest_update.extra:
            trigger_dag = latest_update.extra["triggering_dag"]
            if trigger_dag == "dag_a":
                context['ti'].xcom_push(key='runtime_params', value={"process_mode": "fast", "threshold": 0.8})
            elif trigger_dag == "dag_b":
                context['ti'].xcom_push(key='runtime_params', value={"process_mode": "full", "threshold": 0.5})
        else:
            # 无更新记录时使用默认参数
            context['ti'].xcom_push(key='runtime_params', value={"process_mode": "default", "threshold": 0.6})

def execute_with_params(**context):
    # 从XCom中获取参数并执行任务逻辑
    params = context['ti'].xcom_pull(key='runtime_params', task_ids='fetch_trigger_source')
    print(f"执行任务,参数: {params}")
    # 这里添加你的实际业务逻辑

dag = DAG(
    dag_id='my_dag',
    default_args=default_args,
    schedule=[dataset],  
    catchup=False
)

fetch_trigger_source = PythonOperator(
    task_id="fetch_trigger_source",
    python_callable=fetch_trigger_source_and_set_params,
    provide_context=True,
    dag=dag
)

execute_task = PythonOperator(
    task_id="execute_task",
    python_callable=execute_with_params,
    provide_context=True,
    dag=dag
)

fetch_trigger_source >> execute_task

关键说明

  • Dataset.update()方法的extra字段是自定义扩展字段,用来存储触发源信息,默认不会自动记录,必须手动传入。
  • DatasetEvent表存储了所有Dataset的更新事件,通过查询最新的记录就能拿到最后一次更新的触发源。
  • 用XCom在任务间传递参数,实现触发源与业务参数的关联。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 18:40:00