如何在Python中使用apply_async调用类方法并更新共享字典?
问题:Python多进程调用类方法无响应,无法更新共享字典
刚接触Python多进程,尝试用apply_async并行更新共享字典,但调用类方法fill_contig_matrix时程序无任何反应,方法开头的打印语句也没输出。如何通过多进程正确更新该字典?
原简化代码
import numpy as np import multiprocessing as mp from multiprocessing import Manager import os class FeatureMatrix: index_dict = {'A': (0, 10), 'T': (1, 11), 'G': (2, 12), 'C': (3, 13), 'a': (4, 14), 't': (5, 15), 'g': (6, 16), 'c': (7, 17), '*': (8, 18), '#': (9, 19)} def __init__(self, pileup_file='calls_to_draft_pileup_test_2.txt', sam_file='calls_to_draft_sorted_test.sam'): self.pileup_file = pileup_file self.sam_file = sam_file self.score_count = 0 self.average_score_dict = {'A': [0, 0], 'T': [0, 0], 'G': [0, 0], 'C': [0, 0], 'a': [0, 0], 't': [0, 0], 'g': [0, 0], 'c': [0, 0], '*': [0, 0], '#': [0, 0]} self.contig_dict = self.contig_matrix() def contig_matrix(self): contig_dict = {} m = Manager() shared_contig_dict = m.dict() with open(self.sam_file, 'r', encoding='utf-8') as file: for line in file: if '@SQ' in line: line = line.strip() line_list = line.split('\t') contig = line_list[1].replace('SN:', '') length = int(line_list[2].replace('LN:', '')) contig_dict[contig] = self.initialize_matrices(length) if '@PG' in line: break shared_contig_dict.update(contig_dict) return shared_contig_dict def file_task_generator(self): with open(self.pileup_file, 'r', encoding='utf-8') as file: line_list = [] for line in file: line = line.strip() line = line.split('\t') if len(line_list) == 0 or line[0] == line_list[-1][0]: line_list.append(line) else: task = line_list line_list = [] line_list.append(line) yield task if line_list: task = line_list yield task def fill_contig_matrix(self, contig_list): print(f"Process {os.getpid()} is processing a task.") for line in contig_list: # Update dictionary... contig = line[0] # 补充获取contig的逻辑,原代码缺失 # 这里添加实际更新contig_dict的逻辑 file_name = f"array_{contig}.npz" np.savez(file_name, self.contig_dict[contig][0], self.contig_dict[contig][1]) def main(): matrix = FeatureMatrix() num_cpus = mp.cpu_count() pool = mp.Pool(processes=num_cpus) for contig_list in matrix.file_task_generator(): pool.apply_async(matrix.fill_contig_matrix, args=(contig_list)) pool.close() pool.join() if __name__ == '__main__': main()
问题分析与修复方案
1. 核心问题:args参数格式错误
apply_async的args需要传入元组,原代码中args=(contig_list)不是元组(括号会被忽略,等价于直接传contig_list),导致函数参数不匹配——fill_contig_matrix期望1个参数,但进程池会把contig_list的元素逐个拆解传递,引发未捕获的异常,任务静默失败。
修复:给args添加逗号,声明为单元素元组:
pool.apply_async(matrix.fill_contig_matrix, args=(contig_list,))
2. 补充缺失逻辑与异常处理
- 原代码
fill_contig_matrix中未定义contig变量,需要从contig_list的元素中提取(比如contig = contig_list[0][0]),否则会抛出NameError。 - 建议捕获任务异常,避免静默失败:可以保存每个任务的结果对象,后续调用
get()查看异常:
# 保存结果对象并检查异常 results = [] for contig_list in matrix.file_task_generator(): res = pool.apply_async(matrix.fill_contig_matrix, args=(contig_list,)) results.append(res) # 等待所有任务完成并检查异常 for res in results: try: res.get() except Exception as e: print(f"Task failed with error: {e}")
3. 共享字典的正确使用
当前用Manager.dict()创建的共享字典是正确的,但要注意:
Manager创建的共享对象有性能开销,大量频繁更新时可考虑优化(比如每个进程先处理本地数据,最后批量更新共享字典)。- 在Windows系统中,
Pool会通过pickle序列化类实例,self.contig_dict作为Manager.dict()序列化后仍能指向共享对象,无需额外修改。
修复后的完整main函数示例
def main(): matrix = FeatureMatrix() num_cpus = mp.cpu_count() pool = mp.Pool(processes=num_cpus) results = [] for contig_list in matrix.file_task_generator(): res = pool.apply_async(matrix.fill_contig_matrix, args=(contig_list,)) results.append(res) # 等待所有任务完成并捕获异常 for res in results: try: res.get() except Exception as e: print(f"Task failed with error: {e}") pool.close() pool.join()
内容的提问来源于stack exchange,提问作者rvijay1302
相关产品推荐
相关产品推荐

