【问题标题】:Spring cloud stream - Kafka consumer consuming duplicate messages with StreamListenerSpring Cloud Stream - Kafka 消费者使用 StreamListener 消费重复消息
【发布时间】:2022-02-08 16:20:19
【问题描述】:

使用我们的 Spring Boot 应用程序,我们注意到 kafka 消费者偶尔会在 prod 环境中随机消费两次消息。我们在 PCF 中部署了 6 个实例和 6 个分区。我们捕获了具有相同偏移量的消息,并且在同一主题中收到了两次分区,这会导致重复,这对我们来说是关键业务。 我们在非生产环境中没有注意到这一点,并且在非生产环境中很难重现。我们最近切换到 Kafka,我们无法找出根本问题。

我们正在使用 spring-cloud-stream/spring-cloud-stream-binder-kafka- 2.1.2 这是配置:

spring:
  cloud:
    stream:
      default.consumer.concurrency: 1 
      default-binder: kafka
      bindings:
        channel:
          destination: topic
          content_type: application/json
          autoCreateTopics: false
          group: group
          consumer:
            maxAttempts: 1
      kafka:
        binder:
          autoCreateTopics: false
          autoAddPartitions: false
          brokers: brokers list
        bindings:
          channel:
            consumer:
              autoCommitOnError: true
              autoCommitOffset: true
              configuration:
                max.poll.interval.ms: 1000000
                max.poll.records: 1 
                group.id: group

我们使用@Streamlisteners 来消费消息。

这是我们收到的重复实例和服务器日志中捕获的错误消息。

错误 46 --- [container-0-C-1] o.a.k.c.c.internals.ConsumerCoordinator : [Consumer clientId=consumer-3, groupId=group] 偏移提交失败 在偏移量 1291358 的分区 topic-0 上:协调器不知道 这个成员的。错误 46 --- [容器-0-C-1] os.kafka.listener.LoggingErrorHandler :处理时出错: 空 OUT org.apache.kafka.clients.consumer.CommitFailedException: 提交无法完成,因为组已经重新平衡并且 将分区分配给另一个成员。这意味着时间 对 poll() 的后续调用之间的时间比配置的长 max.poll.interval.ms,这通常意味着轮询循环是 花费太多时间处理消息。你可以解决这个问题 通过增加会话超时或减少最大值 使用 max.poll.records 在 poll() 中返回的批次大小。 在 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator$OffsetCommitResponseHandler.handle(ConsumerCoordinator.java:871) ~[kafka-clients-2.0.1.jar!/:na]

复制时没有崩溃,所有实例都运行良好。错误日志也存在混淆-处理时出错:null,因为消息已成功处理两次。而 max.poll.interval.ms: 100000 大约是 16 分钟,应该有足够的时间来处理系统的任何消息,并且会话超时和 heartbit 配置是默认设置。在大多数情况下,会在 2 秒内收到副本。 我们缺少任何配置吗?非常感谢任何建议/帮助。

【问题讨论】:

  • 通过将 session.timeout.ms 更改为 45 秒(默认为 10 秒)解决了这个问题

标签: spring-boot apache-kafka spring-kafka spring-cloud-stream


【解决方案1】:

由于组已重新平衡,因此无法完成提交

发生了重新平衡,因为您的侦听器花费了太长时间;您应该调整max.poll.recordsmax.poll.interval.ms 以确保您始终可以处理在时限内收到的记录。

在任何情况下,Kafka 不保证完全一次交付,只保证至少一次交付。您需要为您的应用程序添加幂等性并检测/忽略重复项。

【讨论】:

  • 正如我上面提到的,我们有 max.poll.interval.ms = 16 分钟和 max.poll.records =1 。我们可以看到消息完成该过程的时间不超过 5 秒。不知道我们能在这方面做得更好。除了手动检查偏移量之外,还有其他方法可以使消费者幂等吗?我们已经在生产者端启用了幂等性。
  • 这毫无意义。显然发生了再平衡。检查服务器日志。
  • 我在上面附上了服务器日志简而言之 - 协调器不知道这个成员,处理时出错:null 提交无法完成,因为该组已经重新平衡并分配了分区等等。
  • 你需要早点看。显然这不是 Spring 问题。但是,您使用的是非常旧的版本。
  • 支持的最旧的 OSS 版本是 3.2 spring.io/projects/spring-cloud-stream#support
【解决方案2】:

另外,请记住 StreamListener 和基于注解的编程模型已被弃用 3 年多,并且已从当前主版本中删除,这意味着下一个版本将不会有它。请将您的解决方案迁移到functional based programming model

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2016-06-21
    • 2017-06-22
    • 1970-01-01
    • 2022-01-11
    • 2018-03-09
    • 2020-11-02
    • 2020-03-28
    • 2018-01-04
    相关资源
    最近更新 更多