Python 2.7中利用multiprocessing并行化大DNA文件解析函数的方法
并行处理大型DNA文件的Python 2.7解决方案
作为Python 2.7新手,处理4GB的大文件确实会遇到单进程速度慢的问题,咱们一步步来解决这个问题。
首先,先优化下你现有的filerd函数,用enumerate可以省去手动维护count变量的麻烦,代码更简洁易读:
def filerd(f): identifier = [] with open(f,'r') as inputfile: # 用enumerate从1开始计数行号 for idx, line in enumerate(inputfile, 1): if idx % 4 == 1: identifier.append(line) return identifier
接下来咱们说并行化的思路:大文件并行处理的核心是分块处理,但因为你的需求是每4行取第1行,所以分块必须保证不会把一个4行的组拆到两个块里,否则会漏取或者重复取行。下面给你两个可行的方案:
方案一:基于行号分块的并行处理(内存足够时用)
这个方案先统计文件的总行数,然后根据CPU核心数把文件分成多个块,每个块的起始行都是4n+1的格式(保证从目标行开始),每个进程负责处理一个块,最后合并结果。
步骤1:编写统计总行数的函数
虽然统计4GB文件的行数需要一点时间,但这是保证分块准确的基础:
def count_lines(fname): with open(fname, 'r') as f: return sum(1 for _ in f)
步骤2:编写单进程处理块的函数
每个进程负责读取指定行范围内的内容,筛选出符合条件的行:
def process_chunk(fname, start_line, end_line): identifiers = [] with open(fname, 'r') as f: # 跳过起始行之前的所有行 for _ in xrange(start_line - 1): f.readline() current_line = start_line while current_line <= end_line: line = f.readline() if not line: # 提前读到文件结尾就退出 break if current_line % 4 == 1: identifiers.append(line) current_line += 1 return identifiers
步骤3:主函数实现并行调度
因为Python 2.7的multiprocessing.Pool没有starmap方法,咱们用apply_async来异步提交任务并收集结果:
import multiprocessing def parallel_filerd(fname): total_lines = count_lines(fname) num_processes = multiprocessing.cpu_count() # 用CPU核心数作为进程数 chunk_size = total_lines // num_processes chunks = [] for i in range(num_processes): # 计算块的起始行,调整为4n+1格式 start = i * chunk_size + 1 if start % 4 != 1: start += (4 - start % 4) # 计算块的结束行,最后一个块直接到文件末尾,其他块调整为4n格式 if i == num_processes - 1: end = total_lines else: end = (i + 1) * chunk_size end -= end % 4 # 如果起始行超过结束行,跳过这个块 if start > end: continue chunks.append((fname, start, end)) # 创建进程池并提交任务 pool = multiprocessing.Pool(num_processes) results = [] for chunk in chunks: res = pool.apply_async(process_chunk, chunk) results.append(res) # 收集所有进程的结果并合并 all_identifiers = [] for res in results: all_identifiers.extend(res.get()) pool.close() pool.join() return all_identifiers
使用注意
如果是在Windows系统下运行,必须把主程序调用放在if __name__ == '__main__':块里,否则会出现进程重复创建的问题:
if __name__ == '__main__': dna_identifiers = parallel_filerd('your_large_dna_file.fastq') # 这里可以继续处理结果
方案二:分块写入临时文件(内存不足时用)
如果结果列表太大(比如4GB文件处理后可能占1GB以上内存),可以让每个进程把结果写入临时文件,最后再合并所有临时文件,这样不会占用大量内存:
修改块处理函数
import tempfile import os def process_chunk_to_file(fname, start_line, end_line, temp_file): with open(fname, 'r') as f, open(temp_file, 'w') as out_f: for _ in xrange(start_line - 1): f.readline() current_line = start_line while current_line <= end_line: line = f.readline() if not line: break if current_line % 4 == 1: out_f.write(line) current_line += 1
主函数实现
def parallel_filerd_to_file(fname, output_file): total_lines = count_lines(fname) num_processes = multiprocessing.cpu_count() chunk_size = total_lines // num_processes chunks = [] temp_files = [] for i in range(num_processes): start = i * chunk_size + 1 if start % 4 != 1: start += (4 - start % 4) if i == num_processes - 1: end = total_lines else: end = (i + 1) * chunk_size end -= end % 4 if start > end: continue # 创建临时文件 temp_f = tempfile.mktemp() temp_files.append(temp_f) chunks.append((fname, start, end, temp_f)) pool = multiprocessing.Pool(num_processes) results = [] for chunk in chunks: res = pool.apply_async(process_chunk_to_file, chunk) results.append(res) # 等待所有进程完成 for res in results: res.get() pool.close() pool.join() # 合并所有临时文件到输出文件 with open(output_file, 'w') as out_f: for temp_f in temp_files: with open(temp_f, 'r') as in_f: out_f.write(in_f.read()) os.unlink(temp_f) # 删除临时文件
一些额外建议
- 进程数不用设置得比CPU核心数多太多,否则会增加进程切换的开销,反而变慢。
- 如果你的DNA文件是FASTQ格式,刚好就是每4行一组,这个方案完全适配。
- 测试的时候可以先用小文件验证逻辑是否正确,再用大文件跑。
内容的提问来源于stack exchange,提问作者Amjad Syed
相关产品推荐
相关产品推荐

