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

Airflow消费者DAG传入default_args后任务未执行问题求助

问题解答:Airflow数据集感知调度中消费者DAG添加default_args后任务不执行

结论:完全可以为消费者DAG传递default_args,你的问题源于代码中的变量名拼写错误,而非default_args本身的限制。

核心错误分析

你在生产者和消费者DAG的定义中,都犯了一个低级变量名拼写错误:

  • 你定义的默认参数变量是default_args,但在初始化DAG时写的是default_args = default_dags,default_dags是未定义的变量。
  • 这个错误会导致DAG初始化异常,Airflow可能不会直接抛出报错,但会生成一个有缺陷的DAG实例——表现为DAG显示“运行完成”但实际任务未被调度执行。

修复步骤

  1. 修正变量名拼写:将DAG初始化中的default_args = default_dags改为default_args = default_args。
  2. 可选优化:将start_date的元组格式改为datetime对象(Airflow推荐写法,避免潜在的时区或解析问题),需要先导入datetime模块。
  3. 强制要求:生产者与消费者DAG的dag_id必须唯一,否则会互相覆盖引发异常。

修复后的代码示例

生产者DAG

from airflow.datasets import Dataset
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
from airflow.operators.empty import EmptyOperator
from datetime import datetime

MY_DATA = Dataset('bigquery://my-project-name/my-schema/-my-table')

data_set_operator = EmptyOperator(
    task_id="producer",
    outlets=MY_DATA
)

default_args = {
    "start_date": datetime(2024, 11, 20),
    "depends_on_past": False,
    "on_failure_callback": some_function
}

with DAG(
    dag_id="my_dag_producer",
    max_active_runs=1,
    default_args=default_args,
    schedule_interval="30 8 * * *",
) as dag:
    sql_task = SQLExecuteQueryOperator(
        task_id="sql_task_producer",
        query="my_query",
        conn_id="bq_conn_id",
        params=my_dictionary
    )

    sql_task >> data_set_operator

消费者DAG

from airflow.datasets import Dataset
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
from datetime import datetime

MY_DATA = Dataset('bigquery://my-project-name/my-schema/-my-table')

default_args = {
    "start_date": datetime(2024, 11, 20),
    "depends_on_past": False,
    "on_failure_callback": some_function
}

with DAG(
    dag_id="my_dag_consumer",
    max_active_runs=1,
    default_args=default_args,
    schedule=MY_DATA
) as dag:
    sql_task = SQLExecuteQueryOperator(
        task_id="sql_task_consumer",
        query="my_query2",
        conn_id="bq_conn_id",
        params=my_dictionary2
    )

额外注意事项

  • on_failure_callback中引用的some_function必须是已定义或导入的可调用函数,否则会导致DAG初始化失败。
  • 使用Dataset调度时,消费者DAG的start_date需要早于生产者DAG的执行时间,否则不会触发调度。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 22:13:14