如何正确配置Kafka ListDeserializer实现列表消息消费?
解决Kafka ListDeserializer初始化异常问题
我需要配置Kafka监听容器工厂来消费值为对象列表的消息,尝试用ListDeserializer配置工厂,代码如下:
@Configuration public class KafkaConsumerConfig { @Value(value = "${spring.kafka.consumer.bootstrap-servers}") private String bootstrapServers; @Value(value = "${spring.kafka.consumer.group-id}") private String groupId; @Bean public ConcurrentKafkaListenerContainerFactory<String, List<Item>> kafkaListenerContainerFactoryItem() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); ConcurrentKafkaListenerContainerFactory<String, List<Item>> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory( new DefaultKafkaConsumerFactory<>( props, new StringDeserializer(), new ListDeserializer<>(ArrayList.class, new JsonDeserializer<>(Item.class, false)))); return factory; } }
运行时出现如下异常:
org.apache.kafka.common.config.ConfigException: List deserializer was already initialized using a non-default constructor at org.apache.kafka.common.serialization.ListDeserializer.configure(ListDeserializer.java:78) ~[kafka-clients-3.1.1.jar:na] at org.springframework.kafka.core.DefaultKafkaConsumerFactory.lambda$valueDeserializerSupplier$9(DefaultKafkaConsumerFactory.java:199) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.core.DefaultKafkaConsumerFactory.createRawConsumer(DefaultKafkaConsumerFactory.java:483) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.core.DefaultKafkaConsumerFactory.createKafkaConsumer(DefaultKafkaConsumerFactory.java:451) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.core.DefaultKafkaConsumerFactory.createConsumerWithAdjustedProperties(DefaultKafkaConsumerFactory.java:427) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.core.DefaultKafkaConsumerFactory.createKafkaConsumer(DefaultKafkaConsumerFactory.java:394) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.core.DefaultKafkaConsumerFactory.createConsumer(DefaultKafkaConsumerFactory.java:371) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.<init>(KafkaMessageListenerContainer.java:776) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.listener.KafkaMessageListenerContainer.doStart(KafkaMessageListenerContainer.java:352) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.listener.AbstractMessageListenerContainer.start(AbstractMessageListenerContainer.java:461) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.listener.ConcurrentMessageListenerContainer.doStart(ConcurrentMessageListenerContainer.java:209) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.listener.AbstractMessageListenerContainer.start(AbstractMessageListenerContainer.java:461) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.config.KafkaListenerEndpointRegistry.startIfNecessary(KafkaListenerEndpointRegistry.java:347) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.kafka.config.KafkaListenerEndpointRegistry.start(KafkaListenerEndpointRegistry.java:292) ~[spring-kafka-2.8.8.jar:2.8.8] at org.springframework.context.support.DefaultLifecycleProcessor.doStart(DefaultLifecycleProcessor.java:178) ~[spring-context-5.3.22.jar:5.3.22] ... 14 common frames omitted
解决方案
这个异常的根源是:ListDeserializer通过带参构造初始化后,Spring Kafka的DefaultKafkaConsumerFactory会再次调用它的configure方法,而该类的设计不允许这种重复初始化操作。
正确的配置方式是通过配置属性而非带参构造来初始化ListDeserializer,同时指定内部元素的反序列化器。修改后的代码如下:
@Configuration public class KafkaConsumerConfig { @Value(value = "${spring.kafka.consumer.bootstrap-servers}") private String bootstrapServers; @Value(value = "${spring.kafka.consumer.group-id}") private String groupId; @Bean public ConcurrentKafkaListenerContainerFactory<String, List<Item>> kafkaListenerContainerFactoryItem() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // 指定值反序列化器为ListDeserializer props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ListDeserializer.class); // 指定列表元素的反序列化器 props.put(ListDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName()); // 配置JsonDeserializer的信任包(生产环境建议指定具体包名) props.put(JsonDeserializer.TRUSTED_PACKAGES, "*"); // 关闭类型信息头匹配原代码配置 props.put(JsonDeserializer.USE_TYPE_INFO_HEADERS, false); ConcurrentKafkaListenerContainerFactory<String, List<Item>> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>( props, new StringDeserializer(), // 使用无参构造创建ListDeserializer new ListDeserializer<>() )); return factory; } }
关键修改说明
- 放弃
ListDeserializer的带参构造,改用无参构造实例化 - 通过配置属性
ListDeserializer.VALUE_DESERIALIZER_CLASS指定列表元素的反序列化器为JsonDeserializer - 为
JsonDeserializer补充必要配置:TRUSTED_PACKAGES设为*简化配置,USE_TYPE_INFO_HEADERS设为false匹配原代码逻辑 - 让
DefaultKafkaConsumerFactory通过配置属性完成反序列化器初始化,避免重复调用configure方法触发异常
内容的提问来源于stack exchange,提问作者Toni
相关产品推荐
相关产品推荐

