Airflow(Cloud Composer)触发长耗时Cloud Run任务无法接收响应如何解决
问题根因
你遇到的问题本质是长耗时同步HTTP请求的TCP空闲连接被中间网络节点断开:你设置的timeout=(5, 3600)是HTTP请求的总超时限制,但Cloud Composer到Cloud Run的链路中,包括VPC网关、GCP边缘节点、Airflow Worker的TCP栈都有默认的TCP空闲超时机制,通常在30秒到5分钟不等,如果连接在空闲时间内没有任何数据包传输,就会被主动断开,导致Airflow端收不到Cloud Run最终返回的响应。
短耗时任务因为在连接断开前就已经返回结果,所以不会触发这个问题。
解决方案
方案1:同步调用优化(快速修复,适合不想改Cloud Run接口的场景)
你需要调整请求配置和相关服务的超时设置:
- 给requests请求添加TCP keepalive配置,定期发送探测包维持连接不被断开,修改后的函数代码如下:
import requests from requests.adapters import HTTPAdapter import socket from datetime import timedelta class TCPKeepAliveAdapter(HTTPAdapter): def init_poolmanager(self, *args, **kwargs): # 适配Cloud Composer Worker的Linux环境配置TCP保活参数 socket_options = [ (socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1), # 连接空闲30秒后开始发送保活探测包 (socket.IPPROTO_TCP, socket.TCP_KEEPIDLE, 30), # 每10秒发送一次探测包 (socket.IPPROTO_TCP, socket.TCP_KEEPINTVL, 10), # 连续9次探测失败才判定连接断开 (socket.IPPROTO_TCP, socket.TCP_KEEPCNT, 9), ] kwargs['socket_options'] = socket_options return super().init_poolmanager(*args, **kwargs) def make_authorized_get_request(service_url, **kwargs): auth_req = google.auth.transport.requests.Request() id_token = google.oauth2.id_token.fetch_id_token(auth_req, service_url) headers = {"Authorization": f"Bearer {id_token}"} # 使用配置了保活机制的会话发送请求 session = requests.Session() session.mount('https://', TCPKeepAliveAdapter()) req = session.get(service_url, headers=headers, timeout=(5, 3600)) status = req.status_code response_json = req.json() return response_json, status
- 调整关联服务的超时配置:
- 进入Cloud Run服务配置页,将请求超时调整到大于你的任务最大耗时,Cloud Run最高支持60分钟的请求超时。
- 给PythonOperator添加
execution_timeout参数,避免Airflow Worker主动终止长耗时任务:
big3_request = PythonOperator( task_id='big3_request', python_callable=make_authorized_get_request, op_kwargs={"service_url":"cloud run url"}, execution_timeout=timedelta(minutes=40) # 按需调整,要大于任务最大耗时 )
方案2:异步触发+轮询(更推荐,稳定性更高)
长耗时HTTP请求的同步调用本身存在很多不可控风险,比如网络波动、临时中断都会导致任务失败,建议改为异步架构:
- 先改造Cloud Run接口:
- 新增触发接口:收到请求后立刻返回唯一任务ID,后台异步执行任务
- 新增状态查询接口:传入任务ID可以查询当前任务的执行状态(运行中/成功/失败)和结果
- Airflow侧拆分为触发+轮询两个任务,避免维持长连接:
from airflow.sensors.python import PythonSensor from datetime import timedelta def trigger_cloud_run(service_url, **kwargs): auth_req = google.auth.transport.requests.Request() id_token = google.oauth2.id_token.fetch_id_token(auth_req, service_url) headers = {"Authorization": f"Bearer {id_token}"} resp = requests.post(f"{service_url}/trigger", headers=headers, timeout=10) task_id = resp.json()["task_id"] # 将任务ID存入XCom供轮询任务使用 kwargs["ti"].xcom_push(key="cloud_run_task_id", value=task_id) def check_run_status(service_url, **kwargs): task_id = kwargs["ti"].xcom_pull(key="cloud_run_task_id", task_ids="trigger_big3_task") auth_req = google.auth.transport.requests.Request() id_token = google.oauth2.id_token.fetch_id_token(auth_req, service_url) headers = {"Authorization": f"Bearer {id_token}"} resp = requests.get(f"{service_url}/check_status?task_id={task_id}", headers=headers, timeout=10) result = resp.json() if result["status"] == "success": # 任务成功,存入结果到XCom kwargs["ti"].xcom_push(key="run_result", value=result["data"]) return True elif result["status"] == "failed": raise Exception(f"Cloud Run任务执行失败:{result.get('error', '未知错误')}") # 任务仍在运行,返回False触发下一次轮询 return False # 触发任务 trigger_task = PythonOperator( task_id='trigger_big3_task', python_callable=trigger_cloud_run, op_kwargs={"service_url":"你的Cloud Run服务地址"}, execution_timeout=timedelta(minutes=1) ) # 轮询任务:每30秒查询一次状态,最多等待40分钟 poll_task = PythonSensor( task_id='poll_big3_status', python_callable=check_run_status, op_kwargs={"service_url":"你的Cloud Run服务地址"}, poke_interval=30, timeout=timedelta(minutes=40).total_seconds(), mode='reschedule' # 轮询间隔期间释放Worker资源,降低资源消耗 ) trigger_task >> poll_task
其他说明
不需要更换其他Operator,两种方案都可以基于你当前使用的PythonOperator实现,优先推荐方案2,稳定性更强,也支持更长耗时的任务场景。
内容的提问来源于stack exchange,提问作者steak_ale_piethon
相关产品推荐
相关产品推荐

