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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 05:15:04