【Bug已解决】[AsyncGRPO] aiohttp limits to 100 reqs when max_inflight_tasks > 100 解决方案
【Bug已解决】[AsyncGRPO] aiohttp limits to 100 reqs when max_inflight_tasks > 100 解决方案
原始报错:[AsyncGRPO] aiohttp limits to 100 reqs when max_inflight_tasks > 100 场景:异步 GRPO 用 aiohttp 作为 HTTP 客户端,向推理引擎并发发请求(rollout 生成)。配置里
max_inflight_tasks设到 200,期望同时跑 200 个并发任务,但实际同时只有约 100 个在飞,其余排队等待。根因是 aiohttp 的ClientSession默认连接池上限就是 100(TCPConnector(limit=100)),并发超过 100 就被池子卡住,第 101 个得等前面的释放。要支持更多 inflight,必须显式调高连接池上限。 关键词:aiohttp、连接池上限、TCPConnector、limit、concurrency、max_inflight_tasks、信号量、异步并发、100 默认限制。
一、现象长什么样
并发上不去,卡在 100:
- 配置
max_inflight_tasks = 200,期望 200 个任务同时向引擎发请求; - 实际观测并发数稳定在 ~100,多的任务在排队,吞吐上不去;
- 用
asyncio跑,理论上能并发 200,但 aiohttp 这层卡住; - 调大
max_inflight_tasks没用,瓶颈不在任务调度而在 HTTP 连接池; - 改小到 100 以内一切正常,一过 100 就出现"第 101 个起开始等";
- 日志/报错没有明显的"连接超限",只是默默排队——典型隐性限制。
核心问题:aiohttp 默认连接池上限 100,成了并发天花板,与max_inflight_tasks脱节。
二、背景:aiohttp 的默认连接池上限为什么是 100
aiohttp 的ClientSession内部有一个TCPConnector管理连接复用,它有两个默认上限:
limit:整个连接池的总并发连接数,默认100;limit_per_host:单个 host 的并发连接数,默认也是100(实际上是limit的值,旧版默认 100)。
当你的 200 个任务都打同一个推理引擎 host,并发连接请求超过limit/limit_per_host,多出来的请求会在连接器里排队,等已有连接释放。所以即使你用asyncio起了 200 个协程,真正同时建连的只有 100 个。
这不是 bug,是默认值保守——避免无节制建连打爆服务端。但当你明确要跑max_inflight_tasks=200,就必须把连接池上限调到匹配值,否则max_inflight_tasks形同虚设。
三、根因:连接池上限与 max_inflight_tasks 脱节
根因拆解:
- 默认 100:
ClientSession()没传connector,用默认limit=100; - 未联动配置:
max_inflight_tasks调大,但TCPConnector.limit没跟着调; - 单 host 上限:全部请求打同一引擎 host,
limit_per_host也卡在 100; - 隐性排队:超 100 的请求静默排队,无报错,难发现;
- 吞吐天花板:并发被池子限死,加任务不加速;
- 误调任务数:开发者调大
max_inflight_tasks以为够了,瓶颈在更底层。
下面用最小模型复现"并发被池子上限卡在 100",再给"调高连接池上限"的修复。
四、最小可运行复现
import asyncio class FakeConnector: def __init__(self, limit=100): self.limit = limit self._sem = asyncio.Semaphore(limit) # 模拟连接池上限 async def request(self, i): async with self._sem: # 同时最多 limit 个 await asyncio.sleep(0.01) return f"resp{i}" async def run(connector, n): # 错:连接池默认 100,n=200 时只有 100 同时飞 tasks = [asyncio.create_task(connector.request(i)) for i in range(n)] results = await asyncio.gather(*tasks) return results if __name__ == "__main__": # 用信号量计数观测真实并发峰值 peak = {"v": 0, "cur": 0} async def instrumented(connector, n): async def one(i): peak["cur"] += 1 peak["v"] = max(peak["v"], peak["cur"]) try: return await connector.request(i) finally: peak["cur"] -= 1 await asyncio.gather(*[one(i) for i in range(n)]) asyncio.run(instrumented(FakeConnector(limit=100), 200)) print("默认 limit=100 时真实并发峰值:", peak["v"]) # 100,被卡运行可见并发峰值停在 100,正是 aiohttp 默认上限的现场。
五、方案:显式调高 TCPConnector 的 limit 与 limit_per_host
第一层:创建ClientSession时传自定义TCPConnector,把limit和limit_per_host调到与max_inflight_tasks匹配:
import aiohttp def make_session(max_inflight_tasks): # 把连接池上限与 max_inflight_tasks 联动 connector = aiohttp.TCPConnector( limit=max_inflight_tasks, # 总连接上限 limit_per_host=max_inflight_tasks, # 单 host 上限(推理引擎通常单 host) ) return aiohttp.ClientSession(connector=connector) # 真实用法示意: # session = make_session(max_inflight_tasks=200) # async with session.post(ENGINE_URL, json=payload) as resp: ...limit与limit_per_host都设为 200,200 个并发连接不再被池子卡住。
六、方案:用信号量显式控制 inflight,与连接池一致
第二层:除了连接池,再在任务调度层用asyncio.Semaphore(max_inflight_tasks)显式限制同时发起的请求数,两层上限保持一致,避免"任务层想发 200、连接池只放 100"的脱节:
async def bounded_request(session, sem, i, payload): async with sem: # 与连接池上限一致 async with session.post(ENGINE_URL, json=payload) as resp: return await resp.json() def build_semaphore(max_inflight_tasks): return asyncio.Semaphore(max_inflight_tasks) # 启动 200 个任务,信号量保证同时不超过 200,连接池也放 200,二者匹配信号量在"任务发起"层兜底,连接池在"建连"层兜底,双层一致才不脱节。
七、方案:配置一处指定,避免两处不一致
第三层:把max_inflight_tasks作为单一真相,连接池 limit、信号量都从它派生,杜绝两处各写各的数字:
CONFIG = {"max_inflight_tasks": 200} def build_runtime(): n = CONFIG["max_inflight_tasks"] connector = aiohttp.TCPConnector(limit=n, limit_per_host=n) session = aiohttp.ClientSession(connector=connector) sem = asyncio.Semaphore(n) return session, sem # 任何地方想改并发,只改 CONFIG["max_inflight_tasks"] 一处单一来源保证"连接池上限 == 信号量 == 配置",改一处全链路生效,不会再出现"配置 200、池子 100"。
八、验证:把"并发随 limit 提升"锁进测试
def test_default_cap_100(): peak = {"v": 0, "cur": 0} async def inst(n): async def one(i): peak["cur"] += 1; peak["v"] = max(peak["v"], peak["cur"]) async with FakeConnector(limit=100)._sem: await asyncio.sleep(0.001) peak["cur"] -= 1 await asyncio.gather(*[one(i) for i in range(n)]) asyncio.run(inst(200)) assert peak["v"] <= 100 def test_raised_limit_allows_more(): peak = {"v": 0, "cur": 0} async def inst(n, limit): sem = asyncio.Semaphore(limit) async def one(i): peak["cur"] += 1; peak["v"] = max(peak["v"], peak["cur"]) async with sem: await asyncio.sleep(0.001) peak["cur"] -= 1 await asyncio.gather(*[one(i) for i in range(n)]) asyncio.run(inst(200, 200)) assert peak["v"] == 200 if __name__ == "__main__": test_default_cap_100() test_raised_limit_allows_more() print("aiohttp 并发上限测试通过。")九、排查清单("并发卡在 100"按顺序查)
- 默认池:ClientSession 是否没传 connector,用了默认 limit=100?是则卡。
- limit 联动:TCPConnector.limit 是否随 max_inflight_tasks 调高?没调则脱节。
- 单 host:请求是否都打同一引擎 host?limit_per_host 也要调高。
- 隐性排队:超 100 的请求是否静默排队无报错?是则隐性限制。
- 双层一致:任务层信号量与连接池上限是否一致?不一致仍卡。
- 配置单一:max_inflight_tasks 是否单一来源派生各处上限?否则易错配。
- 峰值观测:能否观测真实并发峰值?能确认是否真被 100 卡住。
十、小结
"aiohttp 在 max_inflight_tasks>100 时限制到 100"是aiohttp 默认连接池上限 100 成了并发天花板,且与max_inflight_tasks配置脱节:所有请求打同一引擎 host,TCPConnector默认limit=limit_per_host=100,超出的请求静默排队,调大任务数无效。
修复三层:
- 调高连接池:创建
ClientSession传TCPConnector(limit=, limit_per_host=)与max_inflight_tasks匹配; - 信号量兜底:任务调度层用
asyncio.Semaphore(max_inflight_tasks),与连接池双层一致; - 配置单一:
max_inflight_tasks作单一真相,连接池与信号量都从它派生,改一处全生效。
核心原则:异步 HTTP 客户端的并发上限由连接池决定,不是由你的任务数决定。凡是"设了 max_inflight_tasks>100 却并发卡在 100"的现场,都应把 aiohttp 的TCPConnector.limit/limit_per_host调到匹配值,并与任务信号量保持一致——让配置真正生效。
