Airflow中SQLExecuteQueryOperator执行Starburst查询结果XCom获取失败
解决SQLExecuteQueryOperator无法向XCom推送Starburst查询结果的问题
你的问题核心是SQLExecuteQueryOperator的参数配置导致XCom推送失败,结合Starburst查询结果的结构,按以下步骤修改即可解决:
问题原因
split_statement=True参数会触发SQL语句拆分逻辑,即使是单条SELECT语句,也会进入批量执行分支,导致返回值未被正确推送到XCom。- 未正确解析Starburst返回的嵌套结果结构,即便XCom推送成功,也可能无法拿到预期的0/1数值。
修改方案
1. 调整SQLExecuteQueryOperator参数
关闭split_statement参数(默认值为False,直接移除该配置即可),确保单条SELECT语句的返回值被正确捕获并推送。
2. 完善Python函数解析结果
Starburst返回的COUNT(*)结果是嵌套列表格式(比如[[0]]),需要提取内层的数值才能使用。
修改后的完整代码
import pendulum, os, logging from airflow import DAG from airflow.operators.python import PythonOperator from airflow.models import TaskInstance from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator from datetime import datetime def catch_output(ti, **kwargs): query_result = ti.xcom_pull(task_ids="query_task") # 解析Starburst返回的嵌套结果,提取COUNT(*)的数值 count_value = query_result[0][0] if query_result else None logging.info(f"查询到的记录数: {count_value}") # 可添加业务判断逻辑 if count_value == 1: logging.info("表中存在目标记录") elif count_value == 0: logging.info("表中不存在目标记录") with DAG( 'dag_name', schedule=None, start_date=pendulum.datetime(2025,2,25, tz='UTC'), max_active_runs=1, catchup=False, default_args={...} ) as dag: query_task = SQLExecuteQueryOperator( task_id='query_task', conn_id='xyz', show_return_value_in_logs=True, # 移除split_statement=True,使用默认的False配置 sql='SELECT COUNT(*) FROM TABLE', return_last=True, do_xcom_push=True ) catch_output = PythonOperator( task_id='catch_output', python_callable=catch_output, dag=dag ) query_task >> catch_output
验证要点
- 执行
query_task后,查看任务日志是否显示返回值(因开启show_return_value_in_logs=True) - 执行
catch_output后,查看日志是否输出解析后的count_value,确认是否为预期的0或1
内容的提问来源于stack exchange,提问作者ennezetaqu
相关产品推荐
相关产品推荐

