如何将多份甲基化TSV文件合并为单个Arrow表并规避内存溢出
解决PyArrow处理数千个甲基化TSV文件的内存溢出问题(生成宽表)
你当前的问题在于dataset.to_batches是按文件顺序逐批读取数据,没有实现按Name探针名称的全局合并。要生成包含所有探针并集的宽表(Name为行,每个样本Beta为列),需要分阶段处理,利用列存格式和分组/Join操作避免内存溢出,以下是可行方案:
方案一:PyArrow + 分阶段列存转换 + 分组聚合
步骤1:将TSV转换为带样本标识的Parquet(列存格式)
先把每个TSV文件转换为Parquet,同时为每个样本的Beta列添加样本名(从文件名提取),避免重复解析TSV,提升后续处理效率:
import pyarrow as pa import pyarrow.csv as pv import pyarrow.parquet as pq import os from pathlib import Path # 配置TSV解析选项 parse_options = pv.ParseOptions(delimiter="\t") convert_options = pv.ConvertOptions( include_columns=["Name", "Beta"], # 仅读取需要的两列 column_types={"Name": pa.string(), "Beta": pa.float64()} # 指定类型,加速解析 ) # 创建临时目录存储Parquet文件 temp_parquet_dir = Path("./temp_methylation_parquet") temp_parquet_dir.mkdir(exist_ok=True) # 批量转换TSV到Parquet for tsv_path in filtered_samples: # 从文件名提取样本名(假设文件名格式为sample_id.tsv) sample_name = os.path.splitext(os.path.basename(tsv_path))[0] # 读取TSV并仅保留目标列 table = pv.read_csv( tsv_path, parse_options=parse_options, convert_options=convert_options ) # 重命名Beta列为样本名,便于后续识别 table = table.rename_columns(["Name", sample_name]) # 写入Parquet(列存格式,压缩比高,读取快) pq.write_table( table, temp_parquet_dir / f"{sample_name}.parquet", compression="snappy" # 启用压缩,减少磁盘占用 )
步骤2:按探针名称聚合生成宽表
用PyArrow Dataset读取所有Parquet文件,按Name分组聚合,将每个样本的Beta值映射到对应列:
import pyarrow.dataset as ds import pyarrow.compute as pc # 创建Parquet数据集 dataset = ds.dataset(temp_parquet_dir, format="parquet") # 获取所有样本列名(排除Name列) sample_columns = [col for col in dataset.schema.names if col != "Name"] # 构建聚合规则:按Name分组,收集每个样本的Beta值(每个分组对应一个探针,每个样本列的聚合结果是单元素列表) aggregations = [pc.list_aggregate(pc.field(col), "list") for col in sample_columns] # 执行分组聚合(PyArrow会自动处理内存,超出内存时使用磁盘中间存储) grouped_table = dataset.group_by("Name").aggregate(aggregations) # 将单元素列表展开为实际的Beta值 for col in sample_columns: grouped_table = grouped_table.set_column( grouped_table.schema.get_field_index(col), col, pc.list_element(pc.field(col), 0) # 提取列表第一个元素 ) # 写入最终宽表(优先用Parquet,若需TSV可替换为pv.write_csv) pq.write_table( grouped_table, "methylation_wide_table.parquet", compression="snappy", row_group_size=100000 # 拆分row group,降低内存占用 )
方案二:用DuckDB简化宽表生成(更高效)
如果样本数极多(数千个),PyArrow的分组聚合可能仍有内存压力,可结合DuckDB(列式OLAP数据库)自动处理磁盘溢出,用SQL实现全Outer Join:
import duckdb import pyarrow.dataset as ds from pathlib import Path temp_parquet_dir = Path("./temp_methylation_parquet") dataset = ds.dataset(temp_parquet_dir, format="parquet") sample_columns = [col for col in dataset.schema.names if col != "Name"] # 连接DuckDB(内存不足时自动使用磁盘) con = duckdb.connect() # 动态构建全Outer Join的SQL语句 base_sample = sample_columns[0] sql = f"SELECT t1.Name, t1.{base_sample}" for idx, sample in enumerate(sample_columns[1:], start=2): sql += f", t{idx}.{sample}" sql += f" FROM '{temp_parquet_dir}/{base_sample}.parquet' t1" for idx, sample in enumerate(sample_columns[1:], start=2): sql += f" FULL OUTER JOIN '{temp_parquet_dir}/{sample}.parquet' t{idx} ON t1.Name = t{idx}.Name" # 执行SQL并导出为Arrow表 wide_table = con.execute(sql).arrow() # 写入输出 import pyarrow.parquet as pq pq.write_table(wide_table, "methylation_wide_table.parquet", compression="snappy")
关键注意事项
- 避免直接拼接TSV:TSV是行存格式,多次拼接会频繁复制数据,极易触发OOM;Parquet是列存,支持高效的列操作和压缩。
- 指定数据类型:解析TSV时明确
Name和Beta的类型,避免PyArrow自动推断带来的内存浪费和性能损耗。 - 清理临时文件:处理完成后可删除
temp_methylation_parquet目录,释放磁盘空间。
内容的提问来源于stack exchange,提问作者Ben
相关产品推荐
相关产品推荐

