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

Python多进程结合生成器使用时代码挂起的原因咨询

Python多进程结合生成器使用时代码挂起的原因咨询

嘿,我来帮你排查下代码挂起的问题!你的思路没问题——用生成器减少内存占用+多进程加速处理,但代码里有两个关键问题导致了卡住:

问题分析

  • 队列被单个元素填满,主进程阻塞:你的生成器每次只产出一个整数,主进程会把100万个单个元素挨个往输入队列里塞。虽然队列默认容量不小,但当生产速度远快于worker的消费速度时,队列会被填满,这时候主进程执行input_queue.put(data_chunk)就会一直卡住,根本没法走到后面往队列放终止信号None、等待进程结束的步骤,自然就僵住了。
  • 处理逻辑不匹配,效率极低:process_data是为处理数据块(列表)写的,但你每次传单个整数,不仅没发挥多进程批量处理的优势,还因为频繁的队列通信白白增加了开销。

修复方案

方案一:手动调整队列+批量生成数据块

把生成器改成批量产出数据块,同时限制队列大小避免内存占用过高,代码修改如下:

import multiprocessing

def process_data(data_chunk):
    return [x * 2 for x in data_chunk]  

def data_generator(chunk_size=1000):
    chunk = []
    for i in range(1000000):
        chunk.append(i)
        if len(chunk) == chunk_size:
            yield chunk
            chunk = []
    # 处理最后一批不足chunk_size的数据
    if chunk:
        yield chunk

def worker(input_queue, output_queue):
    while True:
        data_chunk = input_queue.get()
        if data_chunk is None:
            break
        output_queue.put(process_data(data_chunk))

if __name__ == "__main__":
    # 限制队列大小为CPU核心数的2倍,避免内存占用过高
    input_queue = multiprocessing.Queue(maxsize=multiprocessing.cpu_count()*2)
    output_queue = multiprocessing.Queue()
    num_processes = multiprocessing.cpu_count()
    processes = [multiprocessing.Process(target=worker, args=(input_queue, output_queue)) for _ in range(num_processes)]
  
    for p in processes:
        p.start()

    # 生成数据块并放入队列
    for data_chunk in data_generator():
        input_queue.put(data_chunk)

    # 给每个worker发终止信号
    for _ in range(num_processes):
        input_queue.put(None)

    # 收集处理结果(如果需要)
    results = []
    while not output_queue.empty():
        results.extend(output_queue.get())

    for p in processes:
        p.join()

方案二:用multiprocessing.Pool更简洁实现

其实Python的Pool类已经帮我们封装了队列和进程管理,用imap或者imap_unordered可以完美结合生成器,既省内存又不用手动处理队列,代码更简洁:

import multiprocessing

def process_item(x):
    return x * 2

def data_generator():
    for i in range(1000000):
        yield i

if __name__ == "__main__":
    num_processes = multiprocessing.cpu_count()
    with multiprocessing.Pool(num_processes) as pool:
        # chunksize参数控制批量处理的大小,平衡通信开销和内存占用
        results = list(pool.imap(process_item, data_generator(), chunksize=1000))

备注:内容来源于stack exchange,提问作者FredTheRick

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 09:48:08