Apache Airflow:每日DAG指定日期加载及过往日期手动重跑问题
Apache Airflow 问题解答
一、每日DAG加载D-1数据的正确姿势
别用WHERE order_date = today-1,这会踩坑:手动重跑旧日期任务时,today是当前日期,会错误加载当前前一天的数据,而非你要补跑的目标日期的前一天。
正确做法是利用Airflow的execution_date:对于@daily调度的DAG,Airflow默认把execution_date设为任务要处理的日期(比如今天运行的任务,execution_date是昨天,正好对应你需要的D-1数据日期)。
代码修正要点
- 模板变量正确传递:用
{{ ds }}(execution_date的日期字符串,格式为YYYY-MM-DD)或者{{ execution_date.date() }}来过滤,BigQuery可直接识别这种格式。 - format参数顺序修正:你的SQL占位符顺序是「目标表 → 源表 → 过滤日期」,之前的format参数顺序颠倒,会导致表名匹配错误。
- PythonOperator配置:要启用模板功能,需设置
provide_context=True(Airflow 1.x);另外别忘了导入PythonOperator,任务ID不要包含空格。
修正后的核心代码片段:
from airflow.operators.python import PythonOperator # 补上缺失的导入 def report_sales_summary(**context): # 加入context参数获取模板变量 credential = service_account.Credentials.from_service_account_file("/db-config") project_id = dbconfig.project client = bigquery.Client(credentials=credential, project=project_id) # 从上下文获取execution_date的日期字符串 execution_date_str = context['ds'] query = """ INSERT INTO `{}.{}.{}` WITH BASE AS ( SELECT *, EXTRACT(DATE FROM (TIMESTAMP(PARSE_DATETIME("%m/%d/%Y %H:%M", order_created_date),'Asia'))) order_date FROM `{}.{}.{}` WHERE order_status = 'sale' ), A AS ( SELECT order_date, COUNT(customer) total_customer, SUM(CAST(unit_price AS FLOAT64)) total_price FROM BASE GROUP BY order_date ), B AS ( SELECT order_date, order_type product_category, COUNT(product_sku) total_qty FROM BASE GROUP BY order_date,product_category ) SELECT A.order_date, total_customer, product_category, total_qty, total_price FROM A JOIN B ON A.order_date = B.order_date WHERE A.order_date = '{}' # 用单引号包裹日期字符串 """.format( dbconfig.project, dbconfig.target_dataset, dbconfig.target_table, dbconfig.project, dbconfig.source_dataset, dbconfig.source_table, execution_date_str ) client.query(query) # 执行BigQuery查询 # DAG定义修正 sales_summary_dag = DAG( 'report_sales_summary_dag', default_args=default_dag_args, schedule_interval='@daily', # 加引号更规范 catchup=False ) report_sales_summary_task = PythonOperator( task_id='report_sales_summary', # 移除任务ID中的空格 python_callable=report_sales_summary, provide_context=True, # 启用上下文传递以获取模板变量 dag=sales_summary_dag # 指定所属DAG )
二、手动重跑过往日期的DAG
方式1:Airflow UI操作
- 进入目标DAG的详情页,点击右上角的Trigger DAG w/ config。
- 在弹出窗口的Execution Date选项中,选择需要重跑的过往日期(格式为
YYYY-MM-DD),点击Trigger即可触发对应日期的任务实例。 - 批量重跑多个日期:在DAG的Tree View或Graph View中,勾选对应日期的任务实例,点击Clear清除历史运行记录,若DAG设置了
catchup=True会自动补跑;也可直接批量选中后点击Trigger手动触发。
方式2:Airflow CLI命令
- 触发单个日期的任务:
其中airflow dags trigger -e 2024-05-20 report_sales_summary_dag-e指定要重跑的execution_date,后面跟随你的DAG ID。 - 批量补跑多个日期:
先将DAG的catchup设为True,重启Airflow调度器,调度器会自动补跑start_date到当前日期之间所有未执行的任务;补跑完成后改回catchup=False即可避免后续自动补跑。
内容的提问来源于stack exchange,提问作者random student
相关产品推荐
相关产品推荐

