Airflow:如何将装饰器任务返回的数据传递给SimpleHttpOperator?
问题根源
你的报错由两个底层原因共同导致:
@task装饰器的multiple_outputs=True参数会将返回的dict自动拆分为多个独立XCom条目(每个键对应一条XCom),不会存储完整dict到XCom默认key,导致后续算子拉取的请求负载不符合预期SimpleHttpOperator对data参数的处理逻辑为:如果传入对象是dict,会自动编码为表单格式,哪怕你设置了JSON格式的请求头也不生效;同时直接传入Taskflow返回的XComArg对象时,模板渲染会生成对象的默认字符串标识,而非实际的负载内容,最终发往接口的请求体为非法JSON,触发400错误
解决方案
方案1:无需自定义算子,直接适配现有逻辑
首先去掉@task的multiple_outputs=True参数,再通过Jinja模板的tojson过滤器将拉取到的XCom dict自动序列化为JSON字符串,修改后代码如下:
from airflow.decorators import dag, task from airflow.providers.http.operators.http import SimpleHttpOperator import json from datetime import datetime default_args = { "owner": "airflow", "start_date": datetime(2021, 1, 1), } @dag(default_args=default_args, schedule_interval=None, tags=["Http Operators"]) def http_operator(): @task() def Get_payload(**kwargs): # STEP 1: 从数据库拉取数据 # STEP 2: 加工处理数据 # STEP 3: 返回请求负载 data = { "key_1": "Value 1", "key_2": "Value 2", "key_3": "Value 3", "key_4": "Value 4", } return data get_data = Get_payload() ml_api = SimpleHttpOperator( task_id="some_api", http_conn_id="http_conn_id", method="POST", endpoint="/some-path", # 用模板拉取前序任务输出,通过tojson过滤器转为JSON字符串 data="{{ ti.xcom_pull(task_ids='Get_payload') | tojson }}", headers={"Content-Type": "application/json"}, ) http_operator_dag = http_operator()
依赖关系会被Airflow自动识别,无需手动声明get_data >> ml_api。
方案2:完全贴合Taskflow使用习惯
如果不想写模板语法,可以在任务中直接返回序列化后的JSON字符串,直接传入SimpleHttpOperator即可:
@task() def Get_payload(**kwargs): data = { "key_1": "Value 1", "key_2": "Value 2", "key_3": "Value 3", "key_4": "Value 4", } # 直接返回JSON字符串 return json.dumps(data) get_data = Get_payload() ml_api = SimpleHttpOperator( task_id="some_api", http_conn_id="http_conn_id", method="POST", endpoint="/some-path", data=get_data, headers={"Content-Type": "application/json"}, )
内容的提问来源于stack exchange,提问作者Dheemanth Bhat
相关产品推荐
相关产品推荐

