在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
相关产品推荐
相关产品推荐

