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

如何在SimpleHttpOperator响应函数中访问Task Instance实现XCom推送?

解决SimpleHttpOperator响应函数中推送指定键XCom的问题

核心解决方案

要在SimpleHttpOperator的响应处理函数中访问task_instance并推送自定义键的XCom,需要两步关键操作:

  1. 给SimpleHttpOperator添加provide_context=True参数,允许响应函数接收Airflow上下文。
  2. 修改响应处理函数,通过**kwargs获取上下文里的task_instance和所需的dag_run配置。

修改后的代码示例

响应处理函数

def _handler_object_result(response, **kwargs):
    # 从上下文获取文件路径
    file = kwargs['dag_run'].conf['file']
    header_result = response.json()
    task_id = header_result["task"]["id"]

    api = "https://redacted.com/api/task/result/{task_id}".format(task_id=task_id)
    resp = requests.get(api, verify=False).json()
    data = json.loads(resp["data"])
    file_object = json.dumps(data["OBJECT"])
    # 将哈希值转为字符串作为XCom键,避免负数等异常情况
    file_hash = str(hash(file))
    
    # 从上下文获取task_instance并推送XCom
    ti = kwargs['task_instance']
    ti.xcom_push(key=file_hash, value=file_object)
    
    # 验证推送结果
    return ti.xcom_pull(key=file_hash) is not None

SimpleHttpOperator配置

object_result = SimpleHttpOperator(
    task_id="object_result",
    method='POST',
    data=json.dumps({"file": "{{ dag_run.conf['file'] }}", "keyword": "object"}),
    http_conn_id="coma_api",
    endpoint="/api/v1/file/describe",
    headers={"Content-Type": "application/json"},
    extra_options={"verify":False},
    response_check=_handler_object_result,  # 直接绑定处理函数,无需lambda
    do_xcom_push=False,
    provide_context=True,  # 开启上下文传递
    dag=dag,
)

关键说明

  • 上下文传递:SimpleHttpOperator默认不会把Airflow上下文(如task_instance、dag_run)传给response_check函数,必须显式设置provide_context=True才能启用。
  • 模板渲染问题:之前lambda中直接传递的"{{ dag_run.conf['file'] }}"不会被Airflow渲染,因为lambda不属于Airflow的模板字段,正确方式是从上下文的dag_run中直接提取配置值。
  • XCom键的安全性:hash(file)返回的是整数,可能为负数,转为字符串作为XCom键可以避免潜在的键名异常问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 05:45:43