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

Airflow:PythonOperator函数内使用{{ ds }}及传递BigQuery查询结果问题

问题与解决方案

问题描述

目前通过PythonOperator调用自定义函数、借助XCom实现任务间数据传递是可行的,但需要在BigQuery查询语句中使用执行日期宏{{ ds }}时,因语句嵌入函数内部无法被Jinja解析,导致查询试图访问名为some_table__{{ ds }}的表而失败。同时希望了解是否存在无需自定义函数即可传递BigQuery查询结果的方案,要求支持回溯运行、不能使用当前日期。


一、解决PythonOperator中的{{ ds }}解析问题

PythonOperator的自定义函数内部字符串不会被Jinja自动渲染,但可以通过context参数直接获取执行日期,手动替换到SQL语句中:

修改原get_some_values函数:

def get_some_values(**context):
    # 从context中提取执行日期ds
    execution_date = context['ds']
    hook = BigQueryHook(use_legacy_sql=False)
    conn = hook.get_conn()
    cursor = conn.cursor()
    # 用字符串格式化替换表名中的日期占位符
    cursor.execute(
        f"SELECT value1, value2, value3 FROM some_dataset.some_table__{execution_date}"
    )
    results = cursor.fetchone()
    
    # 将结果存入XCom
    if results is not None:
        for i, result in enumerate(results):
            context['ti'].xcom_push(f'value{i+1}', result)

此方式下,回溯运行时会自动使用对应历史执行日期,完全符合需求。


二、无需自定义函数的替代方案

直接使用Airflow原生Operator组合实现,全程无需编写自定义函数:

1. 用BigQueryOperator执行查询并推送结果到XCom

BigQueryOperator默认会将查询结果推送到XCom(key为return_value),只需确保开启do_xcom_push=True:

from airflow.providers.google.cloud.operators.bigquery import BigQueryOperator

get_values = BigQueryOperator(
    task_id='get_values',
    sql="""
        SELECT value1, value2, value3 FROM some_dataset.some_table__{{ ds }}
    """,
    use_legacy_sql=False,
    do_xcom_push=True,
    dag=dag
)

这里的SQL语句在Operator层面会被Jinja正常解析,{{ ds }}会自动替换为执行日期,回溯运行时也会匹配对应历史日期。

2. 用SlackWebhookOperator直接从XCom取值发送消息

通过Jinja模板在message参数中直接拉取前一个任务的XCom结果:

from airflow.providers.slack.operators.slack_webhook import SlackWebhookOperator

send_slack = SlackWebhookOperator(
    task_id='send_slack_message',
    http_conn_id=SLACK_CONN_ID,
    webhook_token=slack_webhook_token,
    channel='#some_channel',
    # 用Jinja语法从XCom拉取结果,[0][0]代表第一行第一列的值
    message="""values returned: {{ ti.xcom_pull(task_ids='get_values')[0][0] }}, {{ ti.xcom_pull(task_ids='get_values')[0][1] }}, {{ ti.xcom_pull(task_ids='get_values')[0][2] }} """,
    username='airflow',
    dag=dag
)

完整DAG示例

from airflow import DAG
from airflow.providers.google.cloud.operators.bigquery import BigQueryOperator
from airflow.providers.slack.operators.slack_webhook import SlackWebhookOperator
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2023, 1, 1)
}

dag = DAG(
    'test',
    default_args=default_args,
    schedule_interval='0 12 * * *',
    catchup=False,
)

get_values = BigQueryOperator(
    task_id='get_values',
    sql="""
        SELECT value1, value2, value3 FROM some_dataset.some_table__{{ ds }}
    """,
    use_legacy_sql=False,
    do_xcom_push=True,
    dag=dag
)

send_slack = SlackWebhookOperator(
    task_id='send_slack_message',
    http_conn_id=SLACK_CONN_ID,
    webhook_token=slack_webhook_token,
    channel='#some_channel',
    message="""values returned: {{ ti.xcom_pull(task_ids='get_values')[0][0] }}, {{ ti.xcom_pull(task_ids='get_values')[0][1] }}, {{ ti.xcom_pull(task_ids='get_values')[0][2] }} """,
    username='airflow',
    dag=dag
)

get_values >> send_slack

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 01:01:51