【发布时间】:2021-07-30 18:56:45
【问题描述】:
我正在使用从 Kafka 到 HDFS 的 Flink bucketing sink。 Flink 的版本是 1.4.2。
我发现每次重新启动作业时都会丢失一些数据,即使使用保存点也是如此。
我发现如果我设置 writer SequenceFile.CompressionType.RECORD 而不是 SequenceFile.CompressionType.BLOCK 可以解决这个问题。似乎在 Flink 尝试保存检查点时,有效长度与实际长度不同,其中应该包括压缩数据。
但是,如果由于磁盘使用情况而无法使用 CompressionType.BLOCK,则可能会出现问题。如何在重新启动作业时使用块压缩来防止数据丢失?
这是 Flink 的已知问题吗?或者有谁知道如何解决这个问题?
【问题讨论】:
标签: hadoop hdfs apache-flink