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

如何解决Arrow write_dataset生成大量小Parquet分区文件的问题

解决Arrow Parquet分区大量小文件问题的方案

针对你遇到的每个Parquet分区内存在大量小文件的问题,可通过以下几种方案优化,均无需一次性加载全部数据至内存:

方案一:写入前按分区键排序(从源头避免小文件)

利用Arrow的延迟计算特性,先对数据集按country和type排序,确保相同分区的数据连续被处理,写入时Arrow会自动合并同分区的连续数据为大文件,而非每个原始txt块生成一个小文件。

# 对数据集按分区键排序(延迟计算,不加载全量数据)
sorted_data <- mydata |> 
  arrange(country, type)

# 写入排序后的数据集
write_dataset(dataset = sorted_data,
              path = path_parquet,
              format = "parquet",
              partitioning = c("country", "type"),
              compression = "snappy",
              use_dictionary = TRUE,
              write_statistics = TRUE,
              # 设定每个文件的最大行数,可根据数据情况调整(如100万行)
              max_rows_per_file = 1000000,
              existing_data_behavior = "delete_matching")

方案二:合并已生成的分区小文件

若已生成大量小文件的Parquet数据集,可直接遍历每个分区,读取该分区所有小文件数据后重新写入,合并为大文件。

# 读取现有Parquet数据集
existing_ds <- open_dataset(path_parquet)

# 获取所有分区组合
partitions <- existing_ds$GetPartitions()

# 遍历每个分区合并文件
for (part in partitions) {
  # 读取当前分区的所有数据
  part_ds <- existing_ds |> filter(!!!part$exprs)
  
  # 定义当前分区的输出路径
  part_path <- file.path(path_parquet, part$path)
  
  # 重新写入当前分区,替换原有小文件
  write_dataset(dataset = part_ds,
                path = part_path,
                format = "parquet",
                compression = "snappy",
                use_dictionary = TRUE,
                write_statistics = TRUE,
                max_rows_per_file = 1000000,
                existing_data_behavior = "delete_matching")
}

方案三:使用repartition函数快速合并

Arrow的repartition函数可直接对数据集重新分区,按指定列聚合数据并分割为指定大小的文件,快速合并小文件。

# 读取现有数据集
existing_ds <- open_dataset(path_parquet)

# 重新分区:按country和type聚合,设定每个文件最大行数
repartitioned_ds <- existing_ds |> 
  repartition(by = c("country", "type"), max_rows_per_file = 1000000)

# 写入重新分区后的数据集(可指定新路径或覆盖原路径)
write_dataset(dataset = repartitioned_ds,
              path = path_parquet_new,
              format = "parquet",
              compression = "snappy",
              use_dictionary = TRUE,
              write_statistics = TRUE,
              existing_data_behavior = "delete_matching")

注意事项

  • max_rows_per_file的数值可根据数据类型和存储需求调整,通常100万~500万行对应几十MB到几百MB的Parquet文件,符合Arrow的最佳实践。
  • 方案一从源头避免小文件问题,若原始数据处理时间允许,优先选用;方案二、三适合优化已生成的小文件数据集。
  • 排序操作基于Arrow的延迟计算,无需加载全量数据至内存,仅在写入时逐步处理排序后的数据。

内容的提问来源于stack exchange,提问作者Jakob

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 13:26:00