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

Python多进程队列未处理完数据即退出的问题及优化需求

问题解决:多进程处理Queue时提前退出的问题

问题背景

有一个包含约500万个对象的JSON文件,用Python解析后将对象放入multiprocessing.Queue,再用多进程处理。但多数进程仅处理少量对象就退出,队列实际未清空,只有启动进程前休眠10秒才能正常处理,但会造成不必要延迟。

原始代码

INPUT_JSON = "file.json"

def process_start(id,q,lock):
        tGetCount = 0
        with lock:
                while not q.empty():
                        try:
                                o = q.get() # also tried o=q.get_nowait()
                                tGetCount += 1
                                if o['data'] == "someCondition":
                                        someFunction()
                        except queue.Empty:
                                time.sleep(1)
                                continue
        print("Process no. {} parsed {} records. Queue size: {}".format(id,tGetCount,q.qsize()))

def queue_dataset():
        cpu_count = 24
        mp_lock = multiprocessing.Lock()
        procs = []
        q = multiprocessing.Queue()

        f = open(INPUT_JSON)
        data = json.load(f)
        f.close()

        # Insert into queue
        tInsertCount = 0
        for obj in data:
                q.put(obj)
                tInsertCount += 1
                if tInsertCount % 100000 == 0:
                        print("Status: Inserted {}".format(tInsertCount))
                        write_log("Status: Inserted {}".format(tInsertCount))

        print("Insertion done")

        for p in range(cpu_count):
                proc = Process(target=process_start, args=(p,q,mp_lock))
                procs.append(proc)
                proc.start()
                time.sleep(10) # Only a couple hundered object processed if program does not sleep 

        for p in procs:
                p.join(timeout=20)

        print("Number of records inserted ({})- current size ({})= parsed: {}".format(tInsertCount,q.qsize(),tInsertCount-q.qsize()))

不同运行场景结果

1. 休眠1秒并使用get_nowait()

Process no. 0 parsed 31862 records. Queue size: 5054794
Process no. 1 parsed 33558 records. Queue size: 5021236
Process no. 2 parsed 34426 records. Queue size: 4986810
Process no. 3 parsed 33051 records. Queue size: 4953759
Process no. 4 parsed 34004 records. Queue size: 4919755
Process no. 5 parsed 33755 records. Queue size: 4886000
Process no. 6 parsed 33435 records. Queue size: 4852565
Process no. 7 parsed 33590 records. Queue size: 4818975
Process no. 8 parsed 35165 records. Queue size: 4783810
Process no. 9 parsed 34140 records. Queue size: 4749670
Process no. 10 parsed 27546 records. Queue size: 4722124
Process no. 11 parsed 33684 records. Queue size: 4688440
Process no. 12 parsed 35285 records. Queue size: 4653155
Process no. 13 parsed 32657 records. Queue size: 4620498
Process no. 14 parsed 35163 records. Queue size: 4585335
Process no. 15 parsed 31848 records. Queue size: 4553487
Process no. 16 parsed 33780 records. Queue size: 4519707
Process no. 17 parsed 34900 records. Queue size: 4484807
Process no. 18 parsed 34030 records. Queue size: 4450777
Process no. 19 parsed 32595 records. Queue size: 4418182
Process no. 20 parsed 35333 records. Queue size: 4382849
Process no. 21 parsed 32484 records. Queue size: 4350365
Process no. 22 parsed 34015 records. Queue size: 4316350
Number of records inserted (5086656)- current size (4247227)= parsed: 839429
Process no. 23 parsed 69393 records. Queue size: 4246957

2. 使用get_nowait()且进程间休眠10秒时的运行结果

Process no. 0 parsed 339800 records. Queue size: 4746856
Process no. 1 parsed 341506 records. Queue size: 4405350
Process no. 2 parsed 334232 records. Queue size: 4071118
Process no. 3 parsed 333835 records. Queue size: 3737283
Process no. 4 parsed 328776 records. Queue size: 3408507
Process no. 5 parsed 332531 records. Queue size: 3075976
Process no. 6 parsed 335032 records. Queue size: 2740944
Process no. 7 parsed 339034 records. Queue size: 2401910
Process no. 8 parsed 343183 records. Queue size: 2058727
Process no. 9 parsed 341165 records. Queue size: 1717562
Process no. 10 parsed 337569 records. Queue size: 1379993
Process no. 11 parsed 345278 records. Queue size: 1034715
Process no. 12 parsed 345037 records. Queue size: 689678
Process no. 13 parsed 334748 records. Queue size: 354930
Process no. 14 parsed 344517 records. Queue size: 10413
Process no. 15 parsed 10413 records. Queue size: 0
Process no. 16 parsed 0 records. Queue size: 0
Process no. 17 parsed 0 records. Queue size: 0
Process no. 18 parsed 0 records. Queue size: 0
Process no. 19 parsed 0 records. Queue size: 0
Process no. 20 parsed 0 records. Queue size: 0
Process no. 21 parsed 0 records. Queue size: 0
Process no. 22 parsed 0 records. Queue size: 0
Process no. 23 parsed 0 records. Queue size: 0
Number of records inserted (5086656)- current size (0)= parsed: 5086656

3. 使用get且不休眠时的运行结果

Process no. 0 parsed 255 records. Queue size: 5086401
Process no. 1 parsed 281 records. Queue size: 5086120
Process no. 2 parsed 0 records. Queue size: 5086120
Process no. 3 parsed 7 records. Queue size: 5086113
Process no. 4 parsed 9 records. Queue size: 5086104
Process no. 5 parsed 6 records. Queue size: 5086098
Process no. 6 parsed 1 records. Queue size: 5086097
Process no. 7 parsed 1 records. Queue size: 5086096
Process no. 8 parsed 11 records. Queue size: 5086085
Process no. 9 parsed 6 records. Queue size: 5086079
Process no. 10 parsed 1 records. Queue size: 5086078
Process no. 11 parsed 2 records. Queue size: 5086076
Process no. 12 parsed 1 records. Queue size: 5086075
Process no. 13 parsed 5 records. Queue size: 5086070
Process no. 14 parsed 8 records. Queue size: 5086062
Process no. 15 parsed 5 records. Queue size: 5086057
Process no. 16 parsed 1 records. Queue size: 5086056
Process no. 17 parsed 6 records. Queue size: 5086050
Process no. 18 parsed 0 records. Queue size: 5086050
Process no. 19 parsed 0 records. Queue size: 5086050
Process no. 20 parsed 1 records. Queue size: 5086049
Process no. 21 parsed 1 records. Queue size: 5086048
Process no. 22 parsed 10 records. Queue size: 5086038
Process no. 23 parsed 11322 records. Queue size: 5074716
Number of records inserted (5086656)- current size (5074716)= parsed: 11940

问题根源

  1. 错误的锁使用:process_start里的with lock会让所有进程争抢同一个锁,只有第一个拿到锁的进程能进入循环处理,其他进程拿不到锁直接退出,这是多数进程处理0条或少量数据的核心原因。
  2. 不可靠的队列空判断:q.empty()不是进程安全的判断方式,多进程场景下,判断队列非空后元素可能已被其他进程取走,导致get()阻塞或get_nowait()抛出异常,进而退出循环。
  3. 缺少退出信号:没有给进程明确的退出信号,进程无法判断队列是否真的处理完毕,只能依赖不可靠的q.empty()。

修正后的代码

import multiprocessing
import json

INPUT_JSON = "file.json"
SENTINEL = None  # 哨兵值,用于通知进程退出

def process_start(id, q):
    tGetCount = 0
    while True:
        o = q.get()
        if o is SENTINEL:
            # 收到哨兵值,退出循环,同时放回哨兵让其他进程接收
            q.put(SENTINEL)
            break
        tGetCount += 1
        if o['data'] == "someCondition":
            # 替换为你的实际处理函数
            someFunction()
    print(f"Process no. {id} parsed {tGetCount} records.")

def queue_dataset():
    cpu_count = multiprocessing.cpu_count()
    procs = []
    q = multiprocessing.Queue()

    # 读取JSON文件(500万对象一次性加载内存压力大的话,可改用流式解析)
    with open(INPUT_JSON) as f:
        data = json.load(f)

    # 插入数据到队列
    tInsertCount = 0
    for obj in data:
        q.put(obj)
        tInsertCount += 1
        if tInsertCount % 100000 == 0:
            print(f"Status: Inserted {tInsertCount}")
            write_log(f"Status: Inserted {tInsertCount}")

    print("Insertion done")

    # 放入与进程数相同的哨兵值,确保每个进程都能收到退出信号
    for _ in range(cpu_count):
        q.put(SENTINEL)

    # 启动所有进程,无需休眠
    for p in range(cpu_count):
        proc = multiprocessing.Process(target=process_start, args=(p, q))
        procs.append(proc)
        proc.start()

    # 等待所有进程结束
    for p in procs:
        p.join()

    print(f"Number of records inserted: {tInsertCount}, all processed.")

if __name__ == "__main__":
    queue_dataset()

关键改动说明

  1. 移除不必要的锁:multiprocessing.Queue本身是进程安全的,额外加锁会导致进程无法并行处理,直接移除即可。
  2. 哨兵值退出机制:所有数据入队后,放入与进程数相等的哨兵值,进程收到哨兵后退出,同时将哨兵放回队列传递给其他进程,保证所有进程都能正常结束。
  3. 依赖q.get()阻塞特性:q.get()会自动阻塞直到队列有元素,无需手动判断队列是否为空,避免多进程下的竞态问题。
  4. 取消启动休眠:队列已提前填满数据,进程启动后可立即处理,无需休眠等待。
  5. 内存优化建议:如果一次性加载500万对象内存占用过高,可使用json.JSONDecoder流式解析JSON文件,边读边入队,降低内存压力。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 12:42:05