Airflow DAG运行启动时间各任务不一致问题及打印方案咨询
Airflow DAG时间变量不一致问题解决方案
问题根因
你当前代码中datetime.now()是在DAG文件被调度器解析时执行的,不是DAG实际运行时计算。Airflow调度器会周期性扫描解析DAG文件,每次解析都会更新该时间值,导致同一个DAG Run的不同任务拿到的时间不一致,甚至出现上下游任务文件名不匹配的运行错误。
正确实现方案
1. 解决时间不一致问题
使用Airflow内置的Jinja模板变量,同一个DAG Run内的所有任务会共享同一个变量值,支持直接在运算符的模板参数中使用:
- 若需要DAG逻辑执行时间(调度的计划时间,同一个DAG Run全局唯一不变),用
{{ execution_date }} - 若需要DAG实际触发运行的时间,用
{{ dag_run.start_date }}
修正后的完整代码如下:
from datetime import datetime from airflow import DAG # 导入你用到的运算符 from airflow.providers.google.cloud.operators.gcs import ContentToGoogleCloudStorageOperator from airflow.providers.google.cloud.operators.bigquery import GoogleCloudStorageToBigQueryOperator dag = DAG( dag_id='maindag', start_date=datetime(2021, 11, 1), max_active_runs=1, catchup=False # 建议加上,避免历史未调度的DAG Run被批量触发 ) transfer_to_gcs = ContentToGoogleCloudStorageOperator( task_id='transfer_to_gcs', content=getdata('people'), dst='people_{{ execution_date.strftime("%Y-%m-%d %H:%M:%S") }}.json', bucket='bucketname', dag=dag ) transfer_to_bq = GoogleCloudStorageToBigQueryOperator( task_id='transfer_to_bq', bucket='bucketname', source_objects=['people_{{ execution_date.strftime("%Y-%m-%d %H:%M:%S") }}.json'], dag=dag ) transfer_to_gcs >> transfer_to_bq
2. 打印DAG启动时间
可以新增一个PythonOperator任务,通过任务上下文获取时间参数并打印,代码示例如下:
from airflow.operators.python import PythonOperator def print_run_time(**context): # 获取逻辑执行时间 exec_time = context["execution_date"].strftime("%Y-%m-%d %H:%M:%S") # 获取DAG实际启动时间 actual_start_time = context["dag_run"].start_date.strftime("%Y-%m-%d %H:%M:%S") print(f"DAG逻辑执行时间:{exec_time}") print(f"DAG实际启动时间:{actual_start_time}") print_time_task = PythonOperator( task_id="print_run_time", python_callable=print_run_time, dag=dag ) # 调整依赖,先执行打印任务 print_time_task >> transfer_to_gcs >> transfer_to_bq
任务运行后可以在该任务的日志中看到打印的时间值。
内容的提问来源于stack exchange,提问作者dariuszewski
相关产品推荐
相关产品推荐

