【发布时间】:2021-12-17 22:38:51
【问题描述】:
我的目标是将我从 kafka 收到的消息转换为 parquet 文件,但我可能错了。你能帮我解决这个问题吗?
private static SinkFunction<String> createFileSink(String outputPath) {
final StreamingFileSink<String> sink = StreamingFileSink
.forRowFormat(new Path(outputPath), new SimpleStringEncoder<String>("UTF-8"))
.withRollingPolicy(
DefaultRollingPolicy.builder()
.withRolloverInterval(TimeUnit.MINUTES.toMillis(15))
.withInactivityInterval(TimeUnit.MINUTES.toMillis(5))
.withMaxPartSize(1024 * 1024)
.build())
.build();
return sink;
}
【问题讨论】:
标签: java architecture apache-flink software-design