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

如何将多份甲基化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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 01:25:07