使用multiprocessing模块时Join()调用导致程序执行冻结
排查multiprocessing join()冻结的问题
我帮你找到了问题的根源,以及对应的解决办法:
核心原因:多余的Lock导致死锁
你代码里手动加了Lock来保护Queue的put()操作,但实际上**multiprocessing.Queue本身就是进程安全的**——它的put()和get()方法内部已经实现了同步机制,不需要额外加锁。
你的死锁场景是这样的:
- 第一个子进程获取锁后,调用
q.put(),如果队列的底层管道缓冲区被占满(哪怕没设置maxsize,系统管道也有默认大小限制),put()会进入阻塞状态 - 这时候锁还没被释放,第二个子进程会一直卡在
lock.acquire()步骤,无法继续执行 - 主进程调用
join()等待两个子进程结束,但两个子进程都处于阻塞状态,永远无法完成,最终导致程序冻结
解决步骤
- 移除多余的Lock:删除代码中所有和
Lock相关的代码(包括导入、初始化、传递参数、acquire/release操作) - 优化分片逻辑:用整数除法
//替代浮点数除法/,避免小数转整数的潜在问题
修改后的完整代码
import time import pandas as pd from multiprocessing import Process, Queue def writeDF(start, end, td, q): print('write from ', start, ' ', end) for i in range(start, end): q.put(td.iloc[i,:]) print('function writeDF completed') if __name__ == '__main__': td = pd.read_csv(r'C:\Users\dorian\Desktop\Analyzer.txt', encoding="ISO-8859-1", index_col="Start", parse_dates=True, sep=',') td = td[0:10] jobs = [] q = Queue() start_time = time.time() total_rows = td.shape[0] chunk_size = total_rows // 2 begin = 0 for i in range(2): stop = begin + chunk_size # 处理最后一个分片可能多出来的行(如果总行数是奇数) if i == 1: stop = total_rows t = Process(target=writeDF, args=(begin, stop, td, q)) t.start() jobs.append(t) begin = stop for p in jobs: print('try to join element:', p) p.join() print('element is joined') l = [] while not q.empty(): l.append(q.get()) df = pd.DataFrame(l) end_time = time.time() print('Value:', df.shape) print('2 process Time taken in seconds -', end_time - start_time)
额外优化建议
如果你的实际数据量很大,不建议把整个DataFrame传递给子进程(Windows下会序列化整个对象,开销很大),可以让每个子进程自己读取文件的对应分片,比如通过skiprows和nrows参数实现,这样能大幅提升性能。
内容的提问来源于stack exchange,提问作者Dorian75000
相关产品推荐
相关产品推荐

