【问题标题】:Creating a parquet file with StreamingFileSink in java在 java 中使用 StreamingFileSink 创建镶木地板文件
【发布时间】: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


    【解决方案1】:

    您应该使用bulk-encoded-format 来编写 Parquet。 RowFormat用于写入文本、csv、json等。

    【讨论】:

    • 非常感谢您的回归。我需要使用 with .RollingPolicy 部分,但我发现这对于 .forBulkFormat 是不可能的。你有这方面的经验吗?
    • 嗨@Emsal。正如我们在文档中看到的(上面的链接):批量格式只能有OnCheckpointRollingPolicy,它(仅)在每个检查点上滚动。
    猜你喜欢
    • 1970-01-01
    • 2021-10-14
    • 2016-10-07
    • 1970-01-01
    • 2022-01-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-09-23
    相关资源
    最近更新 更多