Python函数可在Process中运行,但apply_async无法正常工作
问题:multiprocessing Pool.apply_async 无法适配pybloomfiltermmap3布隆过滤器?
我有一段基于Python multiprocessing库、使用Process类可正常运行的代码,现希望改为异步方式以支持用户设置线程数。简化后的代码如下:
from scripts import classify_reads from multiprocessing import Value, Process, Pool from multiprocessing.managers import BaseManager
shared_k_size = Value('i', k_size) shared_n_consec_matches = Value('i', n_consec_matches) shared_output = Value(c_wchar_p, args["output"]) shared_anchor_proportion_cutoff = Value('f', anchor_proportion_cutoff) # 创建用于共享内存集合的自定义管理器 class CustomManager(BaseManager): pass CustomManager.register('shared_set', set) manager = CustomManager() manager.start() shared_kmer_set = manager.shared_set() def classify_reads(fastq_file, bloom_filter, kmer_size, n_consecutive_matches, output, shared_anchor_kmer_set): kmer_set = set() # 填充kmer_set print(len(kmer_set)) shared_anchor_kmer_set.update(anchor_kmer_set) # 这段代码可以正常运行: processes = [Process(target=classify_reads, args=(i, present_kmers_bf, shared_k_size.value, shared_n_consec_matches.value, shared_output.value, shared_kmer_set)) for i in args["input"]] for i in range(0, len(processes), args["threads"]): for process in processes[i:i + args["threads"]]: process.start() for process in processes: process.join() for process in processes: process.terminate() # 这段代码无法运行: with Pool(processes=args["threads"]) as pool: print(f"starting async") for i in args["input"]: try : pool.apply_async(classify_reads.classify_reads, args=(i, present_kmers_bf, shared_k_size.value, shared_n_consec_matches.value, shared_output.value, shared_kmer_set,)) except Exception as e: print(f"An error occurred: {e}") pool.close() pool.join()
其中present_kmers_bf是pybloomfiltermmap3库提供的布隆过滤器,其余使用类似参数(含自定义管理器)的代码已成功改为apply_async,但此段代码无法正常运行。我推测mmap文件可能与apply_async不兼容,但官方文档显示它支持pickle序列化。
内容的提问来源于stack exchange,提问作者Yama
相关产品推荐
相关产品推荐

