You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何编写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β

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 18:52:57