【问题标题】:How not to lose messages from Kafka when database is offline数据库离线时如何不丢失来自Kafka的消息
【发布时间】:2019-06-07 04:54:14
【问题描述】:

我正在开发微服务,它使用来自 Kaffka 的消息,然后处理这些消息并将输出存储到 MongoDB

我是 kafka 新手,遇到一些丢失消息的问题。

场景很简单:

如果 mongoDB 处于离线状态,微服务收到一条消息,然后尝试将输出保存到 Mongo,然后我收到错误消息,提示 mongo 离线并且消息丢失。

我的问题是 kafka 中有任何机制在这种情况下停止发送消息。应该在 Kafka 中手动提交偏移量吗?处理 Kafka 消费者错误的最佳实践是什么?

【问题讨论】:

    标签: java spring-boot apache-kafka microservices spring-kafka


    【解决方案1】:

    对于这种情况,您应该手动提交偏移量。仅当您的消息处理成功时才提交偏移量。你像下面那样提交它。但是您应该注意,消息具有 ttl,因此在 ttl 过去后,消息会自动从 kafka 代理中删除。

    consumer.commitSync(); 
    

    【讨论】:

      【解决方案2】:

      我认为与其手动提交,不如使用 Kafka Streams 和 Kafka Connect。管理两个系统之间的事务:Apache Kafka 和 MongoDB 可能不是那么容易,所以最好使用已经开发和测试过的工具(您可以阅读更多关于 Kafka Connect:https://kafka.apache.org/documentation/#connecthttps://docs.confluent.io/current/connect/index.html

      你的场景可能是这样的:

      【讨论】:

        【解决方案3】:

        您可以通过在MessageListenerContainer 上使用pauseresume 方法来做到这一点(但您必须使用spring kafka > 2.1.x)spring-kafka-docs

        @KafkaListener 生命周期管理

        @KafkaListener 注解创建的侦听器容器不是应用程序上下文中的bean。相反,它们使用KafkaListenerEndpointRegistry 类型的基础设施bean 进行注册。该 bean 由框架自动声明并管理容器的生命周期;它会自动启动任何将autoStartup 设置为true 的容器。

        所以 Autowire KafkaListenerEndpointRegistry 应用程序中的注册表端点

        @Autowired
        private KafkaListenerEndpointRegistry registry;
        

        从注册表spring-kafka-docs获取MessageListenerContainer

        public MessageListenerContainer getListenerContainer(java.lang.String id)
        

        返回具有指定 id 的 MessageListenerContainer,如果不存在这样的容器,则返回 null。

        参数:

        id - 容器的id

        MessageListenerContainer 上,您可以使用pauseresume 方法spring-kafka-docs

        默认 void pause()

        在下一次 poll() 之前暂停此容器。

        默认 void resume()

        如果暂停,则在下一次 poll() 之后恢复此容器。

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 1970-01-01
          • 2019-08-29
          • 1970-01-01
          • 2017-07-23
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2016-03-16
          相关资源
          最近更新 更多