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

在Airflow中使用BigQueryInsertJobOperator执行多语句BigQuery任务遇阻

解决Airflow中BigQueryInsertJobOperator执行IF ELSE删插脚本的问题

你的代码存在两处关键问题:一是BigQueryInsertJobOperator的查询配置格式错误,二是SQL脚本有语法疏漏,以下是具体修正方案:

1. 修正SQL脚本语法

确保每个独立SQL语句结尾添加分号(比如IF块内的DELETE语句),否则BigQuery会解析失败:

BQ_INSERT_QUERY = '''
    DECLARE current_year STRING;
    DECLARE past_year STRING;
    
    SET (current_year, past_year) = (SELECT AS STRUCT
         CAST(EXTRACT(YEAR FROM CURRENT_DATE()) AS STRING) AS current_year,
         CAST(EXTRACT(YEAR FROM DATE_SUB(CURRENT_DATE(), INTERVAL 1 YEAR)) AS STRING) AS past_year);
    
    IF EXTRACT(MONTH FROM CURRENT_DATE()) IN (1, 2, 3) THEN
        DELETE FROM `data-pipeline.MAIN.Main_Table`
        WHERE CAST(column1 AS STRING) LIKE CONCAT('%', current_year);
        
        INSERT INTO `data-pipeline.MAIN.Main_Table` (column1,column2,column3)
        SELECT column1,column2,column3
        FROM `data-pipeline.MAIN.Hold_Table`
        WHERE CAST(column1 AS STRING) LIKE CONCAT('%', current_year);
    ELSE
        DELETE FROM `data-pipeline.MAIN.Main_Table`
        WHERE CAST(column1 AS STRING) LIKE CONCAT('%', past_year);
        
        INSERT INTO `data-pipeline.MAIN.Main_Table` (column1,column2,column3)
        SELECT column1,column2,column3
        FROM `data-pipeline.MAIN.Hold_Table`
        WHERE CAST(column1 AS STRING) LIKE CONCAT('%', past_year);
    END IF; 
'''

2. 修正BigQueryInsertJobOperator配置

configuration.query.query字段需要传入字符串而非列表,原代码的[BQ_INSERT_QUERY]属于错误写法,修正后的Operator代码:

insert_target_table = BigQueryInsertJobOperator(
    task_id='daily_insert_target_table',
    project_id='data-pipeline',
    location='asia-southeast1',
    configuration={
        'query': {
            'query': BQ_INSERT_QUERY,
            'useLegacySql': False,
            # 可选:如需支持大结果集可添加以下参数
            'allowLargeResults': True
        }
    },
    dag=dag
)

关键说明

  • BigQuery支持多语句脚本,只要useLegacySql设为False,就能正确解析IF ELSE逻辑和多语句操作
  • 必须保证每个SQL语句(DELETE、INSERT等)结尾都带分号,避免语法解析错误
  • 查询内容直接以字符串传入即可,无需包装成列表

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 04:01:15