【问题标题】:How does the kafka handle the messages sent to it while it is stopped?kafka 在停止时如何处理发送给它的消息?
【发布时间】:2020-09-09 20:57:48
【问题描述】:

我是 Kafka 新手,我一直在研究 Kafka 在停止时向其发送消息时的行为。

我面临的情况是我使用“Kubectl delete StatefulSet kafka_kf”来停止 Kafka。然后我使用 java Kafka Producer 向 Kafka 发送了一些消息。然后我再次启动 Kafka,这些发送给 Kafka 的消息在我启动 Kafka 的那一刻立即出现在消费者中。 知道在这种情况下卡夫卡内部会发生什么吗?以及如何防止这些消息出现在消费者中?这些消息稍后会导致重复问题,这就是为什么我需要它们不出现。

我通过使用命令打开的消费者看到消息出现在消费者中:

kubectl exec -ti test -- ./bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --isolation-level read_committed --topic testtopic

用于向 kafka 发送消息的代码和平是: producer.send(message)

【问题讨论】:

    标签: apache-kafka kafka-consumer-api kafka-producer-api


    【解决方案1】:

    首先,我认为了解producer.send() 是一个异步调用很重要,因此它不会阻塞。其次,send() 方法实际上并不将消息推送到代理,而是将消息放在本地内存中的二进制队列中。每个都有一个单独的二进制队列 划分生产者与之通信的主题。记录实际上是由生产者端的内部后台线程推送到代理的,该线程将由可配置的批处理阈值触发。等待来自代理的确认(由 acks 设置配置)的正是此操作,而不是 send() 方法。

    [来源:Confluent 培训 - 构建 Apache Kafka 的开发人员技能]

    当 Kafka 不可用时,您将在您的生产者中获得 TimeoutException。但是,这个异常可以通过重试来处理,生产者配置retries默认设置为2147483647。

    一旦您使 Kafka 可用,您的生产者就可以实际将消息发送到 Kafka,而您的消费者将收到它们。

    如果你不想接收这些消息你需要设置KafkaProducer配置retries=0

    要了解有关生产者回调异常的更多信息,您可以查看我的另一个 answer

    编辑评论中的新问题:

    有什么方法可以判断一条消息(或所有消息)是否发送成功?

    您可以在发送数据时定义如下所示的自定义回调类。如果消息的生成出现问题,此回调将抛出异常。

    class ProducerCallback extends Callback {
    
      @Override
      override def onCompletion(recordMetadata: RecordMetadata, e: Exception): Unit = {
        if (e != null) {
          e.printStackTrace()
        }
      }
    
    }
    
    producer.send(message, new ProducerCallback)
    

    作为替代方案,您可以简单地调用

    producer.send(message).get()
    

    因为这将阻塞,直到您收到来自 Kafka 代理的所有确认(请参阅 KafkaProducer 配置 acks)。

    【讨论】:

    • 嗨,迈克,感谢您的帮助,有什么方法可以查看消息(或所有消息)是否已成功发送?所以我可以根据producer.send(message)的结果继续编写代码谢谢
    • 谢谢,完成。你能在这里得到一些帮助吗:stackoverflow.com/questions/63829279/…
    猜你喜欢
    • 2018-07-30
    • 2015-05-16
    • 2015-06-13
    • 2019-12-06
    • 1970-01-01
    • 1970-01-01
    • 2019-06-24
    • 1970-01-01
    • 2021-12-05
    相关资源
    最近更新 更多