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

Airflow自定义插件:SQL执行结果日志可见但未写入XCOM求助

Airflow XCom未推送SQL执行结果问题解决

问题现象

自定义SqlExecuteOperator执行SQL后,日志中能看到Record: (1,),但XCom中没有对应结果,Python任务日志显示Done. Returned value was: None。

问题原因

  1. Python任务未返回Operator执行结果:在@task装饰的函数中调用sql_task_op.execute()后,没有将结果作为函数返回值返回,导致Python任务本身返回None,XCom记录的是Python任务的返回值。
  2. 传入空Context:调用execute(context=dict())时传入空字典,Operator无法获取task_instance等关键信息,无法正常触发XCom推送逻辑。
  3. Operator的do_xcom_push被上层覆盖:在Python任务内部调用Operator时,do_xcom_push参数的作用会被Python任务的XCom逻辑覆盖,实际生效的是Python任务的返回值。

解决方案

方案一:直接将自定义Operator作为独立任务(推荐)

无需用@task包裹,直接把SqlExecuteOperator作为DAG的任务节点,由Airflow直接调度:

sql_task_op = SqlExecuteOperator(
    task_id="sql_task",
    conn_id=SQL_CONNECTION,
    sql="SELECT 1",
    database='TEST',
    do_xcom_push=True,
)

# 配置DAG任务依赖示例
# sql_task_op >> downstream_task

此方式下,Operator的execute方法返回值会自动通过do_xcom_push推送到XCom,对应任务ID为sql_task。

方案二:在Python任务内调用Operator(保留@task场景)

修改Python任务函数,传入完整Context并返回执行结果:

@task
def get_results_from_sql(**kwargs):
    sql_task_op = SqlExecuteOperator(
        task_id="sql_task",
        conn_id=SQL_CONNECTION,
        sql="SELECT 1",
        database='TEST',
        do_xcom_push=True,
    )
    # 传入完整的kwargs作为执行上下文
    result = sql_task_op.execute(context=kwargs)
    # 返回结果,让Python任务将其推送到XCom
    return result

此时XCom会记录get_results_from_sql任务的返回值,日志会显示Done. Returned value was: 1。

验证

执行完成后,在Airflow UI的XCom页面,对应任务ID下可查看推送的结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 12:50:55