如何在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
相关产品推荐
相关产品推荐

