【问题标题】:Apache Flink: KafkaProducer Data LossApache Flink:KafkaProducer 数据丢失
【发布时间】:2018-12-13 22:33:16
【问题描述】:

我们的 Flink 流式工作流将消息发布到 Kafka。 KafkaProducer 的“重试”机制在将消息添加到其内部缓冲区之前不会启动。

如果在此之前有异常,KafkaProducer 将抛出该异常,并且似乎 Flink 没有处理该异常。在这种情况下,将会丢失数据。

相关Flink代码(FlinkKafkaProducerBase):

if (logFailuresOnly) {
            callback = new Callback() {
                @Override
                public void onCompletion(RecordMetadata metadata, Exception e) {
                    if (e != null) {
                        LOG.error("Error while sending record to Kafka: " + e.getMessage(), e);
                    }
                    acknowledgeMessage();
                }
            };
        }
        else {
            callback = new Callback() {
                @Override
                public void onCompletion(RecordMetadata metadata, Exception exception) {
                    if (exception != null && asyncException == null) {
                        asyncException = exception;
                    }
                    acknowledgeMessage();
                }
            };
        }

以下是我们已确定会导致数据丢失的场景:

  1. 所有的 kafka 代理都挂了。

    在这种情况下,在将消息附加到其缓冲区之前,KafkaProducer 会尝试获取元数据。如果 KafkaProducer 在配置的超时时间内无法获取元数据,则会引发异常。

  2. -内存记录不可写(kafka 0.9.0.1 库中存在错误)

https://issues.apache.org/jira/browse/KAFKA-3594

以上两种情况,KafkaProducer 都不会重试,Flink 会忽略这些消息。甚至没有记录消息。例外是,但不是失败的消息。

可能的解决方法(Kafka 设置):

  1. 元数据超时值非常高 (metadata.fetch.timeout.ms)
  2. 缓冲区过期值非常高 (request.timeout.ms)

我们仍在调查更改上述 kafka 设置可能产生的副作用。

那么,我们的理解正确吗?或者有没有办法通过修改一些 Flink 设置来避免这种数据丢失?

谢谢。

【问题讨论】:

    标签: apache-kafka apache-flink producer


    【解决方案1】:

    这就是我对您的问题的看法。 首先查看 Kafka 保证之一:

    对于复制因子为 N 的主题,我们最多可以容忍 N-1 个服务器故障而不会丢失任何提交到日志的记录。

    首先,它关心提交到日志的消息或记录。任何未能交付的记录都不会被视为已提交。其次,如果你所有的经纪人都挂了,会有一些数据丢失。

    以下设置是我们用来防止生产者端数据丢失的设置:

    • block.on.buffer.full = true
    • acks = 全部
    • 重试次数 = MAX_VALUE
    • max.in.flight.requests.per.connection = 1
    • 使用 KafkaProducer.send(record, callback) 代替 send(record)
    • unclean.leader.election.enable=false
    • replication.factor > min.insync.replicas
    • min.insync.replicas > 1

    【讨论】:

    • 嗨紫水晶。感谢您的答复。是的,我们有上述所有设置。因此,即使所有经纪人都关闭了,我们也不应该丢失数据。这实际上是 Flink 中的一个 bug。请查收:apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/…
    • 在版本 (
    • 这些设置来自新的生产者。
    • @Ninad 我阅读了上面提到的邮件线程,但没有看到结论。这仍然是一个问题,还是在 Flink 的某些更高版本中解决了?
    猜你喜欢
    • 2023-03-25
    • 1970-01-01
    • 1970-01-01
    • 2016-07-05
    • 2023-04-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多