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
相关产品推荐
相关产品推荐

