如何编写Kafka XML值反序列化器并注入配置(非仅指定类名)
自定义Kafka XML值反序列化器及配置注入方案
一、编写XML值反序列化器
自定义Kafka值反序列化器需实现org.apache.kafka.common.serialization.Deserializer<T>接口,以下以Jackson XML处理器为例实现:
1. 引入依赖(Maven)
<dependency> <groupId>com.fasterxml.jackson.dataformat</groupId> <artifactId>jackson-dataformat-xml</artifactId> <version>2.15.2</version> <!-- 使用与项目兼容的版本 --> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.4.0</version> <!-- 匹配你的Kafka集群版本 --> </dependency>
2. 实现反序列化器类
假设目标实体类为User(需与XML结构匹配,可添加Jackson XML注解):
import com.fasterxml.jackson.dataformat.xml.XmlMapper; import org.apache.kafka.common.serialization.Deserializer; import java.util.Map; public class XmlValueDeserializer<T> implements Deserializer<T> { private static final XmlMapper XML_MAPPER = new XmlMapper(); private Class<T> targetClass; // 通过configure方法接收目标实体类配置 @Override public void configure(Map<String, ?> configs, boolean isKey) { String className = (String) configs.get("xml.target.class"); try { this.targetClass = (Class<T>) Class.forName(className); } catch (ClassNotFoundException e) { throw new IllegalArgumentException("无法加载目标实体类", e); } } @Override public T deserialize(String topic, byte[] data) { if (data == null) { return null; } try { return XML_MAPPER.readValue(data, targetClass); } catch (Exception e) { throw new org.apache.kafka.common.errors.SerializationException("XML反序列化失败", e); } } @Override public void close() { // 无资源需释放,空实现即可 } }
对应的User实体类示例:
import com.fasterxml.jackson.dataformat.xml.annotation.JacksonXmlRootElement; @JacksonXmlRootElement(localName = "user") public class User { private String id; private String name; // getter、setter、构造函数省略 }
二、配置注入自定义反序列化器
1. 原生Kafka消费者配置
直接在消费者Properties中指定值反序列化器,并传入目标实体类参数:
import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; import java.util.Collections; import java.util.Properties; public class XmlConsumer { public static void main(String[] args) { Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "xml-consumer-group"); // 指定自定义值反序列化器全类名 props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, XmlValueDeserializer.class.getName()); // 传入目标实体类全类名 props.put("xml.target.class", User.class.getName()); KafkaConsumer<String, User> consumer = new KafkaConsumer<>(props); consumer.subscribe(Collections.singletonList("xml-topic")); // 消费逻辑省略 } }
2. Spring Kafka配置(非delegate方式)
方式一:配置类注册实例
import org.apache.kafka.clients.consumer.ConsumerConfig; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import java.util.HashMap; import java.util.Map; @Configuration public class KafkaConfig { @Bean public ConsumerFactory<String, User> consumerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); configProps.put(ConsumerConfig.GROUP_ID_CONFIG, "xml-spring-group"); // 直接传入自定义反序列化器实例 return new DefaultKafkaConsumerFactory<>(configProps, new org.apache.kafka.common.serialization.StringDeserializer(), new XmlValueDeserializer<User>() {{ try { this.targetClass = (Class<User>) Class.forName(User.class.getName()); } catch (ClassNotFoundException e) { throw new RuntimeException(e); } }}); } // 可后续注册KafkaListenerContainerFactory等组件 }
方式二:配置文件指定类名
在application.yml中直接配置:
spring: kafka: consumer: bootstrap-servers: localhost:9092 group-id: xml-spring-group value-deserializer: com.example.XmlValueDeserializer properties: xml.target.class: com.example.User
内容的提问来源于stack exchange,提问作者αƞjiβ
相关产品推荐
相关产品推荐

