如何在SimpleHttpOperator响应函数中访问Task Instance实现XCom推送?
解决SimpleHttpOperator响应函数中推送指定键XCom的问题
核心解决方案
要在SimpleHttpOperator的响应处理函数中访问task_instance并推送自定义键的XCom,需要两步关键操作:
- 给
SimpleHttpOperator添加provide_context=True参数,允许响应函数接收Airflow上下文。 - 修改响应处理函数,通过
**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
相关产品推荐
相关产品推荐

