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无效)
如果服务端强制断开空闲连接,就做更智能的重连,避免重复处理数据,同时减少频繁重试:
- 记录最后处理成功的事件ID(如果API支持传递
last_event_id这类参数),重连时带上,跳过已处理的数据。 - 采用指数退避重试,避免短时间内多次重试影响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.
相关产品推荐
相关产品推荐

