【问题标题】:Kafka Fails to Process all the messages - Java Spring BootKafka 无法处理所有消息 - Java Spring Boot
【发布时间】:2020-04-27 08:55:51
【问题描述】:

我有一个 Spring Boot 应用程序(Spring 版本 2.2.2.RELEASE),我在其中配置了 Kafka 消费者,它处理来自 Kafka 的数据并服务于多个 Web 套接字。订阅 kafka 成功,但并非选定 Kafka 主题的所有消息都由消费者处理。很少有消息被延迟,很少有消息被错过。但是生产者正在发送完全确保的数据。下面我分享了我使用过的配置属性。

@Bean
public ConsumerFactory<String, String> consumerFactory() {
    final String BOOTSTRAP_SERVERS = kafkaBootstrapServer;
    Map<String, Object> props = new HashMap<>();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS);
    props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroupId);
    return new DefaultKafkaConsumerFactory<>(props);
}

我是否缺少任何配置?

【问题讨论】:

  • 您如何知道数据正在到达您在 UI 中查看的主题?

标签: spring-boot apache-kafka kafka-consumer-api spring-kafka


【解决方案1】:

对于新的消费者(从不提交 group.id 的偏移量),您必须将 AUTO_OFFSET_RESET 设置为 earliest 以避免丢失主题中的任何现有记录(默认为 latest)。

【讨论】:

  • 我不需要现有的消息,但我需要从启动应用程序的那一刻起发送的所有消息。在这种情况下,将 AUTO_OFFSET_RESET 设置为最早的帮助吗?
  • 如果消费者 group.id 已经有一个提交的偏移量,这不会有任何区别。它仅在 Consumer 第一次启动时应用(直到它提交至少一个偏移量)。
  • 消费者不会“错过”记录。
  • 好的,我正在对消息进行一些处理,此处理每条消息大约需要 4-5 秒。在这种情况下,在我的处理完成之前,下一条消息是否会等待处理?或者它将在另一个线程上?在使用消息后将此处理移动到单独的线程中是个好主意吗?
  • 假设您使用的是spring-kafka,这取决于您的并发设置以及主题中有多少个分区。同一分区内的所有消息将在同一线程上处理。您必须max.poll.interval.ms 内处理 max.poll.records 以避免重新平衡。
猜你喜欢
  • 2018-07-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-04-11
  • 2017-02-09
  • 1970-01-01
  • 2021-11-19
  • 2020-08-21
相关资源
最近更新 更多