如何查询更新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
相关产品推荐
相关产品推荐

