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

Airflow中SQLExecuteQueryOperator执行Starburst查询结果XCom获取失败

解决SQLExecuteQueryOperator无法向XCom推送Starburst查询结果的问题

你的问题核心是SQLExecuteQueryOperator的参数配置导致XCom推送失败,结合Starburst查询结果的结构,按以下步骤修改即可解决:

问题原因

  1. split_statement=True参数会触发SQL语句拆分逻辑,即使是单条SELECT语句,也会进入批量执行分支,导致返回值未被正确推送到XCom。
  2. 未正确解析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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 01:23:11