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

Spark Parquet批量写入生成多文件:如何输出单一Parquet文件

问题:批量写入Parquet生成多文件,如何高效生成单文件?

从数据库读取数据,通过ResultSet创建org.apache.spark.sql.Row后写入Parquet文件。因数据量较大,采用20万条的批量处理逻辑。原本想用coalesce(1)确保单分区以生成单个Parquet文件,但实际每个批次都会生成一个约2000KB的小文件,最终总大小不足100MB,但多文件会拖累读取性能。试过最后读取所有文件用repartition(1)重写,但担心影响写入性能。

原始批量写入代码

while (rs.next()) {
    Object[] fieldValues = new Object[columnCount];
    for (int i = 1; i <= columnCount; i++) {
        Object val = rs.getObject(i);
        if(val != null){
            fieldValues[i - 1] = val;
        }  
    }
    rows.add(RowFactory.create(fieldValues));

    // 达到批次大小则写入Parquet
    if (rows.size() >= batchSize) {
        Dataset<Row> batchDf = spark.createDataFrame(rows, schema);
        batchDf.coalesce(1).write().mode("append").parquet(parquetFileName);
        doesDataExist = true;
        rows.clear();
    }
}

// 写入剩余不足批次大小的数据
if (!rows.isEmpty()) {
    Dataset<Row> remainingDf = spark.createDataFrame(rows, schema);
    remainingDf.coalesce(1).write().mode("append").parquet(parquetFileName);
    doesDataExist = true;
}

尝试的合并代码

// 读取所有批次文件并重写为单文件
Dataset<Row> batchDf = spark.read().parquet(parquetFileName);
batchDf.repartition(1).write().mode("overwrite").parquet(parquetFileName);       

解决方案

1. 优先使用Spark原生JDBC读取(最优方案)

Spark内置支持JDBC数据源,可直接读取数据库表为DataFrame,无需手动处理ResultSet和分批逻辑,底层会优化数据读取与写入流程,配合coalesce(1)可直接生成单文件。

Dataset<Row> df = spark.read()
    .format("jdbc")
    .option("url", "jdbc:your_db_url")
    .option("dbtable", "target_table")
    .option("user", "username")
    .option("password", "password")
    .option("fetchsize", "200000") // 对应原逻辑的批量拉取大小
    .load();

df.coalesce(1).write().mode("overwrite").parquet(parquetFileName);

优势:代码简洁,性能更优,完全避免多文件问题。

2. 手动处理ResultSet时,避免分批写入

不要每个批次都写入Parquet,而是将所有Row缓存(内存不足时可分批缓存到磁盘),最后一次性生成DataFrame写入:

List<Row> allRows = new ArrayList<>();
List<Dataset<Row>> tempDfs = new ArrayList<>();
while (rs.next()) {
    Object[] fieldValues = new Object[columnCount];
    for (int i = 1; i <= columnCount; i++) {
        Object val = rs.getObject(i);
        if(val != null){
            fieldValues[i - 1] = val;
        }  
    }
    allRows.add(RowFactory.create(fieldValues));

    // 内存不足时,分批缓存到磁盘
    if (allRows.size() >= 1000000) { // 可根据内存调整阈值
        Dataset<Row> tempDf = spark.createDataFrame(allRows, schema);
        tempDf.persist(StorageLevel.MEMORY_AND_DISK());
        tempDfs.add(tempDf);
        allRows.clear();
    }
}

// 合并所有数据
Dataset<Row> finalDf = spark.createDataFrame(allRows, schema);
for (Dataset<Row> tempDf : tempDfs) {
    finalDf = finalDf.union(tempDf);
    tempDf.unpersist();
}

finalDf.coalesce(1).write().mode("overwrite").parquet(parquetFileName);

注意:若数据量远超内存,优先选择Spark JDBC方案,避免内存溢出。

3. 已生成多文件时的合并优化

如果已经生成了多个小文件,用coalesce(1)替代repartition(1),前者不会触发shuffle,性能更优:

Dataset<Row> batchDf = spark.read().parquet(parquetFileName);
batchDf.coalesce(1).write().mode("overwrite").parquet(parquetFileName);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 09:25:07