【问题标题】:Spring-Kafka Concurrency PropertySpring-Kafka 并发属性
【发布时间】:2019-09-01 18:58:21
【问题描述】:

我正在使用 Spring-Kafka 编写我的第一个 Kafka Consumer。查看了框架提供的不同选项,并且对此几乎没有疑问。如果您已经研究过,有人可以在下面澄清一下。

问题 - 1:根据 Spring-Kafka 文档,有 2 种方法可以实现 Kafka-Consumer; “您可以通过配置 MessageListenerContainer 并提供消息侦听器或使用 @KafkaListener 注释来接收消息”。谁能告诉我什么时候应该选择一个选项而不是另一个选项?

问题 - 2:我选择了 KafkaListener 方法来编写我的应用程序。为此,我需要初始化一个容器工厂实例,并且在容器工厂内部可以选择控制并发性。只是想仔细检查一下我对并发的理解是否正确。

假设,我有一个主题名称 MyTopic,其中有 4 个分区。为了使用来自 MyTopic 的消息,我启动了我的应用程序的 2 个实例,这些实例是通过将并发设置为 2 来启动的。因此,理想情况下,根据 kafka 分配策略,2 个分区应该转到 consumer1,2 个其他分区应该转到 consumer2 .由于并发设置为2,是否每个消费者将启动2个线程,并会并行消费来自主题的数据?如果我们并行消费,我们还应该考虑什么。

问题 3 - 我选择了手动确认模式,而不是在外部管理偏移量(不将其持久化到任何数据库/文件系统)。那么我是否需要编写自定义代码来处理重新平衡,或者框架会自动管理它?我认为没有,因为我只有在处理完所有记录后才承认。

问题 - 4 : 另外,在手动 ACK 模式下,哪个 Listener 会提供更好的性能? BATCH 消息侦听器或普通消息侦听器。我想如果我使用普通消息侦听器,则在处理每条消息后都会提交偏移量。

粘贴下面的代码供您参考。

批量确认消费者

    public void onMessage(List<ConsumerRecord<String, String>> records, Acknowledgment acknowledgment,
          Consumer<?, ?> consumer) {
      for (ConsumerRecord<String, String> record : records) {
          System.out.println("Record : " + record.value());
          // Process the message here..
          listener.addOffset(record.topic(), record.partition(), record.offset());
       }
       acknowledgment.acknowledge();
    }

初始化容器工厂:

@Bean
public ConsumerFactory<String, String> consumerFactory() {
    return new DefaultKafkaConsumerFactory<String, String>(consumerConfigs());
}

@Bean
public Map<String, Object> consumerConfigs() {
    Map<String, Object> configs = new HashMap<String, Object>();
    configs.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootStrapServer);
    configs.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
    configs.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, enablAutoCommit);
    configs.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, maxPolInterval);
    configs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset);
    configs.put(ConsumerConfig.CLIENT_ID_CONFIG, clientId);
    configs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    configs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    return configs;
}

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<String, String>();
    // Not sure about the impact of this property, so going with 1
    factory.setConcurrency(2);
    factory.setBatchListener(true);
    factory.getContainerProperties().setAckMode(AckMode.MANUAL);
    factory.getContainerProperties().setConsumerRebalanceListener(RebalanceListener.getInstance());
    factory.setConsumerFactory(consumerFactory());
    factory.getContainerProperties().setMessageListener(new BatchAckConsumer());
    return factory;
}

【问题讨论】:

    标签: spring apache-kafka spring-kafka


    【解决方案1】:

    第一季度:

    从文档中,

    @KafkaListener 注解用于将 bean 方法指定为 侦听器容器的侦听器。豆子被包裹在一个 MessagingMessageListenerAdapter 配置了各种功能,例如 作为转换器来转换数据,如有必要,以匹配方法 参数。

    您可以使用 SpEL 配置注释上的大多数属性 “#{…​} 或属性占位符 (${…​})。有关更多信息,请参阅 Javadoc。”

    这种方法对于简单的 POJO 侦听器很有用,您不需要实现任何接口。您还可以使用注释以声明性方式侦听任何主题和分区。您还可以潜在地返回您收到的值,而在 MessageListener 的情况下,您会受到接口签名的约束。

    第二季度:

    理想情况下是的。如果您有多个主题可供使用,那么它会变得更加复杂。默认情况下,Kafka 使用 RangeAssignor,它有自己的行为(您可以更改它——查看更多详细信息 under)。

    第三季度:

    如果您的消费者死亡,就会出现再平衡。如果您手动确认并且您的消费者在提交偏移量之前死亡,您不需要做任何事情,Kafka 会处理。但你最终可能会收到一些重复的消息(至少一次)

    第四季度:

    这取决于您所说的“性能”。如果您的意思是延迟,那么尽可能快地使用每条记录将是可行的方法。如果要实现高吞吐量,那么批量消费效率更高。

    我使用 Spring kafka 和各种侦听器编写了一些示例 - 查看 this repo

    【讨论】:

      【解决方案2】:
      1. @KafkaListener 是一个消息驱动的“POJO”,它添加了负载转换、参数匹配等内容。如果你实现了MessageListener,你只能从 Kafka 获取原始的ConsumerRecord。见@KafkaListener Annotation

      2. 是的,并发代表线程数;每个线程创建一个Consumer;它们并行运行;在您的示例中,每个将获得 2 个分区。

      如果我们并行消费,我们还应该考虑什么。

      您的侦听器必须是线程安全的(没有共享状态或任何此类状态需要被锁保护。

      1. 不清楚您所说的“处理重新平衡事件”是什么意思。当发生再平衡时,框架将提交所有未决的偏移量。

      2. 没有区别;消息侦听器 Vs。批处理侦听器只是一种偏好。即使使用消息侦听器,使用手动确认模式,当轮询的所有结果都已处理后,偏移量也会被提交。在 MANUAL_IMMEDIATE 模式下,偏移量被逐个提交。

      【讨论】:

      • 再次感谢加里!关于问题#2 - 我的问题有点不同。如果我有 2 个相同应用程序的实例(具有相同的组 ID)怎么办?每个实例是否会启动 2 个线程并并行消耗(来自每个分区)?关于问题#3 - 是的,我的意思是重新平衡。但是如果我使用批处理模式并且仅在处理完所有记录后才确认,那么框架如何知道在重新平衡期间要提交哪个偏移量?
      • 如果您有 2 个实例,每个实例的并发性为 2,则每个线程将获得一个分区。这取决于重新平衡的原因。如果是因为另一个组成员加入,则在您的侦听器退出(并提交偏移量)之前不会发生。如果发生这种情况是因为您的侦听器超出了max.poll.interval.ms,则不会提交任何偏移量,并且该批次将重播给重新平衡后获得分区的任何消费者。
      • 目前不支持使用批处理侦听器提交单个偏移量。传递给批处理侦听器的 Acknowledgment 参数将提交整个批处理的偏移量。
      • 不过,您可以添加Consumer 作为参数并自己提交偏移量。
      • 您可以在Acknowledgment中拨打nack()。容器将提交到索引的偏移量,执行查找并重新传递失败的记录。或者您可以使用RecoveringBatchErrorHandler。见docs.spring.io/spring-kafka/docs/current/reference/html/…
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-01-06
      • 2017-11-24
      • 2021-04-06
      • 2018-09-05
      • 1970-01-01
      • 2016-06-21
      相关资源
      最近更新 更多