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

Airflow 2.0如何在Operator/Sensor外部访问execution_date等宏变量

问题原因

方案1 KeyError原因

  • 核心错误:你给python_callable参数传的是wait_for_data()(带括号),相当于DAG解析阶段就直接执行了该函数,此时任务上下文还未生成,自然拿不到ds等宏变量。正确写法是只传函数名wait_for_data,把函数对象交给Operator在运行时调用。
  • 上下文传参对象错误:provide_context(Airflow 2.x已默认开启,无需手动配置)会把上下文参数传递给python_callable对应的执行函数(即wait_for_data、run_aggregation),你需要给这些执行函数加上**kwargs参数,而不是给包裹Operator的外层函数加。
  • 直接写在Python函数f-string中的{{}}不会被Airflow渲染,只有Operator声明的模板字段才会触发宏解析,你在普通Python字符串里写的{{ prev_ds }}只会被当成普通文本输出。

方案2 TypeError原因

  • 字典取值语法错误:你写的kwargs.get(['templates_dict'])传入了列表作为key,字典的key必须是字符串,正确写法为kwargs.get('templates_dict')。
  • 取值key不匹配:你在templates_dict中定义的日期key是end_date对应{{ ds }},后续你尝试用ds作为key取值自然拿不到对应值。

关于Operator外调用宏的说明

Airflow的宏渲染仅发生在任务运行前的模板解析阶段,仅Operator声明的template_fields字段内容会被解析,无法在普通Python函数中直接调用宏,必须通过上下文把渲染完成的参数传递进去。

正确实现方案

方案1:适配现有代码的最小修改

你当前的代码混用了TaskFlow API的@task装饰器和传统Operator实例化,属于冗余写法,直接用@task系列装饰器即可,不需要手动创建PythonOperator/PythonSensor:

def create_db_engine():
    [redacted]
    return engine

def run_query(sql):
    engine = create_db_engine()
    with engine.connect() as connection:
        data = connection.execute(sql)
        return data

# 传感器用@sensor装饰器,直接写执行逻辑即可
@task.sensor(poke_interval=30, timeout=3600)
def data_sensor(**kwargs):
    # 直接从上下文拿到渲染完成的日期参数
    ds = kwargs["ds"]
    sensor_query = f'''
    select blah
    from table
    where dt = '{ds}'
    limit 1
    '''
    return run_query(sensor_query).rowcount >= 1

@task
def agg_data(**kwargs):
    prev_ds = kwargs["prev_ds"]
    ds = kwargs["ds"]
    agg_query = f'''
    delete from table
    where datefield = '{prev_ds}'::DATE;
    
    insert into table(datefield, metric)
    select date_trunc('day', timefield), sum(metric)
    from sourcetable
    where timefield >= '{prev_ds}'::DATE
    and timefield < '{ds}'::DATE
    group by 1;
    '''
    run_query(agg_query)

# DAG内定义依赖
with DAG(...) as dag:
    data_sensor() >> agg_data()

方案2:更规范的Redshift原生算子实现

Airflow的Redshift provider原生支持IAM角色链认证,不需要自己维护SQLAlchemy连接逻辑,配置完成后直接用RedshiftSQLOperator即可原生支持宏渲染,代码更简洁:

  1. 先在Airflow连接管理中创建Redshift连接:连接类型选择Amazon Redshift,勾选Use IAM Authentication,填写集群ID、IAM角色ARN、Region、数据库名等参数即可。
  2. 直接调用算子写SQL,宏可以直接写在SQL字符串中,Airflow会自动渲染:
from airflow.providers.amazon.aws.operators.redshift_sql import RedshiftSQLOperator
from airflow.providers.amazon.aws.sensors.redshift_sql import RedshiftSQLSensor

with DAG(...) as dag:
    check_data = RedshiftSQLSensor(
        task_id="check_data",
        redshift_conn_id="你配置的Redshift连接ID",
        sql="select blah from table where dt = '{{ ds }}' limit 1",
        poke_interval=30,
        timeout=3600
    )

    delete_old_data = RedshiftSQLOperator(
        task_id="delete_old_data",
        redshift_conn_id="你配置的Redshift连接ID",
        sql="delete from table where datefield = '{{ prev_ds }}'::DATE"
    )

    insert_agg_data = RedshiftSQLOperator(
        task_id="insert_agg_data",
        redshift_conn_id="你配置的Redshift连接ID",
        sql="""
        insert into table(datefield, metric)
        select date_trunc('day', timefield), sum(metric)
        from sourcetable
        where timefield >= '{{ prev_ds }}'::DATE
        and timefield < '{{ ds }}'::DATE
        group by 1;
        """
    )

    check_data >> delete_old_data >> insert_agg_data

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 04:30:05