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

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接口的场景)

你需要调整请求配置和相关服务的超时设置:

  1. 给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
  1. 调整关联服务的超时配置:
  • 进入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请求的同步调用本身存在很多不可控风险,比如网络波动、临时中断都会导致任务失败,建议改为异步架构:

  1. 先改造Cloud Run接口:
    • 新增触发接口:收到请求后立刻返回唯一任务ID,后台异步执行任务
    • 新增状态查询接口:传入任务ID可以查询当前任务的执行状态(运行中/成功/失败)和结果
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 12:09:04