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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.23 14:47:40