aiohttp批量发送大POST请求:序列化阻塞与低延迟优化问询
优化大体积POST请求的发送延迟与EventLoop阻塞问题
一、绕过JSON序列化阻塞,实现请求就绪即发送
JSON序列化(尤其是5MB级别的大payload)会同步占用EventLoop,导致请求排队发送。解决核心是把序列化移出EventLoop线程,让每个请求的序列化完成后立即启动发送:
1. 用线程池异步完成序列化
借助concurrent.futures.ThreadPoolExecutor,将JSON序列化任务丢到后台线程执行,EventLoop可以同时处理请求的其他准备工作,序列化完成一个就安排一个请求发送,无需等待全部序列化完成:
import asyncio import json from concurrent.futures import ThreadPoolExecutor import aiohttp # 初始化线程池,根据CPU核心数调整worker数量 serializer_executor = ThreadPoolExecutor(max_workers=4) async def send_request(session, url, payload_bytes): # 直接传递序列化好的bytes,指定Content-Type async with session.post( url, data=payload_bytes, headers={"Content-Type": "application/json"} ) as resp: return await resp.text() async def main(): target_url = "https://your-api-target.com/endpoint" # 构造50个大体积payload payloads = [{"large_field": "x" * 5*1024*1024} for _ in range(50)] loop = asyncio.get_running_loop() # 并行执行序列化任务,不阻塞EventLoop serialized_tasks = [ loop.run_in_executor(serializer_executor, json.dumps, payload) for payload in payloads ] serialized_strings = await asyncio.gather(*serialized_tasks) # 转为bytes,方便直接发送 payload_bytes_list = [s.encode("utf-8") for s in serialized_strings] # 此时每个payload都已准备好,批量发送时会逐个就绪逐个发送 async with aiohttp.ClientSession() as session: send_tasks = [ send_request(session, target_url, payload_bytes) for payload_bytes in payload_bytes_list ] responses = await asyncio.gather(*send_tasks) print("所有请求完成") if __name__ == "__main__": asyncio.run(main())
2. 手动构造JSON bytes(适合固定结构场景)
如果你的payload结构固定,可以直接拼接JSON格式的字符串并转为bytes,完全跳过json.dumps的开销。比如:
# 假设payload结构固定为{"data": <大字符串>} payload_bytes = b'{"data": "' + large_string.encode("utf-8") + b'"}'
注意:这种方式需要自己处理转义(比如双引号、特殊字符),仅适合结构简单且可控的场景。
二、优化aiohttp的EventLoop等待延迟
针对你提到的connector等待EventLoop的环节,可从以下几个方向加速:
1. 替换为高性能EventLoop
默认的asyncio EventLoop性能有限,改用uvloop(基于libuv的C实现)可大幅提升调度效率,降低延迟:
import uvloop import asyncio # 在程序启动时设置EventLoop策略 asyncio.set_event_loop_policy(uvloop.EventLoopPolicy())
2. 调优aiohttp Connector参数
- 增大并发连接限制:默认
TCPConnector(limit=100),如果服务器允许更高并发,可适当调大(比如limit=200),减少请求排队等待连接的时间; - 开启连接复用:保持
force_close=False(默认),复用TCP连接,避免重复握手的开销; - 优化DNS解析:使用
AsyncResolver异步解析DNS,避免阻塞:connector = aiohttp.TCPConnector(resolver=aiohttp.AsyncResolver())
3. 控制并发请求数量
一次性提交50个请求可能导致EventLoop调度压力过大,用asyncio.Semaphore控制并发数,让请求分批有序执行,反而能降低整体延迟:
async def send_request_with_semaphore(session, url, payload_bytes, semaphore): async with semaphore: return await send_request(session, url, payload_bytes) # 在main函数中初始化信号量,比如限制20个并发 semaphore = asyncio.Semaphore(20) send_tasks = [ send_request_with_semaphore(session, target_url, payload_bytes, semaphore) for payload_bytes in payload_bytes_list ]
4. 替代方案
如果aiohttp的延迟仍无法满足要求,可尝试:
- httpx异步客户端:底层实现更轻量化,调度效率更高,API与aiohttp类似;
- trio/curio框架:替代asyncio的异步框架,调度模型更适合低延迟场景;
- 直接使用socket:极端场景下,可基于asyncio直接封装TCP请求,跳过aiohttp的上层封装开销,但开发成本较高。
内容的提问来源于stack exchange,提问作者Bala
相关产品推荐
相关产品推荐

