如何基于Java、Apache Arrow与Parquet实现高效分区数据摄入方案
本地Parquet分区数据集写入实现方案(适配Java 11+、DataFusion查询场景)
核心逻辑说明
Parquet-mr本身没有内置分区数据集写入API——Parquet格式仅定义单文件的存储结构,分区数据集是Hive、Spark、DataFusion等查询引擎共同遵循的Hive风格目录约定,不属于Parquet底层格式的能力范畴,因此不需要寻找专门的分区写入接口,你当前已经验证过性能达标的AvroParquetWriter可以直接复用,无需替换组件。
分区落地的核心逻辑是:给每条记录提取分区字段值,将同分区的记录写入对应独立目录下的Parquet文件即可。DataFusion会自动识别这种目录结构,查询时直接跳过不符合过滤条件的分区目录,实现分区剪枝提速。
针对你99万条/秒的IoT采集场景,注意时间分区不要直接用原始采集的毫秒/微秒时间戳作为分区值,否则会产生海量小目录,反而拖慢查询效率,推荐按小时粒度做时间分区(如果单传感器每小时数据量超过10GB,可调整为15分钟粒度),最终目录结构如下:
数据存储根目录/ ├── ts_hour=2024-05-20-10/ │ ├── sensor=vibration_sensor_01/ │ │ ├── part-1716199200000-a1b2c3.parquet │ │ └── part-1716200000000-d4e5f6.parquet │ └── sensor=temperature_sensor_03/ │ └── part-1716199200000-g7h8i9.parquet └── ts_hour=2024-05-20-11/ └── sensor=pressure_sensor_02/ └── part-1716202800000-j0k1l2.parquet
具体代码实现
整体实现只需要在现有写入逻辑上层加一层分区路由和Writer缓存,不需要修改原有Avro序列化、Parquet压缩编码等已经验证过性能的配置。核心逻辑分为4步:
- 每条记录进入后,提取两个分区字段的值:将原始时间戳按选定粒度格式化成分区字符串,对传感器名做路径非法字符转义
- 拼接分区目录路径,路径不存在时递归创建
- 用线程安全的Map维护分区路径到ParquetWriter实例的缓存,不存在对应Writer时就创建新的Parquet文件和Writer实例
- 单文件大小达到128MB-256MB阈值时,关闭当前Writer,下次写入同分区数据时自动创建新文件,避免产生过大或过小的文件
核心Java 11代码示例:
import org.apache.avro.generic.GenericRecord; import org.apache.hadoop.fs.Path; import org.apache.parquet.avro.AvroParquetWriter; import org.apache.parquet.hadoop.ParquetWriter; import org.apache.parquet.hadoop.metadata.CompressionCodecName; import java.io.IOException; import java.nio.file.Files; import java.time.Instant; import java.time.ZoneOffset; import java.time.format.DateTimeFormatter; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; public class PartitionedParquetWriter { private final Path dataLakeRoot; private final org.apache.avro.Schema avroSchema; private final ConcurrentHashMap<String, ParquetWriter<GenericRecord>> writerCache = new ConcurrentHashMap<>(); private final ConcurrentHashMap<String, Long> writerFileSizeCounter = new ConcurrentHashMap<>(); private static final long FILE_ROLL_THRESHOLD = 128 * 1024 * 1024; // 单文件128MB滚动 private static final DateTimeFormatter TS_PARTITION_FORMAT = DateTimeFormatter.ofPattern("yyyy-MM-dd-HH").withZone(ZoneOffset.UTC); public PartitionedParquetWriter(String rootPath, org.apache.avro.Schema schema) throws IOException { this.dataLakeRoot = new Path(rootPath); this.avroSchema = schema; Files.createDirectories(java.nio.file.Path.of(rootPath)); // 注册JVM关闭钩子,异常退出时自动关闭所有Writer避免文件损坏 Runtime.getRuntime().addShutdownHook(new Thread(this::closeAllWriters)); } public void write(GenericRecord record) throws IOException { // 提取分区值 long rawTimestamp = (long) record.get("timestamp"); String tsPartition = TS_PARTITION_FORMAT.format(Instant.ofEpochMilli(rawTimestamp)); String sensorName = record.get("sensor_name").toString() .replaceAll("[\\\\/:*?\"<>|]", "_"); // 转义路径非法字符 // 拼接分区路径 String partitionRelPath = String.format("ts_hour=%s/sensor=%s", tsPartition, sensorName); Path partitionDir = new Path(dataLakeRoot, partitionRelPath); Files.createDirectories(java.nio.file.Path.of(partitionDir.toUri())); // 获取或创建对应分区的Writer ParquetWriter<GenericRecord> writer = writerCache.get(partitionRelPath); if (writer == null) { String fileName = String.format("part-%d-%s.parquet", System.currentTimeMillis(), UUID.randomUUID()); Path parquetFilePath = new Path(partitionDir, fileName); // 这里的压缩、编码等参数全部复用你现有已经验证过性能的配置即可 writer = AvroParquetWriter.<GenericRecord>builder(parquetFilePath) .withSchema(avroSchema) .withCompressionCodec(CompressionCodecName.ZSTD) .build(); writerCache.put(partitionRelPath, writer); writerFileSizeCounter.put(partitionRelPath, 0L); } // 写入记录 writer.write(record); // 估算记录大小,达到阈值滚动文件 long currentSize = writerFileSizeCounter.addAndGet(partitionRelPath, estimateRecordSize(record)); if (currentSize >= FILE_ROLL_THRESHOLD) { writer.close(); writerCache.remove(partitionRelPath); writerFileSizeCounter.remove(partitionRelPath); } } public void closeAllWriters() { writerCache.values().forEach(w -> { try { w.close(); } catch (IOException e) { // 按你的业务场景处理异常 } }); writerCache.clear(); writerFileSizeCounter.clear(); } // 简单估算单条Avro记录的序列化后大小,可根据你的数据结构调整估算逻辑 private long estimateRecordSize(GenericRecord record) { return 48; // 示例值,按实际单条测量值平均大小调整即可,不需要100%精准 } }
性能与查询优化注意事项
- 分区字段顺序保持「时间分区在前、传感器分区在后」即可,符合IoT场景绝大多数查询带时间范围过滤的特点,DataFusion可以在第一层目录扫描时就剪枝掉所有不满足时间范围的分区,剪枝效率最高
- 不需要把分区字段从Avro Schema中移除,Parquet文件中保留分区字段不会影响查询正确性,DataFusion读取时会自动用目录中的分区值覆盖文件内的同名字段,省得修改现有数据结构
- 不要随意替换现有AvroParquetWriter实现:你已经验证过当前写入性能满足99万条/秒的要求,换用其他Parquet写入实现(比如Arrow Parquet Writer)反而需要额外做跨格式数据拷贝,会降低写入吞吐量
- AvroParquetWriter默认会开启行组级别的统计信息收集,不需要手动配置,DataFusion扫描Parquet文件时可以借助这些统计信息直接跳过不满足过滤条件的行组,进一步提升查询速度
- 如果服务重启频繁导致同分区下产生大量小文件,可以在业务低峰期跑后台定时任务合并同分区的小文件,不影响正常写入流程
DataFusion侧读取配置
读取时直接指向数据存储根目录即可,不需要手动枚举所有分区文件,DataFusion会自动识别Hive风格的分区结构:
# Python DataFusion示例 import datafusion ctx = datafusion.SessionContext() # 传入分区字段顺序和写入时保持一致即可 df = ctx.read_parquet("/path/to/your/datalake/root", partition_cols=["ts_hour", "sensor"])
内容的提问来源于stack exchange,提问作者João Paraná
相关产品推荐
相关产品推荐

