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

如何将Airflow Jinja变量值传入BigQuery Operator内的函数

问题分析与解决方案

核心问题

你直接将'{{ ds_nodash }}'作为参数传给gen_pg_count_update_query函数,导致函数拿到的是原始模板字符串而非实际日期值——因为Airflow的Jinja模板渲染发生在任务运行阶段,但函数在DAG解析阶段就会执行,此时模板还未被渲染。另外,ds_nodash格式为YYYYMMDD,但你函数里用的是%Y%m%dT%H%M%S(对应ts_nodash的格式),格式不匹配也会引发strptime报错。


解决方案

方案1:直接用Jinja模板实现日期计算

将日期逻辑移到Jinja模板字符串中,Airflow会在任务运行时自动渲染变量:

from airflow.utils.dates import timedelta

# 全局配置变量(根据实际值调整)
max_pipeline_delay = 30
frequency_hour = 24

update_pg_count = BigQueryInsertJobOperator(
    task_id="pg_count_insert_job",
    configuration={
        "query": {
            "query": """
                {% set current_datetime = macros.datetime.strptime(ds_nodash, '%Y%m%d') %}
                {% set datetime_to = current_datetime - macros.timedelta(minutes=max_pipeline_delay) %}
                {% set datetime_from = datetime_to - macros.timedelta(hours=frequency_hour) %}
                {% set unix_from = (datetime_from - macros.datetime(1970,1,1)).total_seconds() * 1000 %}
                {% set unix_to = (datetime_to - macros.datetime(1970,1,1)).total_seconds() * 1000 %}
                INSERT INTO table1 (col1, col2)
                SELECT
                    (SELECT COUNT(*) FROM tablex WHERE updateat BETWEEN '{{ datetime_from }}' AND '{{ datetime_to }}'),
                    (SELECT COUNT(*) FROM tabley WHERE updateat_unix BETWEEN {{ unix_from }} AND {{ unix_to }})
            """,
            "useLegacySql": False,
        }
    },
    location=bq_region,
    dag=dag
)

方案2:用PythonOperator生成查询后传给BigQuery

先通过Python任务计算查询语句并存入XCom,再由BigQuery任务读取:

from airflow.operators.python import PythonOperator
import datetime
import time

def gen_pg_count_update_query(**context):
    # 从运行上下文获取实际日期值
    ds_nodash_val = context['ds_nodash']
    current_datetime = datetime.datetime.strptime(ds_nodash_val, "%Y%m%d")
    datetime_to = current_datetime - datetime.timedelta(minutes=max_pipeline_delay)
    datetime_from = datetime_to - datetime.timedelta(hours=frequency_hour)
    unix_from = time.mktime(datetime_from.timetuple()) * 1000
    unix_to = time.mktime(datetime_to.timetuple()) * 1000

    return f"""
        INSERT INTO table1 (col1, col2)
        SELECT
            (SELECT COUNT(*) FROM tablex WHERE updateat BETWEEN '{datetime_from}' AND '{datetime_to}'),
            (SELECT COUNT(*) FROM tabley WHERE updateat_unix BETWEEN {unix_from} AND {unix_to})
    """

# 生成查询的Python任务
generate_query = PythonOperator(
    task_id="generate_pg_count_query",
    python_callable=gen_pg_count_update_query,
    provide_context=True,
    dag=dag
)

# 执行BigQuery插入任务
update_pg_count = BigQueryInsertJobOperator(
    task_id="pg_count_insert_job",
    configuration={
        "query": {
            "query": "{{ ti.xcom_pull(task_ids='generate_pg_count_query') }}",
            "useLegacySql": False,
        }
    },
    location=bq_region,
    dag=dag
)

# 设置任务依赖
generate_query >> update_pg_count

方案3:在函数内部获取运行时上下文

通过Airflow内置方法获取运行时变量,避免直接传模板字符串:

from airflow.operators.python import get_current_context
import datetime
import time

def gen_pg_count_update_query():
    context = get_current_context()
    ds_nodash_val = context['ds_nodash']
    current_datetime = datetime.datetime.strptime(ds_nodash_val, "%Y%m%d")
    datetime_to = current_datetime - datetime.timedelta(minutes=max_pipeline_delay)
    datetime_from = datetime_to - datetime.timedelta(hours=frequency_hour)
    unix_from = time.mktime(datetime_from.timetuple()) * 1000
    unix_to = time.mktime(datetime_to.timetuple()) * 1000

    return f"""
        INSERT INTO table1 (col1, col2)
        SELECT
            (SELECT COUNT(*) FROM tablex WHERE updateat BETWEEN '{datetime_from}' AND '{datetime_to}'),
            (SELECT COUNT(*) FROM tabley WHERE updateat_unix BETWEEN {unix_from} AND {unix_to})
    """

update_pg_count = BigQueryInsertJobOperator(
    task_id="pg_count_insert_job",
    configuration={
        "query": {
            "query": "{{ gen_pg_count_update_query() }}",
            "useLegacySql": False,
        }
    },
    location=bq_region,
    dag=dag,
    # 将函数加入模板全局变量
    templates_dict={"gen_pg_count_update_query": gen_pg_count_update_query}
)

注意事项

  • 若需要包含时间的日期值,替换ds_nodash为ts_nodash(格式YYYYMMDDTHHMMSS),并同步调整格式字符串。
  • 确保max_pipeline_delay和frequency_hour为全局可访问变量,或在函数内部定义。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 00:25:35