【问题标题】:how to set kafka consumer concurrency using spring boot如何使用spring boot设置kafka消费者并发
【发布时间】:2018-10-20 04:25:58
【问题描述】:

我正在编写一个基于 Java 的 Kafka Consumer 应用程序。我正在为我的应用程序使用 kafka-clients、Spring Kafka 和 Spring boot。虽然 Spring boot 让我可以轻松编写 Kafka 消费者(无需真正编写 ConcurrentKafkaListenerContainerFactory、ConsumerFactory 等),但我希望能够为这些消费者定义/自定义一些属性。但是,我找不到使用 Spring boot 的简单方法。例如:我有兴趣设置的一些属性是 -

ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG

我查看了 Spring Boot 预定义属性here

另外,基于之前的问题here,我想在消费者上设置并发,但找不到配置、application.properties 驱动的方式来使用 Spring Boot。

一个明显的方法是在我的 Spring 上下文中再次定义 ConcurrentKafkaListenerContainerFactory, ConsumerFactory 类并从那里开始工作。我想了解是否有更清洁的方法,特别是因为我使用的是 Spring Boot。

版本-

  • kafka-clients - 0.10.0.0-SASL
  • spring-kafka - 1.1.0.RELEASE
  • 弹簧靴 - 1.5.10.RELEASE

【问题讨论】:

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


    【解决方案1】:

    在您引用的 URL 处,向下滚动到

    spring.kafka.listener.concurrency= # Number of threads to run in the listener containers.
    

    spring-kafka - 1.1.0.RELEASE

    我建议至少升级到 1.3.5;由于 KIP-62,它有一个更简单的线程模型。

    编辑

    使用 Boot 2.0,您可以设置任意生产者、消费者、管理员、公共属性,如 in the boot documentation 所述。

    spring.kafka.consumer.properties.heartbeat.interval.ms
    

    在 Boot 1.5 中,只有 spring.kafka.properties,如 here 所述。

    这会为生产者和消费者设置属性,但您可能会在日志中看到一些关于生产者未使用/不受支持的属性的噪音。

    或者,您可以简单地覆盖 Boot 的消费者工厂并根据需要添加属性...

    @Bean
    public ConsumerFactory<?, ?> kafkaConsumerFactory(KafkaProperties properties) {
        Map<String, Object> consumerProps = properties.buildConsumerProperties();
        consumerProps.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 5_000);
        return new DefaultKafkaConsumerFactory<Object, Object>(consumerProps);
    }
    

    【讨论】:

    • 啊,谢谢你的并发属性,Gary!但是,您能否阐明我如何配置 Spring Boot 不支持的属性? 我尝试使用 application.properties 中的一些属性和一些通过代码定义的属性来初始化 ConsumerConfig,但是方法不起作用。至于升级版本,我依赖于 kafka-clients 版本,因为我们的代理支持该版本。因此,我必须根据兼容性矩阵使用上述版本 - projects.spring.io/spring-kafka/#quick-start
    • 查看我的答案的编辑。无论如何,我强烈建议您升级到更新的 kafka-clients(和 spring-kafka 1.3.5)。感谢KIP-35,新客户可以与老经纪人交谈,只要您不尝试使用经纪人不支持的功能。
    【解决方案2】:

    【讨论】:

      猜你喜欢
      • 2020-02-27
      • 2016-12-09
      • 2016-06-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-02-18
      • 2019-08-25
      • 1970-01-01
      相关资源
      最近更新 更多