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

aiohttp实现无限流传输遇长等待报错问题求助

解决aiohttp流式连接10分钟断开的问题

原因分析

这个错误大概率是中间件(负载均衡、代理)或者服务器在连接空闲10分钟后主动断开TCP连接,aiohttp检测到连接中断后抛出ClientPayloadError。

方案1:开启TCP Keepalive(最直接)

通过设置socket的TCP保活参数,让底层定期发送探测包,维持连接活跃,避免被中间件判定为空闲断开。

需要导入socket模块,修改请求代码:

import socket

async with self.client_session.get(
    self._url(path),
    headers=self._headers,
    timeout=aiohttp.ClientTimeout(total=0, connect=0, sock_connect=0, sock_read=0),
    socket_options=[
        (socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1),
        (socket.IPPROTO_TCP, socket.TCP_KEEPIDLE, 300),  # 空闲5分钟后开始发送保活探测
        (socket.IPPROTO_TCP, socket.TCP_KEEPINTVL, 60),  # 每隔1分钟发送一次探测
        (socket.IPPROTO_TCP, socket.TCP_KEEPCNT, 5),      # 连续5次探测无响应则断开连接
    ]
) as response:
    async for message in response.content:
        # 处理消息逻辑
        pass

方案2:优化重连机制(如果Keepalive无效)

如果服务端强制断开空闲连接,就做更智能的重连,避免重复处理数据,同时减少频繁重试:

  1. 记录最后处理成功的事件ID(如果API支持传递last_event_id这类参数),重连时带上,跳过已处理的数据。
  2. 采用指数退避重试,避免短时间内多次重试影响API服务。

示例代码:

import asyncio
from functools import wraps

def retry_with_backoff(max_retries=5, initial_delay=1):
    def decorator(func):
        @wraps(func)
        async def wrapper(*args, **kwargs):
            retries = 0
            delay = initial_delay
            last_event_id = kwargs.get('last_event_id')
            while retries < max_retries:
                try:
                    # 如果有last_event_id,传递给API请求
                    if last_event_id:
                        kwargs['params'] = kwargs.get('params', {})
                        kwargs['params']['last_event_id'] = last_event_id
                    # 调用原连接处理函数
                    last_event_id = await func(*args, **kwargs)
                    delay = initial_delay  # 成功后重置延迟
                    retries = 0
                except aiohttp.ClientPayloadError:
                    retries += 1
                    if retries >= max_retries:
                        raise
                    await asyncio.sleep(delay)
                    delay = min(delay * 2, 30)  # 指数退避,最大延迟30秒
            return last_event_id
        return wrapper
    return decorator

# 定义你的连接处理函数
@retry_with_backoff()
async def stream_events(self, last_event_id=None):
    async with self.client_session.get(
        self._url(path),
        headers=self._headers,
        timeout=aiohttp.ClientTimeout(total=0, connect=0, sock_connect=0, sock_read=0),
    ) as response:
        async for message in response.content:
            # 处理消息,更新last_event_id
            event = parse_message(message)
            last_event_id = event['id']
            # 其他处理逻辑
        return last_event_id

方案3:确认API的心跳机制

检查API文档,看是否支持服务器主动发送心跳包(比如空分块、特定的ping消息),如果支持,确保客户端能正确识别并忽略这些心跳,维持连接。如果API没有这个机制,可联系服务端开启。

内容的提问来源于stack exchange,提问作者Daniel Mühlbachler-P.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 04:00:23