如何在Apache Beam Java中实现按小时分区的Parquet文件写入
Java 实现 Apache Beam 按小时分区写入 Parquet(Kafka 数据源)
前置依赖
首先在 pom.xml 中引入核心依赖,以下为 2.49.0 版本示例,可根据实际需求调整版本:
<dependencies> <!-- Beam 核心依赖 --> <dependency> <groupId>org.apache.beam</groupId> <artifactId>beam-sdks-java-core</artifactId> <version>2.49.0</version> </dependency> <!-- Kafka IO 依赖 --> <dependency> <groupId>org.apache.beam</groupId> <artifactId>beam-sdks-java-io-kafka</artifactId> <version>2.49.0</version> </dependency> <!-- Parquet IO 依赖 --> <dependency> <groupId>org.apache.beam</groupId> <artifactId>beam-sdks-java-io-parquet</artifactId> <version>2.49.0</version> </dependency> <!-- Hadoop 兼容依赖,写入文件系统需要 --> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-common</artifactId> <version>2.8.5</version> <exclusions> <exclusion> <groupId>org.slf4j</groupId> <artifactId>slf4j-log4j12</artifactId> </exclusion> </exclusions> </dependency> </dependencies>
注意要保证 Beam 版本和 Hadoop 版本兼容,避免类冲突问题
核心实现步骤
1. 定义数据结构与Schema
推荐使用 Avro Schema 适配 Parquet 存储,你可以先定义 Avro schema 文件,或者直接用代码生成 Schema 示例:
// 示例数据结构,包含时间戳字段 ts(毫秒级时间戳)、id、content Schema avroSchema = SchemaBuilder.record("EventData") .namespace("com.example.beam") .fields() .requiredLong("ts") .requiredString("id") .requiredString("content") .endRecord();
2. 读取Kafka数据并分配事件时间
从 Kafka 消费数据后反序列化为 Avro GenericRecord,同时提取数据中的时间戳作为事件时间,用于后续窗口划分:
Pipeline pipeline = Pipeline.create(options); pipeline // 读取Kafka数据 .apply(KafkaIO.<String, GenericRecord>read() .withBootstrapServers("kafka-broker:9092") .withTopic("your-topic-name") .withKeyDeserializer(StringDeserializer.class) .withValueDeserializer(ConfluentSchemaRegistryDeserializer.of(avroSchema, "http://schema-registry:8081")) .withoutMetadata()) // 提取数据中的时间戳作为事件时间 .apply(WithTimestamp.of(record -> new Instant((Long) record.get("ts"))));
3. 按小时划分固定窗口
将数据按小时维度划分窗口,实现每小时生成分区的逻辑:
.apply(Window.into(FixedWindows.of(Duration.standardHours(1))) // 可根据业务需求配置允许迟到时间 .withAllowedLateness(Duration.standardMinutes(10)) .discardingFiredPanes());
4. 动态写入分区Parquet文件
使用 FileIO.writeDynamic 按小时生成动态路径,写入Parquet文件:
.apply(FileIO.<Instant, GenericRecord>writeDynamic() // 按窗口的结束时间作为分区键,生成分区路径格式 yyyy/MM/dd/HH .by((record, windowCtx) -> windowCtx.window().maxTimestamp()) .withDestinationCoder(InstantCoder.of()) .via(ParquetIO.sink(avroSchema).withCompressionCodec(CompressionCodecName.SNAPPY)) // 输出根路径 .to("/your/output/path/parquet") // 分区路径生成规则 .withNaming(key -> FileIO.Write.defaultNaming( String.format("dt=%s/", DateTimeFormat.forPattern("yyyy/MM/dd/HH").print(key)), ".parquet" )) // 每个分区的文件分片数量,避免小文件过多 .withNumShards(2));
注意事项
- 时间戳提取逻辑需要和你数据内的时间戳格式匹配,如果是字符串格式的时间需要先转换为毫秒级时间戳再赋值给事件时间
- 可根据业务需求调整
withNumShards的值,平衡文件数量和文件大小 - 如果需要写入对象存储(OSS/S3/GCS),只需要修改
to方法的输出路径前缀即可,Beam 会自动适配对应的存储系统 - 生产环境建议开启 Parquet 压缩,降低存储成本
内容的提问来源于stack exchange,提问作者manoveg
相关产品推荐
相关产品推荐

