【问题标题】:Spring Kafka with Dynamic @KafkaListener带有动态 @KafkaListener 的 Spring Kafka
【发布时间】:2019-08-24 05:56:59
【问题描述】:

我正在使用带有 spring-kafka 的 Spring Boot 2.x(不是 spring-integration-kafka)

我有多个用 @KafkaListener 注释的 bean ...每个都消耗一个主题...所以既然我有 12 个主题,那么我还需要有 12 个 KafkaConsumers bean ...我想知道我是否可以以编程方式/动态创建这些 bean ...也许使用 KafkaListenerEndpointRegistry 来动态创建消费者容器。

注意:我需要批量消费消息...所以也许我可以使用 BatchMessageListener?

当前代码:

@KafkaListener(
        id = COUNTRY,
        containerFactory = KAFKA_LISTENER_FACTORY_BEAN_NAME,
        topics = {TOPIC},
        groupId = GROUP_ID,
        clientIdPrefix = CLIENT_ID,
        errorHandler = VALIDATION_ERROR_HANDLER_BEAN_NAME
    )
    @Override
    public void consume(final List<MessageDTO> messages,
        @Header(KafkaHeaders.RECEIVED_TOPIC) final List<String> topics,
        @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) final List<String> messagesKey,
        @Header(KafkaHeaders.RECEIVED_PARTITION_ID) final List<Integer> partitionIds,
        @Header(KafkaHeaders.RECEIVED_TIMESTAMP) final List<Long> timestamps,
        @Header(KafkaHeaders.OFFSET) final List<Long> offsets) {
            (...)
    }

每个主题消费者都有自己的实现,具体取决于主题。请你们指点我的博客/伪代码/git线程/答案吗?

【问题讨论】:

  • 每个主题都在同一个集群中?那么有效载荷呢?它们有什么不同吗?
  • 有效载荷结构是相同的......每个国家都有一个主题。然后每个主题至少有一个消费者,因为实施取决于国家/地区

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


【解决方案1】:
【解决方案2】:

如果你的主题有一些模式,你也可以试试这个:

      kafka:
        bindings:
            input.consumer.destination-is-pattern: true

【讨论】:

  • 这仅适用于 Spring Cloud Stream;使用@KafkaListener 时,您将使用topicPattern 注释属性。
猜你喜欢
  • 2018-06-27
  • 2019-09-16
  • 1970-01-01
  • 2021-01-15
  • 2020-10-03
  • 2019-01-20
  • 2020-09-10
  • 2017-11-30
  • 1970-01-01
相关资源
最近更新 更多