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

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数据日期)。

代码修正要点

  1. 模板变量正确传递:用{{ ds }}(execution_date的日期字符串,格式为YYYY-MM-DD)或者{{ execution_date.date() }}来过滤,BigQuery可直接识别这种格式。
  2. format参数顺序修正:你的SQL占位符顺序是「目标表 → 源表 → 过滤日期」,之前的format参数顺序颠倒,会导致表名匹配错误。
  3. 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操作

  1. 进入目标DAG的详情页,点击右上角的Trigger DAG w/ config。
  2. 在弹出窗口的Execution Date选项中,选择需要重跑的过往日期(格式为YYYY-MM-DD),点击Trigger即可触发对应日期的任务实例。
  3. 批量重跑多个日期:在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 14:07:50