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

如何在Airflow BigQuery Operator中动态指定SQL文件并执行(GCP)

解决Airflow动态选择SQL文件给BigQueryOperator的问题

核心问题分析

你当前的代码存在三个关键问题:

  1. 全局变量SQL_to_run在DAG解析阶段就被赋值为空字符串,PythonOperator运行时修改的是函数内的局部变量,无法影响BigQueryOperator的参数。
  2. condition_check函数未正确获取执行日期的星期信息,且错误使用赋值运算符=而非比较运算符==。
  3. 未设置任务间的依赖关系,也未正确通过XCom传递动态值。

解决方案:用XCom传递动态SQL文件名

以下是修正后的完整代码,包含XCom推送/拉取、正确的日期判断和任务依赖逻辑:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.google.cloud.operators.bigquery import BigQueryOperator
from airflow.operators.short_circuit import ShortCircuitOperator
from datetime import datetime

# Airflow变量
SQL1 = 'SQL1.sql'
SQL2 = 'SQL2.sql'
project_name = "your-project-id"
dataset = "your-dataset"
default_args = {
    'owner': 'airflow',
    # 补充你的其他默认参数
}

def condition_check(**context):
    # 从上下文获取执行日期
    execution_date = context['execution_date']
    # 获取英文星期名(如'Sunday'、'Monday')
    weekday = execution_date.strftime("%A")
    
    if weekday == 'Sunday':
        sql_file = SQL1
    elif weekday == 'Monday':
        sql_file = SQL2
    else:
        sql_file = None
    
    # 将选中的SQL文件名推送到XCom
    context['ti'].xcom_push(key='sql_to_run', value=sql_file)
    return sql_file

def should_run(**context):
    # 拉取XCom值,判断是否需要执行后续SQL任务
    sql_file = context['ti'].xcom_pull(task_ids='type_of_extract', key='sql_to_run')
    return sql_file is not None

with DAG(dag_id="DAG_NM",
        start_date=datetime(2024,2,21),
        template_searchpath=['/usr/local/airflow/dags/resources/'],
        schedule_interval="0 10 * * 7#1,7#2",
        default_args=default_args,
        catchup=False) as dag:

    type_of_sql_task = PythonOperator(
        task_id="type_of_extract",
        python_callable=condition_check,
        provide_context=True,
    )

    check_run_condition = ShortCircuitOperator(
        task_id="check_run_condition",
        python_callable=should_run,
        provide_context=True,
    )

    execute_sql_data_task = BigQueryOperator(
        task_id="Execute_data",
        # 使用Jinja模板拉取XCom中的动态SQL文件名
        sql="{{ ti.xcom_pull(task_ids='type_of_extract', key='sql_to_run') }}",
        params={'projectid': project_name, 'dataset': dataset},
        use_legacy_sql=False,
    )

    # 设置任务执行顺序
    type_of_sql_task >> check_run_condition >> execute_sql_data_task

关键说明

  • XCom传递逻辑:condition_check函数通过context['ti'].xcom_push将选中的SQL文件名存入XCom,BigQueryOperator用Jinja模板{{ ti.xcom_pull(...) }}读取该值,实现动态参数传递。
  • 日期判断:通过execution_date.strftime("%A")获取执行日期的英文星期名,确保条件判断准确。
  • 任务跳过机制:ShortCircuitOperator会在无符合条件的SQL文件时,跳过后续的BigQuery任务,避免无效执行。
  • 依赖关系:明确任务执行顺序,确保先完成SQL文件的判断,再进行条件检查和SQL执行。

内容的提问来源于stack exchange,提问作者Mani Shankar.S

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 21:52:43