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

Python实现服务器并发请求处理:空闲时即时获取新请求

异步服务器请求调度优化:单请求完成后立即补新请求

需求与问题

现有场景:5台服务器,每台最多可并发处理2个请求;初始有10个待处理请求,要求服务器处理完任意一个请求、出现空闲时,立即从请求池获取新请求。

当前代码存在的问题:服务器worker会一次性取出最多2个请求,等待全部处理完成后才会取下一批,无法实现单请求完成后立即补新请求的需求。

修改后的代码

import asyncio
import random

async def process_request(server_id, request_id):
    processing_time = random.randint(10, 30)
    print(f"Server {server_id} is processing request {request_id} for {processing_time} seconds")
    await asyncio.sleep(processing_time)
    print(f"Server {server_id} finished processing request {request_id}")

async def server_worker_single(server_id, queue):
    # 单个并发处理协程:持续从队列取请求,处理完一个就取下一个
    while True:
        request_id = await queue.get()
        try:
            await process_request(server_id, request_id)
            # 处理完成后添加新请求到队列(模拟请求持续产生)
            await queue.put(random.randint(1, 100))
        finally:
            queue.task_done()

async def server_worker(server_id, queue, num_concurrent_requests_per_server):
    # 启动对应并发数的单请求处理协程
    tasks = [asyncio.create_task(server_worker_single(server_id, queue)) 
             for _ in range(num_concurrent_requests_per_server)]
    await asyncio.gather(*tasks)

async def main():
    num_servers = 5
    num_concurrent_requests_per_server = 2
    # 全局请求队列,所有服务器从同一请求池取请求
    request_queue = asyncio.Queue()

    # 初始添加10个请求(5台服务器×每台2并发)
    for _ in range(num_servers * num_concurrent_requests_per_server):
        await request_queue.put(random.randint(1, 100))

    # 启动所有服务器worker
    server_tasks = []
    for i in range(num_servers):
        task = asyncio.create_task(server_worker(i, request_queue, num_concurrent_requests_per_server))
        server_tasks.append(task)

    # 模拟运行2分钟后停止,可根据需求替换为固定请求数处理完成的逻辑
    await asyncio.sleep(120)

    # 取消所有服务器任务并等待结束
    for task in server_tasks:
        task.cancel()
    await asyncio.gather(*server_tasks, return_exceptions=True)

if __name__ == "__main__":
    asyncio.run(main())

关键改动说明

  • 拆分并发处理逻辑:将原有的批量取请求、批量等待完成的逻辑,拆分为单个请求的独立处理协程。每台服务器启动对应并发数的server_worker_single协程,每个协程持续从队列取请求,处理完一个立即取下一个,实现空闲即补新请求的效果。
  • 全局请求池优化:将原有的单服务器队列改为全局请求队列,更贴合"请求池"的设计逻辑,所有服务器共享请求资源。若需保留单服务器独立队列,只需将request_queue改为每个服务器单独实例即可,核心处理逻辑不变。
  • 任务完成标记:添加queue.task_done()确保队列的任务计数正确,若后续需要用await request_queue.join()等待所有请求处理完成,可调整停止逻辑适配。

内容的提问来源于stack exchange,提问作者xan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 21:45:30