【问题标题】:When to use ConcurrentKafkaListenerContainerFactory?何时使用 ConcurrentKafkaListenerContainerFactory?
【发布时间】:2019-07-28 03:02:12
【问题描述】:

我是 kafka 的新手,我通过了documentation,但我什么都不懂。有人可以解释何时使用ConcurrentKafkaListenerContainerFactory 类吗?我使用了Kafkaconsumer 类,但我看到ConcurrentKafkaListenerContainerFactory 正在我当前的项目中使用。请说明它的用途。

【问题讨论】:

  • ConcurrentKafkaListenerContainerFactory 来自 spring 框架,可以在 spring 生态系统中使用。 KafkaConsumer 来自 Apache 的 Kafka 的 Java sdk。两者只是在 Java 上实现 Kafka 消费者的不同工具/api。只是提供的功能可能有所不同。
  • 感谢@MadhuBhat 提供的信息,但这就是我想知道的,我什么时候想在 kafkaconsumer 上使用 ConcurrentKafkaListenerContainerFactory 以及它提供的额外功能。

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


【解决方案1】:

Kafka 消费者不是线程安全的。所有网络 I/O 都发生在进行调用的应用程序的线程中。确保多线程访问正确同步是用户的责任。非同步访问会导致ConcurrentModificationException.

如果为消费者分配了多个分区来获取数据,它将尝试同时从所有分区中消费,从而有效地为这些分区提供相同的消费优先级。然而,在某些情况下,消费者可能希望首先专注于全速从已分配分区的某个子集中获取数据,并且仅在这些分区几乎没有或没有数据可使用时才开始获取其他分区。

Spring-kafka

ConcurrentKafkaListenerContainerFactory 用于为带有@KafkaListener 注释的方法创建容器

spring kafka中有两个MessageListenerContainer

KafkaMessageListenerContainer
ConcurrentMessageListenerContainer

KafkaMessageListenerContainer 在单个线程上接收来自所有主题或分区的所有消息。 ConcurrentMessageListenerContainer 委托给一个或多个 KafkaMessageListenerContainer 实例以提供多线程消费。

使用 ConcurrentMessageListenerContainer

@Bean
KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Integer, String>>
                    kafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<Integer, String> factory =
                            new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setConcurrency(3);
    factory.getContainerProperties().setPollTimeout(3000);
    return factory;
  }

它具有并发属性。例如,container.setConcurrency(3) 创建了三个KafkaMessageListenerContainer 实例。

如果你提供了六个TopicPartition实例并且并发为3;每个容器有两个分区。对于五个 TopicPartition 实例,两个容器获得两个分区,第三个获得一个。如果并发大于 TopicPartition 的个数,则调低并发,使每个容器得到一个分区。

这是带有文档here的明确示例

【讨论】:

【解决方案2】:

Kafka Consumer API 不是线程安全的。 ConcurrentKafkaListenerContainerFactory api 提供了使用 Kafka Consumer API 的并发方式以及设置其他 kafka 消费者属性。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-10-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-11-08
    • 1970-01-01
    相关资源
    最近更新 更多