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

如何实现支持动态扩容任务队列的异步/并行处理方案?

动态任务队列的并行处理解决方案

Python的multiprocessing库完全支持动态添加任务的场景,你之前没找到合适方案是因为不能直接用普通列表做任务队列——多进程间内存相互隔离,子进程对普通列表的修改无法同步到主进程。正确的做法是用进程安全的队列管理任务,结合worker进程实现动态任务处理。

核心思路

  • 用multiprocessing.JoinableQueue存储任务:它支持多进程安全读写,还能让主进程等待所有任务处理完成。
  • 启动多个worker进程:每个worker循环从队列取任务执行,处理过程中生成的新任务直接放入队列。
  • 用共享计数器保证新任务编号唯一:避免多进程下的编号冲突。
  • 任务完成后发送结束信号:让worker进程有序退出。

实现代码

import multiprocessing

def process_job_logic(job):
    # 替换为你的实际任务处理逻辑
    print(f"Processing {job}...")

def worker(task_queue, counter):
    while True:
        job = task_queue.get()
        if job is None:
            # 收到结束信号,终止当前worker
            task_queue.task_done()
            break
        
        # 执行任务处理
        process_job_logic(job)
        
        # 检查是否需要生成新任务
        job_num = int(job.split('job')[-1])
        if job_num % 3 == 0:
            # 更新共享计数器,生成新任务
            with counter.get_lock():
                counter.value += 1
                new_job = f'new job{counter.value}'
            task_queue.put(new_job)
            print(f"Added new task: {new_job}")
        
        task_queue.task_done()

if __name__ == '__main__':
    # 初始化共享计数器(初始任务编号到5,新任务从6开始)
    counter = multiprocessing.Manager().Value('i', 5)
    
    # 创建可等待的任务队列
    task_queue = multiprocessing.JoinableQueue()
    
    # 添加初始任务
    for i in range(1, 6):
        task_queue.put(f'job{i}')
    
    # 启动worker进程(数量建议和CPU核心数一致)
    worker_count = multiprocessing.cpu_count()
    workers = []
    for _ in range(worker_count):
        p = multiprocessing.Process(target=worker, args=(task_queue, counter))
        p.start()
        workers.append(p)
    
    # 等待队列中所有任务处理完成
    task_queue.join()
    
    # 给每个worker发送结束信号
    for _ in range(worker_count):
        task_queue.put(None)
    
    # 等待所有worker进程退出
    for p in workers:
        p.join()

关键细节说明

  • JoinableQueue:通过join()方法,主进程会阻塞直到队列中所有任务都被task_done()标记完成,确保没有任务遗漏。
  • 共享计数器:multiprocessing.Manager().Value创建的计数器是进程安全的,get_lock()保证多进程修改时的原子性,避免编号重复。
  • worker退出逻辑:所有任务处理完成后,主进程向队列放入与worker数量相同的None,每个worker取到None后就会退出,避免无限阻塞。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 18:25:09