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

如何基于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步:

  1. 每条记录进入后,提取两个分区字段的值:将原始时间戳按选定粒度格式化成分区字符串,对传感器名做路径非法字符转义
  2. 拼接分区目录路径,路径不存在时递归创建
  3. 用线程安全的Map维护分区路径到ParquetWriter实例的缓存,不存在对应Writer时就创建新的Parquet文件和Writer实例
  4. 单文件大小达到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á

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 09:30:45