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

Airflow中PythonSensor返回的PokeReturnValue无法通过XCom存储/读取

问题定位与解决方案

核心问题

PythonSensor 默认不会自动推送任意返回值到XCom,需满足两个前提才会生成XCom记录:

  • 返回的 PokeReturnValue 必须包含 data 字段
  • 传感器的 do_xcom_push 参数需保持为 True(默认值,但若被手动修改会失效)

另外,若传感器轮询多次,仅当最终返回 is_done=True 时,才会触发XCom推送(除非显式手动推送)。

具体修复步骤

1. 正确构造返回的PokeReturnValue

在第一个PythonSensor的poke函数中,必须将需要传递的数据放入PokeReturnValue的data字段:

from airflow.sensors.python import PokeReturnValue
import random

def first_sensor_poke(**context):
    random_list = [random.randint(1, 100) for _ in range(5)]
    success_flag = random.choice([True, False])
    # 日志记录生成值
    context['logger'].info(f"生成随机列表: {random_list}, 成功标记: {success_flag}")
    # 正确返回带数据的PokeReturnValue
    return PokeReturnValue(
        is_done=success_flag,
        data={"random_list": random_list, "success": success_flag}
    )

2. 确保传感器开启XCom推送

定义第一个PythonSensor时,显式指定do_xcom_push=True(避免默认配置被修改):

from airflow.sensors.python import PythonSensor

first_sensor = PythonSensor(
    task_id="first_python_sensor",
    poke_interval=30,
    timeout=300,
    python_callable=first_sensor_poke,
    do_xcom_push=True,  # 显式开启XCom推送
    dag=dag
)

3. 后两个任务正确读取XCom

在第二个PythonSensor和PythonOperator的函数中,通过task_instance读取指定task_id的XCom数据:

def read_xcom_function(**context):
    # 读取第一个传感器的XCom数据
    xcom_data = context['task_instance'].xcom_pull(
        task_ids="first_python_sensor",
        key="return_value"  # PokeReturnValue的data会被存在return_value键下
    )
    if xcom_data:
        context['logger'].info(f"读取到XCom数据: {xcom_data}")
    else:
        context['logger'].warning("未读取到XCom数据")

4. 可选:显式手动推送XCom(更可控)

如果需要无论传感器成功与否都推送数据,可在poke函数中直接调用xcom_push,不受is_done状态限制:

def first_sensor_poke(**context):
    random_list = [random.randint(1, 100) for _ in range(5)]
    success_flag = random.choice([True, False])
    context['logger'].info(f"生成随机列表: {random_list}, 成功标记: {success_flag}")
    # 显式推送XCom,指定自定义key
    context['task_instance'].xcom_push(
        key="sensor_data",
        value={"random_list": random_list, "success": success_flag}
    )
    return PokeReturnValue(is_done=success_flag)

此时读取时需指定对应的key:

xcom_data = context['task_instance'].xcom_pull(task_ids="first_python_sensor", key="sensor_data")

验证方法

执行DAG后,可通过以下方式确认XCom是否生成:

  • 在Airflow UI的任务详情页,切换到XCom标签查看记录
  • 直接查询Airflow数据库的xcom表,筛选对应dag_id和task_id的记录

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 12:45:35