Python多进程加join后访问Manager.list报无connection属性
编写多进程DNA基因组互相比对脚本时,采用multiprocessing实现并行计算,所有子进程会向公共共享列表genome_score_avgs追加计算结果。
主进程代码
if __name__ == "__main__": start = time.perf_counter() with Manager() as manager: genome_score_avgs = manager.list() processes = [Process(target=compareGenomes, args=(chunk, genome_score_avgs,)) for chunk in divideGenomes('TEST_DIR')] for p in processes: p.start() for p in processes: p.join() print(genome_score_avgs) print(*createTimeline(genome_score_avgs), sep='\n') print(f'Finished in {time.perf_counter() - start} seconds')
运行报错信息
Traceback (most recent call last): File "/Library/Frameworks/Python.framework/Versions/3.9/lib/python3.9/multiprocessing/managers.py", line 801, in _callmethod conn = self._tls.connection AttributeError: 'ForkAwareLocal' object has no attribute 'connection' During handling of the above exception, another exception occurred: Traceback (most recent call last): File "/Users/ayushpal/Coding/PythonStuff/C4DInter/main.py", line 59, in <module> print(*createTimeline(genome_score_avgs), sep='\n') File "/Users/ayushpal/Coding/PythonStuff/C4DInter/main.py", line 42, in createTimeline min_score = min(score_avgs, key=lambda x: x[2]) File "<string>", line 2, in __getitem__ File "/Library/Frameworks/Python.framework/Versions/3.9/lib/python3.9/multiprocessing/managers.py", line 805, in _callmethod self._connect() File "/Library/Frameworks/Python.framework/Versions/3.9/lib/python3.9/multiprocessing/managers.py", line 792, in _connect conn = self._Client(self._token.address, authkey=self._authkey) File "/Library/Frameworks/Python.framework/Versions/3.9/lib/python3.9/multiprocessing/connection.py", line 507, in Client c = SocketClient(address) File "/Library/Frameworks/Python.framework/Versions/3.9/lib/python3.9/multiprocessing/connection.py", line 635, in SocketClient s.connect(address) FileNotFoundError: [Errno 2] No such file or directory <ListProxy object, typeid 'list' at 0x7fc04ea36bb0; '__str__()' failed>
此前查询到同类报错的诱因为主进程早于子进程结束,导致共享列表被销毁,解决方案是对所有进程调用p.join()等待执行完成,但代码中已经实现join逻辑,仍然抛出相同错误。
子进程函数代码
def compareGenomes(genome_pairings, genome_score_avgs): scores = [] for genome1, genome2 in genome_pairings: print(genome1, genome2) for i, seq in enumerate(genome1.protein_seqs): for j, seq2 in enumerate(genome2.protein_seqs[i::]): alignment = align.globalxx(seq, seq2) scores.append(alignment) top_scores = [] for i in range(len(genome1.protein_seqs)): top_scores.append(max(scores, key=lambda x: x[0][2] / len(x[0][1]))) scores.remove(max(scores, key=lambda x: x[0][2] / len(x[0][1]))) avg_score = sum([i[0][2] / len(i[0][1]) for i in top_scores]) / len(top_scores) with open(f'alignments/{genome1.name}x{genome2.name}.txt', 'a') as file: file.writelines([format_alignment(*i[0]) for i in top_scores]) genome_score_avgs.append((genome1, genome2, avg_score))
报错的核心原因是访问genome_score_avgs的代码位于with Manager() as manager上下文管理器的作用域外。Manager()的with块在执行到缩进范围外时,会自动停止Manager服务、销毁进程间通信用的套接字文件和共享资源。即使已经通过p.join()等待所有子进程执行完毕,只要退出with块,manager.list()生成的ListProxy代理对象就失去了后端服务支撑,此时再调用print(genome_score_avgs)、将其传入createTimeline()做遍历取值操作时,代理对象无法连接到已经关闭的Manager服务,就会抛出看到的FileNotFoundError。
之前查到的“未调用join导致主进程提前退出”是同类报错的另一种触发场景,本质都是访问代理对象时Manager服务已经终止,和当前场景的触发原因一致,只是终止Manager的时机不同。
- 方案1:将所有访问
genome_score_avgs的逻辑移到with Manager()的代码块内部,确保访问共享列表时Manager服务处于运行状态:
if __name__ == "__main__": start = time.perf_counter() with Manager() as manager: genome_score_avgs = manager.list() processes = [Process(target=compareGenomes, args=(chunk, genome_score_avgs,)) for chunk in divideGenomes('TEST_DIR')] for p in processes: p.start() for p in processes: p.join() # 所有访问代理列表的操作都放在with块内 print(genome_score_avgs) print(*createTimeline(genome_score_avgs), sep='\n') print(f'Finished in {time.perf_counter() - start} seconds')
- 方案2:退出with块之前,将Manager托管的列表转换为Python原生列表,后续直接操作原生列表即可脱离对Manager服务的依赖:
if __name__ == "__main__": start = time.perf_counter() with Manager() as manager: genome_score_avgs = manager.list() processes = [Process(target=compareGenomes, args=(chunk, genome_score_avgs,)) for chunk in divideGenomes('TEST_DIR')] for p in processes: p.start() for p in processes: p.join() # 转换为原生列表,脱离Manager依赖 result_list = list(genome_score_avgs) # 后续直接操作原生列表 print(result_list) print(*createTimeline(result_list), sep='\n') print(f'Finished in {time.perf_counter() - start} seconds')
优先推荐第二种方案:转换为原生列表后,后续的遍历、最值计算等操作不需要走Manager的进程间通信代理,执行速度远快于直接操作ListProxy对象,数据量越大性能差距越明显。
内容的提问来源于stack exchange,提问作者Ayush

