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

RabbitMQ单队列多事件类型消费异常解决方案求助

RabbitMQ跨服务多事件消费问题解决方案

问题背景

作为RabbitMQ新手,正在构建跨语言微服务通信Demo:两个Spring Boot项目连接同一RabbitMQ,单队列处理多种事件,要求消息顺序一致,每个生产者对应一个消费者。通过自定义eventType消息头区分事件类型,手动反序列化时出现ClassNotFoundException,异常显示消费者找不到生产者端的事件类全限定名。

异常原因

Spring AMQP默认的Jackson2JsonMessageConverter会自动在消息头中添加__TypeId__字段,存储生产者事件类的完整包路径(如com.jchaaban.rabbitmqproducer.event.MessageDeletedEvent)。消费者端的事件类包名不同(com.jchaaban.rabbitmqcosumer.event),框架尝试自动转换时找不到对应类,导致异常——即使你想手动处理Message对象,框架的自动转换逻辑依然会先执行。

解决方案

方案1:禁用自动类型头,自定义类型映射

生产者配置修改

修改Jackson2JsonMessageConverter,禁用自动添加类型头,只保留自定义的eventType头:

@Bean
public MessageConverter jsonMessageConverter() {
    Jackson2JsonMessageConverter converter = new Jackson2JsonMessageConverter();
    // 禁用默认的类型头
    converter.setTypeIdPropertyName(null);
    return converter;
}

消费者配置修改

配置Jackson2JsonMessageConverter,根据自定义的eventType头映射到消费者本地的事件类:

@Bean
public MessageConverter jsonMessageConverter(ObjectMapper objectMapper) {
    Jackson2JsonMessageConverter converter = new Jackson2JsonMessageConverter(objectMapper);
    DefaultJackson2JavaTypeMapper typeMapper = new DefaultJackson2JavaTypeMapper();
    // 自定义类型映射:eventType值 -> 消费者本地类
    Map<String, Class<?>> typeMappings = new HashMap<>();
    typeMappings.put("MessageCreatedEvent", MessageCreatedEvent.class);
    typeMappings.put("MessageDeletedEvent", MessageDeletedEvent.class);
    typeMapper.setIdClassMapping(typeMappings);
    // 指定用自定义的eventType头作为类型标识
    typeMapper.setTypeIdHeader("eventType");
    converter.setJavaTypeMapper(typeMapper);
    return converter;
}

消费者消费逻辑简化

此时可以直接接收具体事件类型,无需手动反序列化:

@Component
public class RabbitMqReceiver {
    private static final Logger logger = LoggerFactory.getLogger(RabbitMqReceiver.class);

    @RabbitListener(queues = "${spring.rabbitmq.eventsQueueName}")
    public void handleMessageCreated(MessageCreatedEvent event) {
        logger.info("Received Message Creation Event: {}", event);
    }

    @RabbitListener(queues = "${spring.rabbitmq.eventsQueueName}")
    public void handleMessageDeleted(MessageDeletedEvent event) {
        logger.info("Received Message Deletion Event: {}", event);
    }
}

方案2:直接接收原始消息,关闭自动转换

如果坚持手动处理反序列化,可以让@RabbitListener接收byte[]或Message,并关闭框架的自动转换:

消费者监听方法修改

@RabbitListener(queues = "${spring.rabbitmq.eventsQueueName}", messageConverter = "rawMessageConverter")
public void receiveMessage(Message message) {
    try {
        String eventType = message.getMessageProperties().getHeader("eventType");
        Object event = deserializeEvent(message.getBody(), eventType);
        handleEvent(event);
    } catch (Exception e) {
        logger.error("Failed to process message", e);
    }
}

// 新增一个不做转换的MessageConverter Bean
@Bean
public MessageConverter rawMessageConverter() {
    return new MessageConverter() {
        @Override
        public Message toMessage(Object object, MessageProperties messageProperties) throws MessageConversionException {
            return null;
        }

        @Override
        public Object fromMessage(Message message) throws MessageConversionException {
            // 直接返回原始Message,不做转换
            return message;
        }
    };
}

方案3:使用事件基类+条件监听

定义一个公共事件基类,所有事件继承它,然后通过condition参数根据eventType头匹配不同处理方法:

定义基类

public abstract class BaseEvent {
    // 公共字段(如事件ID、时间戳等)
}

// 事件类继承基类
public class MessageCreatedEvent extends BaseEvent {
    // 自定义字段
}

public class MessageDeletedEvent extends BaseEvent {
    // 自定义字段
}

消费者监听方法

@Component
public class RabbitMqReceiver {
    private static final Logger logger = LoggerFactory.getLogger(RabbitMqReceiver.class);

    @RabbitListener(queues = "${spring.rabbitmq.eventsQueueName}", condition = "headers['eventType'] == 'MessageCreatedEvent'")
    public void handleCreationEvent(MessageCreatedEvent event) {
        logger.info("Received Message Creation Event: {}", event);
    }

    @RabbitListener(queues = "${spring.rabbitmq.eventsQueueName}", condition = "headers['eventType'] == 'MessageDeletedEvent'")
    public void handleDeletionEvent(MessageDeletedEvent event) {
        logger.info("Received Message Deletion Event: {}", event);
    }
}

消费者配置修改

同样需要配置Jackson2JsonMessageConverter,禁用默认类型头并设置自定义类型映射(同方案1的消费者配置)。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 03:20:55