【问题标题】:ConcurrentKafkaListenerContainerFactory message converter is ignored when configuring listeners automatically自动配置监听器时 ConcurrentKafkaListenerContainerFactory 消息转换器被忽略
【发布时间】:2022-01-13 00:33:58
【问题描述】:

我需要在运行时创建 Kafka 侦听器,一切似乎都正常,除了消息转换器属性似乎被忽略(或者这可能是一个设计功能或者我做错了什么)。

使用@KafkaListener 时,它可以正常工作,但手动创建侦听器时,我的消息没有转换为所需的对象,并且出现错误:

Caused by: java.lang.ClassCastException: class java.lang.String cannot be cast to class com.my.company.model.MyPojo (java.lang.String is in module java.base of loader 'bootstrap'; com.my.company.model.MyPojo is in unnamed module of loader 'app')
    at com.my.company.config.MyPojo.kafka.KafkaConfig.lambda$createListenerContainers$2(KafkaConfig.java:142)

我的配置:

@Bean
ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory() {    
    var factory = new ConcurrentKafkaListenerContainerFactory<String, Object>();
    factory.setConsumerFactory(consumerFactory());
    factory.setMessageConverter(new StringJsonMessageConverter());       
    return factory;
}

@Bean
MessageListenerContainer createListenerContainer1() {
    ContainerProperties containerProperties = new ContainerProperties(topicConfig("my_topic"));
    var container = new KafkaMessageListenerContainer<>(consumerFactory(), containerProperties);
    //tried this too...
    //var container = kafkaListenerContainerFactory().createContainer(topicConfig("my_topic")); 
    container.setupMessageListener((MessageListener<String, MyPojo>) data -> getDataService.process(data.value()););
    container.start();

    return container;
}

工作中的 Kafka 监听器:

@KafkaListener(id = "1", topics = "my_topic)
public void listenGetDataTopic(@Payload MyPojo message) {
    log.info(message);
}

我尝试了很多不同的配置并对其进行深入调试,当然我看到了使用 @KafkaListener 和手动创建侦听器时处理消息之间的区别,但我不知道如何应用消息转换为手动创建的侦听器。有没有可能做到这一点?

【问题讨论】:

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


    【解决方案1】:

    消息转换器不是容器的属性,它是侦听器适配器的属性,用于调用@KafkaListener 的pojo 方法。

    直接使用容器时,您的侦听器必须实现MessageListener 或其子接口之一。

    您可以自己在侦听器中调用转换器(例如,创建一个轻量级适配器),或者您需要使用另一种技术来动态创建 @KafkaListeners。

    Kafka Spring: How to create Listeners dynamically or in a loop?

    Kafka Consumer in spring can I re-assign partitions programmatically?

    Can i add topics to my @kafkalistener at runtime

    这些技术的一些例子。

    【讨论】:

    • 谢谢加里。将来是否可以扩展配置以允许全局配置消息转换器? (什么,顺便说一句,因为我们将它分配给工厂,但实际上它的工作方式有点不同)
    猜你喜欢
    • 2013-09-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-12-20
    • 1970-01-01
    • 1970-01-01
    • 2019-07-23
    相关资源
    最近更新 更多