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

Spring JMS中如何处理两种不同消息类型?

问题描述

在Spring应用中使用嵌入式ActiveMQ时,如何处理DealCreateDto和DealUpdateDto两种不同的消息数据类型?是否需要为每种消息类型创建一个JmsTemplate实例,或是有其他方案?

我已经实现了两个MessageConverter:

@Component
@RequiredArgsConstructor
public class DealCreateDtoConverter implements MessageConverter {

    private final ObjectMapper objectMapper;

    @NonNull
    @Override
    @SneakyThrows
    public Message toMessage(@NonNull Object object, @NonNull Session session) {
        TextMessage message = session.createTextMessage();
        message.setText(objectMapper.writeValueAsString(object));
        return message;
    }

    @NonNull
    @Override
    @SneakyThrows
    public DealCreateDto fromMessage(@NonNull Message message) {
        return objectMapper.readValue(((TextMessage) message).getText(), DealCreateDto.class);
    }
}

@Component
@RequiredArgsConstructor
public class DealUpdateDtoConverter implements MessageConverter {

    private final ObjectMapper objectMapper;

    @NonNull
    @Override
    @SneakyThrows
    public Message toMessage(@NonNull Object object, @NonNull Session session) {
        TextMessage message = session.createTextMessage();
        message.setText(objectMapper.writeValueAsString(object));
        return message;
    }

    @NonNull
    @Override
    @SneakyThrows
    public DealUpdateDto fromMessage(@NonNull Message message) {
        return objectMapper.readValue(((TextMessage) message).getText(), DealUpdateDto.class);
    }
}

但我了解到一个JmsTemplate只能配置一个消息转换器,当我尝试同时使用这两种消息类型时:

private final JmsTemplate jmsTemplate;

@Override
public void create(@NonNull DealCreateDto deal) {
    jmsTemplate.convertAndSend(ActiveMqConfig.DEAL_CREATE_QUEUE, deal);
}

@JmsListener(destination = ActiveMqConfig.DEAL_CREATE_QUEUE, containerFactory = "jmsFactory")
public void receiveCreate(@Payload @NonNull DealCreateDto dealDto) {
}

@Override
public void update(@NonNull DealUpdateDto deal) {
    jmsTemplate.convertAndSend(ActiveMqConfig.DEAL_UPDATE_QUEUE, deal);
}

@JmsListener(destination = ActiveMqConfig.DEAL_UPDATE_QUEUE, containerFactory = "jmsFactory")
public void receiveUpdate(@Payload @NonNull DealUpdateDto dealDto) {
}

出现了如下异常:

2022-08-16 08:41:19:131 ERROR o.a.c.c.C.[.[.[.[dispatcherServlet] - Servlet.service() for servlet [dispatcherServlet] in context with path [] threw exception [Request processing failed; nested exception is org.springframework.jms.support.converter.MessageConversionException: Cannot convert object of type [ltd.mydomain.dto.deal.DealCreateDto] to JMS message. Supported message payloads are: String, byte array, Map<String,?>, Serializable object.] with root cause
org.springframework.jms.support.converter.MessageConversionException: Cannot convert object of type [ltd.mydomain.dto.deal.DealCreateDto] to JMS message. Supported message payloads are: String, byte array, Map<String,?>, Serializable object.
    at org.springframework.jms.support.converter.SimpleMessageConverter.toMessage(SimpleMessageConverter.java:79)
    ...(省略中间栈信息)
    at java.base/java.lang.Thread.run(Thread.java:833)

我的JMS配置如下:

@EnableJms
@Configuration
public class ActiveMqConfig {

    public static final String DEAL_CREATE_QUEUE = "deal-create-queue";

    public static final String DEAL_UPDATE_QUEUE = "deal-update-queue";

    @Bean
    public JmsListenerContainerFactory<?> jmsFactory(
            @NonNull ConnectionFactory connectionFactory,
            @NonNull DefaultJmsListenerContainerFactoryConfigurer configurer) {
        DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
        configurer.configure(factory, connectionFactory);
        return factory;
    }
}
解决方案

方案一:使用通用Jackson消息转换器(推荐)

无需为每个DTO单独实现转换器,直接用Spring提供的MappingJackson2MessageConverter,它能自动处理多类型DTO的序列化/反序列化,通过消息头传递类型信息,接收端自动匹配目标类型。

配置通用转换器

@Configuration
public class JmsConfig {
    // 配置基于Jackson的通用消息转换器
    @Bean
    public MessageConverter jacksonJmsMessageConverter(ObjectMapper objectMapper) {
        MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter();
        converter.setTargetType(MessageType.TEXT);
        // 设置类型标识的消息头,接收端通过该头识别DTO类型
        converter.setTypeIdPropertyName("_type");
        converter.setObjectMapper(objectMapper);
        return converter;
    }

    // 配置绑定通用转换器的JmsTemplate
    @Bean
    public JmsTemplate jmsTemplate(ConnectionFactory connectionFactory, MessageConverter jacksonJmsMessageConverter) {
        JmsTemplate template = new JmsTemplate(connectionFactory);
        template.setMessageConverter(jacksonJmsMessageConverter);
        return template;
    }

    // 配置绑定通用转换器的监听容器工厂
    @Bean
    public JmsListenerContainerFactory<?> jmsFactory(
            ConnectionFactory connectionFactory,
            DefaultJmsListenerContainerFactoryConfigurer configurer,
            MessageConverter jacksonJmsMessageConverter) {
        DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
        configurer.configure(factory, connectionFactory);
        factory.setMessageConverter(jacksonJmsMessageConverter);
        return factory;
    }
}

使用方式

发送端直接调用convertAndSend即可,无需额外配置:

@Override
public void create(@NonNull DealCreateDto deal) {
    jmsTemplate.convertAndSend(ActiveMqConfig.DEAL_CREATE_QUEUE, deal);
}

@Override
public void update(@NonNull DealUpdateDto deal) {
    jmsTemplate.convertAndSend(ActiveMqConfig.DEAL_UPDATE_QUEUE, deal);
}

接收端保持原有代码,转换器会根据消息头自动反序列化到对应DTO:

@JmsListener(destination = ActiveMqConfig.DEAL_CREATE_QUEUE, containerFactory = "jmsFactory")
public void receiveCreate(@Payload @NonNull DealCreateDto dealDto) {
    // 处理创建逻辑
}

@JmsListener(destination = ActiveMqConfig.DEAL_UPDATE_QUEUE, containerFactory = "jmsFactory")
public void receiveUpdate(@Payload @NonNull DealUpdateDto dealDto) {
    // 处理更新逻辑
}

方案二:配置多个JmsTemplate(适配自定义转换器)

如果坚持使用自己实现的转换器,可以为每个DTO类型单独配置JmsTemplate,绑定对应的转换器:

配置多个JmsTemplate和监听容器工厂

@Configuration
public class JmsConfig {
    // 绑定DealCreateDto转换器的JmsTemplate
    @Bean
    public JmsTemplate dealCreateJmsTemplate(ConnectionFactory connectionFactory, DealCreateDtoConverter createConverter) {
        JmsTemplate template = new JmsTemplate(connectionFactory);
        template.setMessageConverter(createConverter);
        return template;
    }

    // 绑定DealUpdateDto转换器的JmsTemplate
    @Bean
    public JmsTemplate dealUpdateJmsTemplate(ConnectionFactory connectionFactory, DealUpdateDtoConverter updateConverter) {
        JmsTemplate template = new JmsTemplate(connectionFactory);
        template.setMessageConverter(updateConverter);
        return template;
    }

    // 对应创建队列的监听容器工厂
    @Bean
    public JmsListenerContainerFactory<?> createQueueFactory(
            ConnectionFactory connectionFactory,
            DefaultJmsListenerContainerFactoryConfigurer configurer,
            DealCreateDtoConverter createConverter) {
        DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
        configurer.configure(factory, connectionFactory);
        factory.setMessageConverter(createConverter);
        return factory;
    }

    // 对应更新队列的监听容器工厂
    @Bean
    public JmsListenerContainerFactory<?> updateQueueFactory(
            ConnectionFactory connectionFactory,
            DefaultJmsListenerContainerFactoryConfigurer configurer,
            DealUpdateDtoConverter updateConverter) {
        DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
        configurer.configure(factory, connectionFactory);
        factory.setMessageConverter(updateConverter);
        return factory;
    }
}

使用方式

发送端注入对应的JmsTemplate:

private final JmsTemplate dealCreateJmsTemplate;
private final JmsTemplate dealUpdateJmsTemplate;

@Override
public void create(@NonNull DealCreateDto deal) {
    dealCreateJmsTemplate.convertAndSend(ActiveMqConfig.DEAL_CREATE_QUEUE, deal);
}

@Override
public void update(@NonNull DealUpdateDto deal) {
    dealUpdateJmsTemplate.convertAndSend(ActiveMqConfig.DEAL_UPDATE_QUEUE, deal);
}

接收端指定对应容器工厂:

@JmsListener(destination = ActiveMqConfig.DEAL_CREATE_QUEUE, containerFactory = "createQueueFactory")
public void receiveCreate(@Payload @NonNull DealCreateDto dealDto) {
    // 处理创建逻辑
}

@JmsListener(destination = ActiveMqConfig.DEAL_UPDATE_QUEUE, containerFactory = "updateQueueFactory")
public void receiveUpdate(@Payload @NonNull DealUpdateDto dealDto) {
    // 处理更新逻辑
}

方案三:发送时手动指定转换器

不想配置多个JmsTemplate的话,可在发送时直接传入对应转换器:

发送端代码

private final JmsTemplate jmsTemplate;
private final DealCreateDtoConverter createConverter;
private final DealUpdateDtoConverter updateConverter;

@Override
public void create(@NonNull DealCreateDto deal) {
    jmsTemplate.convertAndSend(ActiveMqConfig.DEAL_CREATE_QUEUE, deal, createConverter);
}

@Override
public void update(@NonNull DealUpdateDto deal) {
    jmsTemplate.convertAndSend(ActiveMqConfig.DEAL_UPDATE_QUEUE, deal, updateConverter);
}

接收端配置

接收端需为每个队列配置对应转换器的监听容器工厂(同方案二的监听容器配置)。

原异常原因分析

原代码中未给JmsTemplate配置自定义转换器,Spring默认使用SimpleMessageConverter,它仅支持String、字节数组、Map、Serializable类型的对象转换。而你的DealCreateDto和DealUpdateDto未实现Serializable接口,因此触发转换异常。

内容的提问来源于stack exchange,提问作者Pavel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 16:03:37