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

如何在AirFlow下游任务中读取HTTP状态码

如何在AirFlow的SimpleHttpOperator中获取HTTP响应状态码?

问题场景

以下是我编写的AirFlow DAG代码,目前check_req分支任务只能获取到HTTP请求的响应体,无法拿到响应状态码,需要修改代码以实现状态码的检查:

@task.branch(task_id="check_req")
def check_req(**kwargs):
    response = kwargs['ti'].xcom_pull(task_ids='get_req')
    print(f"Response data: {response}")
    if response == 200:
        return "good"
    else:
        return "bad"

get_req = SimpleHttpOperator(
    task_id='get_req',
    http_conn_id="http_conn",
    endpoint='endpoint/resource',
    method='GET',
    headers={"Content-Type": "application/json"},
)

check_req_op = check_req()
begin = EmptyOperator(task_id="begin")
end = EmptyOperator(task_id="end", trigger_rule="none_failed_min_one_success")
good = EmptyOperator(task_id="good")
bad = EmptyOperator(task_id="bad")

begin >> get_req >> check_req_op
check_req_op >> good >> end
check_req_op >> bad >> end

解决方案

核心原理

SimpleHttpOperator默认仅将响应体推送到XCom,我们可以通过配置response_filter参数,自定义返回的内容——该参数接收requests.Response对象,通过它可以获取状态码、响应体等所有响应信息。

修改后的代码

方式1:同时获取状态码和响应体

修改get_req任务的定义,添加response_filter让它返回包含状态码和响应体的字典:

get_req = SimpleHttpOperator(
    task_id='get_req',
    http_conn_id="http_conn",
    endpoint='endpoint/resource',
    method='GET',
    headers={"Content-Type": "application/json"},
    # 自定义返回内容,包含状态码和响应体
    response_filter=lambda response: {
        "status_code": response.status_code,
        "content": response.json()  # 响应为JSON格式时用json(),非JSON用text
    }
)

然后修改check_req任务,从XCom中取出字典并判断状态码:

@task.branch(task_id="check_req")
def check_req(**kwargs):
    response_data = kwargs['ti'].xcom_pull(task_ids='get_req')
    status_code = response_data['status_code']
    content = response_data['content']
    print(f"Response status code: {status_code}")
    print(f"Response content: {content}")
    if status_code == 200:
        return "good"
    else:
        return "bad"

方式2:仅获取状态码

如果只需要状态码,可简化response_filter直接返回状态码:

get_req = SimpleHttpOperator(
    task_id='get_req',
    http_conn_id="http_conn",
    endpoint='endpoint/resource',
    method='GET',
    headers={"Content-Type": "application/json"},
    response_filter=lambda response: response.status_code
)

对应的check_req任务可以直接判断:

@task.branch(task_id="check_req")
def check_req(**kwargs):
    status_code = kwargs['ti'].xcom_pull(task_ids='get_req')
    print(f"Response status code: {status_code}")
    if status_code == 200:
        return "good"
    else:
        return "bad"

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 12:30:27