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
问题根源
- 错误的锁使用:
process_start里的with lock会让所有进程争抢同一个锁,只有第一个拿到锁的进程能进入循环处理,其他进程拿不到锁直接退出,这是多数进程处理0条或少量数据的核心原因。 - 不可靠的队列空判断:
q.empty()不是进程安全的判断方式,多进程场景下,判断队列非空后元素可能已被其他进程取走,导致get()阻塞或get_nowait()抛出异常,进而退出循环。 - 缺少退出信号:没有给进程明确的退出信号,进程无法判断队列是否真的处理完毕,只能依赖不可靠的
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()
关键改动说明
- 移除不必要的锁:
multiprocessing.Queue本身是进程安全的,额外加锁会导致进程无法并行处理,直接移除即可。 - 哨兵值退出机制:所有数据入队后,放入与进程数相等的哨兵值,进程收到哨兵后退出,同时将哨兵放回队列传递给其他进程,保证所有进程都能正常结束。
- 依赖
q.get()阻塞特性:q.get()会自动阻塞直到队列有元素,无需手动判断队列是否为空,避免多进程下的竞态问题。 - 取消启动休眠:队列已提前填满数据,进程启动后可立即处理,无需休眠等待。
- 内存优化建议:如果一次性加载500万对象内存占用过高,可使用
json.JSONDecoder流式解析JSON文件,边读边入队,降低内存压力。
内容的提问来源于stack exchange,提问作者bigpancakestack
相关产品推荐
相关产品推荐

