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

如何在Spring Cloud Stream中对单个Kafka主题的不同数据类型消息进行反序列化?

如何在Spring Cloud Stream中对单个Kafka主题的不同数据类型消息进行反序列化?

我之前在做Spring Boot + Spring Cloud Stream对接Kafka的项目时,刚好遇到过和你一模一样的需求——单个主题里混着不同类型的JSON消息,靠header里的type字段区分,要自动反序列化成对应的Java实体类。折腾了一番后,总结出了几种可行的方案,下面分享给你:

一、先准备基础实体类

首先得把你要处理的消息实体类定义好,比如Customer和Address:

public class Customer {
    private String id;
    private String name;

    // 记得加上无参构造、全参构造,以及getter/setter方法
    // 这里为了简洁省略,实际项目里可不能忘哦
}

public class Address {
    private String street;
    private String city;

    // 同样,构造方法、getter/setter不能少
}

二、方案一:自定义消息转换器(优雅统一处理)

这种方式是全局配置一个转换器,所有进入的消息都会先经过它,根据type header自动选择对应的实体类反序列化,后续消费端直接拿实体类用就行,非常省心。

1. 实现自定义消息转换器

继承Spring的AbstractMessageConverter,核心逻辑就是读取header里的type,然后用Jackson反序列化:

import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.messaging.Message;
import org.springframework.messaging.converter.AbstractMessageConverter;
import org.springframework.messaging.converter.MessageConversionException;
import org.springframework.util.MimeType;

public class DynamicJsonMessageConverter extends AbstractMessageConverter {

    private final ObjectMapper objectMapper;
    private static final String TYPE_HEADER = "type";

    public DynamicJsonMessageConverter(ObjectMapper objectMapper) {
        super(MimeType.valueOf("application/json"));
        this.objectMapper = objectMapper;
    }

    @Override
    protected boolean supports(Class<?> clazz) {
        // 支持我们定义的所有实体类,这里直接把Customer和Address加进去
        return Customer.class.isAssignableFrom(clazz) || Address.class.isAssignableFrom(clazz);
    }

    @Override
    protected Object convertFromInternal(Message<?> message, Class<?> targetClass, Object conversionHint) {
        // 从header里拿type字段
        String typeHeader = message.getHeaders().get(TYPE_HEADER, String.class);
        if (typeHeader == null) {
            throw new MessageConversionException("消息里缺少必填的'type'头信息");
        }

        // 根据type匹配对应的实体类
        Class<?> actualClass;
        switch (typeHeader) {
            case "customer":
                actualClass = Customer.class;
                break;
            case "address":
                actualClass = Address.class;
                break;
            default:
                throw new MessageConversionException("不支持的消息类型: " + typeHeader);
        }

        // 反序列化payload
        try {
            byte[] payload = (byte[]) message.getPayload();
            return objectMapper.readValue(payload, actualClass);
        } catch (Exception e) {
            throw new MessageConversionException("反序列化消息失败", e);
        }
    }
}

2. 注册自定义转换器到Spring容器

写个配置类,把上面的转换器注册进去,让Spring Cloud Stream使用它:

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.messaging.converter.MessageConverter;
import org.springframework.cloud.stream.config.StreamMessageConverterCustomizer;
import com.fasterxml.jackson.databind.ObjectMapper;

@Configuration
public class StreamConfig {

    @Bean
    public DynamicJsonMessageConverter dynamicJsonMessageConverter(ObjectMapper objectMapper) {
        // 用Spring容器里的ObjectMapper,避免重复创建
        return new DynamicJsonMessageConverter(objectMapper);
    }

    @Bean
    public StreamMessageConverterCustomizer messageConverterCustomizer(DynamicJsonMessageConverter converter) {
        return (messageConverters) -> {
            // 把自定义转换器放到最前面,确保优先使用
            messageConverters.add(0, converter);
        };
    }
}

三、方案二:结合@StreamListener的条件路由

配置好自定义转换器后,就可以用@StreamListener配合条件表达式,把不同类型的消息路由到对应的处理方法,方法参数直接写实体类就行:

import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.cloud.stream.messaging.Sink;
import org.springframework.stereotype.Component;

@Component
public class MessageConsumer {

    // 只处理type为customer的消息
    @StreamListener(target = Sink.INPUT, condition = "headers['type']=='customer'")
    public void handleCustomer(Customer customer) {
        // 这里直接拿到反序列化后的Customer对象,做业务处理
        System.out.println("收到Customer消息: " + customer);
    }

    // 只处理type为address的消息
    @StreamListener(target = Sink.INPUT, condition = "headers['type']=='address'")
    public void handleAddress(Address address) {
        System.out.println("收到Address消息: " + address);
    }
}

四、方案三:手动反序列化(简单直接)

如果你的项目比较简单,不想自定义转换器,也可以在消费端手动处理反序列化,直接拿byte[]类型的payload,结合type header自己转:

import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.cloud.stream.messaging.Sink;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;

@Component
public class SimpleMessageConsumer {

    private final ObjectMapper objectMapper;

    // 注入Spring容器的ObjectMapper
    public SimpleMessageConsumer(ObjectMapper objectMapper) {
        this.objectMapper = objectMapper;
    }

    @StreamListener(Sink.INPUT)
    public void handleMessage(@Header("type") String type, byte[] payload) {
        try {
            if ("customer".equals(type)) {
                Customer customer = objectMapper.readValue(payload, Customer.class);
                handleCustomer(customer);
            } else if ("address".equals(type)) {
                Address address = objectMapper.readValue(payload, Address.class);
                handleAddress(address);
            } else {
                System.err.println("未知的消息类型: " + type);
            }
        } catch (Exception e) {
            System.err.println("反序列化消息失败: " + e.getMessage());
        }
    }

    private void handleCustomer(Customer customer) {
        // 处理Customer业务
        System.out.println("处理Customer: " + customer);
    }

    private void handleAddress(Address address) {
        // 处理Address业务
        System.out.println("处理Address: " + address);
    }
}

五、一些注意事项

  • 生产者发送消息时,一定要记得设置type header,并且把消息的content-type设为application/json,不然转换器可能不生效。
  • 如果你后续要加新的消息类型,只需要新增对应的实体类,然后在自定义转换器的switch里加个case就行,扩展性很好。
  • 尽量用Spring容器里的ObjectMapper,不要自己new,这样可以复用Spring的全局配置(比如日期格式化、空值处理等)。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.08 12:50:30