使用ThreadPoolExecutor并行化处理mzML谱图数据时返回空DataFrame的问题解决及代码修改建议
使用ThreadPoolExecutor并行化处理mzML谱图数据时返回空DataFrame的问题解决及代码修改建议
嘿,我看你在尝试用ThreadPoolExecutor并行处理mzML文件时遇到了空DataFrame的问题,这大概率是因为pymzml的run.Reader对象不是线程安全的——多线程共享同一个Reader实例会导致读取指针混乱、数据丢失,最终返回空结果。下面我给你梳理问题根源和具体的修改方案:
核心问题分析
pymzml的run.Reader是为单线程读取设计的,它内部维护着文件读取的状态指针。当多个线程同时操作同一个Reader时,会出现指针跳转冲突,导致部分扫描无法被正确读取,甚至整个Reader实例失效,最终返回空的DataFrame。
具体修改步骤
1. 让谱图生成器支持独立读取子区间
修改_spectrum_generator,让它接收扫描起始和结束索引,并且每个线程创建独立的Reader实例,避免共享冲突:
def _spectrum_generator(filename, start_scan_idx, end_scan_idx): # 每个线程初始化自己的Reader实例,保证线程隔离 run = pymzml.run.Reader(filename, MS_precisions=MS_precisions) for spec_index in range(start_scan_idx, end_scan_idx): try: spec = run[spec_index] yield spec except Exception as e: print(f"读取扫描{spec_index}时出错: {e}") continue
2. 添加获取总扫描数的辅助函数
先获取文件的总扫描数,方便后续拆分并行任务:
def get_total_scans(filename): run = pymzml.run.Reader(filename, MS_precisions=MS_precisions) return len(run)
3. 重构数据收集函数为可并行的子任务
把原来的_gather_lcms_data拆分为处理单个扫描区间的子任务函数,每个子任务独立处理一部分扫描,返回该区间的结果片段:
def _process_scan_range(filename, start_idx, end_idx, min_mz, max_mz, polarity_filter="None", top_spectrum_peaks=100, include_polarity=False): all_mz = [] all_rt = [] all_polarity = [] all_i = [] all_scan = [] all_index = [] number_spectra = 0 all_msn_mz = [] all_msn_rt = [] all_msn_polarity = [] all_msn_scan = [] all_msn_level = [] for spec in _spectrum_generator(filename, start_idx, end_idx): rt = spec.scan_time_in_minutes() # 极性过滤逻辑 if polarity_filter != "None": scan_polarity = _get_scan_polarity(spec) if polarity_filter != scan_polarity: continue if spec.ms_level == 1: number_spectra += 1 try: # MZ范围过滤 if min_mz <= 0 and max_mz >= 2000: peaks = spec.peaks("raw") else: peaks = spec.reduce(mz_range=(min_mz, max_mz)) # 过滤无效峰(mz或强度小于1) peaks = peaks[~np.any(peaks < 1.0, axis=1)] # 取Top N高强度峰 peaks = peaks[peaks[:,1].argsort()][-top_spectrum_peaks:] mz, intensity = zip(*peaks) all_mz.extend(mz) all_i.extend(intensity) all_rt.extend([rt]*len(mz)) all_scan.extend([spec.ID]*len(mz)) all_index.extend([number_spectra]*len(mz)) # 添加极性信息 if include_polarity: scan_polarity = _get_scan_polarity(spec) polarity_val = POLARITY_POS if scan_polarity == "Positive" else POLARITY_NEG all_polarity.extend([polarity_val]*len(mz)) except Exception as e: print(f"处理MS1扫描{spec.ID}时出错: {e}") continue elif spec.ms_level > 1: try: msn_mz = spec.selected_precursors[0]["mz"] if msn_mz < min_mz or msn_mz > max_mz: continue all_msn_mz.append(msn_mz) all_msn_rt.append(rt) all_msn_scan.append(spec.ID) all_msn_level.append(spec.ms_level) # 添加极性信息 if include_polarity: scan_polarity = _get_scan_polarity(spec) polarity_val = POLARITY_POS if scan_polarity == "Positive" else POLARITY_NEG all_msn_polarity.append(polarity_val) except Exception as e: print(f"处理MSn扫描{spec.ID}时出错: {e}") continue # 组装当前区间的结果 ms1_results = { "mz": all_mz, "rt": all_rt, "i": all_i, "scan": all_scan, "index": all_index } msn_results = { "precursor_mz": all_msn_mz, "rt": all_msn_rt, "scan": all_msn_scan, "level": all_msn_level } if include_polarity: ms1_results["polarity"] = all_polarity msn_results["polarity"] = all_msn_polarity return pd.DataFrame(ms1_results), number_spectra, pd.DataFrame(msn_results)
4. 修改主函数实现并行处理
在_save_lcms_data_feather中,拆分扫描区间为多个子任务,用ThreadPoolExecutor并行执行,最后合并所有子任务的结果:
def _save_lcms_data_feather(filename): output_ms1_filename, output_msn_filename = _get_feather_filenames(filename) start = time.time() # 获取文件总扫描数 total_scans = get_total_scans(filename) print(f"待处理总扫描数: {total_scans}") # 拆分扫描区间(可根据机器性能调整chunk_size) chunk_size = 100 scan_ranges = [] for i in range(0, total_scans, chunk_size): end_idx = min(i + chunk_size, total_scans) scan_ranges.append( (filename, i, end_idx, 0, 10000, "None", 100000, True) ) # 并行处理所有区间 ms1_frames = [] msn_frames = [] total_spectra = 0 # 线程数建议用CPU核心数,避免IO竞争 with concurrent.futures.ThreadPoolExecutor(max_workers=os.cpu_count()) as executor: futures = [executor.submit(_process_scan_range, *args) for args in scan_ranges] # 遍历完成的任务,收集结果 for future in tqdm(concurrent.futures.as_completed(futures), total=len(futures)): ms1_df, num_spec, msn_df = future.result() ms1_frames.append(ms1_df) msn_frames.append(msn_df) total_spectra += num_spec # 合并所有子任务的结果 ms1_results = pd.concat(ms1_frames, ignore_index=True) msn_results = pd.concat(msn_frames, ignore_index=True) # 输出验证信息 print(ms1_results.head(10)) print(f"数据收集完成,耗时: {time.time() - start:.2f}秒") print(f"处理的MS1谱图总数: {total_spectra}") # 保存到feather文件 ms1_results = ms1_results.sort_values(by='i', ascending=False).reset_index(drop=True) ms1_results.to_feather(output_ms1_filename) msn_results.to_feather(output_msn_filename)
额外注意事项
- 线程数调整:不要设置过大的线程数,建议使用
os.cpu_count(),过多线程会导致磁盘IO竞争,反而降低处理速度。 - chunk_size优化:如果你的mzML文件单扫描数据量小,可以适当调大chunk_size;如果单扫描数据量大,调小chunk_size避免单个线程占用过多内存。
- 异常处理:子任务中添加了异常捕获,避免单个扫描处理失败导致整个任务崩溃,同时可以定位出错的扫描。
备注:内容来源于stack exchange,提问作者Astro
相关产品推荐
相关产品推荐

