如何在aiohttp服务中合并多个POST请求批量GPU处理后分别返回响应
aiohttp 批量合并请求调用GPU解决方案
不需要更换aiohttp框架,通过异步条件变量+Future对象即可实现需求,完整实现方案如下:
完整可运行代码
import asyncio from aiohttp import web # 批量处理配置 BATCH_SIZE = 5 # 待处理请求缓冲:存储结构为 (请求id, 请求data, 接收结果的Future对象) pending_requests = [] # 异步条件变量,控制缓冲并发访问和批次处理通知 cond = asyncio.Condition() # 替换为你的实际GPU处理逻辑,输入为data列表,输出为对应顺序的answer列表 async def gpu_process(data_list: list) -> list: # 示例逻辑,实际替换为GPU运算代码 return [f"answer_{i+1}" for i in range(len(data_list))] # 后台批量处理协程 async def batch_processor(): while True: async with cond: # 等待缓冲满足批量大小 await cond.wait_for(lambda: len(pending_requests) >= BATCH_SIZE) # 取出当前批次所有请求,清空缓冲 current_batch = pending_requests.copy() pending_requests.clear() # 提取批次内的所有data,批量调用GPU data_batch = [item[1] for item in current_batch] answer_batch = await gpu_process(data_batch) # 为每个请求写入对应结果,唤醒等待的请求协程 for idx, (req_id, _, future) in enumerate(current_batch): if not future.done(): future.set_result({ "id": req_id, "answers": answer_batch[idx] }) # 请求处理接口 async def my_func(request): post_data = await request.json() req_id = post_data["id"] req_data = post_data["data"] # 为当前请求创建专属Future对象,用于接收批量处理结果 loop = asyncio.get_running_loop() result_future = loop.create_future() # 加锁将当前请求加入缓冲 async with cond: pending_requests.append((req_id, req_data, result_future)) # 满足批量大小则通知处理器执行 if len(pending_requests) >= BATCH_SIZE: cond.notify() # 等待批量处理完成返回结果 response_data = await result_future return web.json_response(response_data) # 服务启动逻辑 async def main(): app = web.Application() app.add_routes([web.post('/', my_func)]) # 启动后台批量处理任务 app['batch_task'] = asyncio.create_task(batch_processor()) # 服务退出时清理后台任务 async def on_shutdown(app): app['batch_task'].cancel() await app['batch_task'] app.on_shutdown.append(on_shutdown) runner = web.AppRunner(app) await runner.setup() site = web.TCPSite(runner, '0.0.0.0', 8080) await site.start() print("服务启动在 http://0.0.0.0:8080") await asyncio.Event().wait() if __name__ == "__main__": asyncio.run(main())
核心实现说明
- 步骤1(请求攒批):通过全局缓冲列表存储所有进来的请求,每个请求加入后检查长度是否达到设定的批量大小,触发批量处理通知。
- 步骤3(对应返回结果):每个请求创建专属的
Future对象用来接收结果,批量处理时按请求在批次中的顺序匹配GPU返回的结果,写入对应Future后,等待的请求协程会自动唤醒并返回对应响应。
可选优化建议
- 增加超时 fallback逻辑:如果缓冲长时间没凑够批量大小,超过设定阈值(比如200ms)也强制触发处理,避免请求长时间挂起,仅需修改
batch_processor中的等待逻辑即可。 - 如果GPU处理是同步阻塞逻辑,可以用
asyncio.to_thread包裹调用,避免阻塞整个事件循环。 - 高并发场景下可以拆分多个缓冲队列,绑定不同的处理协程,避免单队列锁竞争。
内容的提问来源于stack exchange,提问作者Inq
相关产品推荐
相关产品推荐

