使用Java中StreamingFileSink创建Parquet文件及Kafka消息转存咨询
现有方案问题分析
你当前的实现存在两个核心问题,完全无法满足输出Parquet的需求:
- 你使用的是
forRowFormat行格式输出,搭配的SimpleStringEncoder是纯文本编码,只能输出普通字符串文本文件,不支持Parquet列式存储格式 - 缺少Parquet序列化所需的schema定义、数据类型转换逻辑,Parquet是结构化存储格式,不能直接写入字符串类型的Kafka原始消息
正确实现方案
前置依赖
首先需要在项目中引入Flink Parquet相关依赖(以Maven为例,版本对应你使用的Flink版本即可):
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-parquet_${scala.binary.version}</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.parquet</groupId> <artifactId>parquet-hadoop</artifactId> <version>1.12.2</version> </dependency>
实现步骤
- 先定义数据结构和Parquet Schema:你需要先将Kafka接收到的字符串消息反序列化为结构化的POJO类,再根据POJO定义对应的Parquet Schema。如果你的消息是JSON格式,可以直接用Avro或者Protobuf定义结构,也可以手动编写Parquet Schema。
- 替换为Parquet Bulk格式的Sink:StreamingFileSink要输出Parquet需要使用
forBulkFormat方法,搭配ParquetWriterFactory实现。
示例代码
假设你已经定义了名为MessagePOJO的结构化类,并且已经将Kafka的字符串流转换为DataStream<MessagePOJO>,Parquet Sink的实现如下:
import org.apache.flink.formats.parquet.ParquetWriterFactory; import org.apache.flink.formats.parquet.avro.ParquetAvroWriters; import org.apache.flink.streaming.api.functions.sink.filesystem.StreamingFileSink; import org.apache.flink.streaming.api.functions.sink.filesystem.bucketassigners.DateTimeBucketAssigner; import org.apache.hadoop.fs.Path; import java.util.concurrent.TimeUnit; // 注意这里泛型改为你实际的结构化POJO类型 private static SinkFunction<MessagePOJO> createParquetFileSink(String outputPath) { // 如果用Avro生成的类,直接用这个方法创建Writer工厂,自定义POJO可以用ParquetPojoWriters ParquetWriterFactory<MessagePOJO> writerFactory = ParquetAvroWriters.forSpecificRecord(MessagePOJO.class); final StreamingFileSink<MessagePOJO> sink = StreamingFileSink .forBulkFormat(new Path(outputPath), writerFactory) // 分桶策略,这里按小时分桶,可根据需求调整 .withBucketAssigner(new DateTimeBucketAssigner<>("yyyy-MM-dd/HH")) .withRollingPolicy( DefaultRollingPolicy.builder() .withRolloverInterval(TimeUnit.MINUTES.toMillis(15)) .withInactivityInterval(TimeUnit.MINUTES.toMillis(5)) .withMaxPartSize(128 * 1024 * 1024) // Parquet建议块大小128MB .build()) .build(); return sink; }
额外注意点
- 如果你不需要使用Avro,也可以直接使用
ParquetPojoWriters.forReflectRecord(MessagePOJO.class)直接为普通Java POJO生成Parquet写入器,不需要额外定义Avro schema - 原始Kafka字符串消息必须先完成反序列化和数据清洗,转换为结构化的POJO实例之后才能写入Parquet,直接写入字符串会导致Parquet结构异常
- Parquet属于批量写入的列式格式,滚动策略的单文件大小建议调整为64MB~256MB之间,不要用你原来的1MB,过小的文件会导致大量小文件问题,影响后续查询性能
内容的提问来源于stack exchange,提问作者Emsal Cengiz
相关产品推荐
相关产品推荐

