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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 20:24:56