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

Spring+Axon项目无法消费RabbitMQ队列消息问题求助

RabbitMQ消息无法被Axon+Spring应用消费排查与解决

我正在开发一个基于Spring和Axon框架的CQRS应用,负责消费RabbitMQ队列中的事件消息,但消息进入队列后应用完全没有消费动作,相关日志也没有任何输出。以下是项目的核心代码和配置:

UserEventHandler.java

package com.jasper.ecommerce.query.handler;

import com.jasper.ecommerce.query.model.UserReadModel;
import com.jasper.ecommerce.query.query.FindAllUsersQuery;
import com.jasper.ecommerce.query.repository.UserRepository;
import com.jasper.ecommerce.shared.events.UserCreatedEvent;
import org.axonframework.eventhandling.EventHandler;
import org.axonframework.queryhandling.QueryHandler;
import org.springframework.stereotype.Component;

import java.util.List;

@Component
public class UserEventHandler {
    private final UserRepository userRepository;

    public UserEventHandler(UserRepository userRepository) {
        this.userRepository = userRepository;
    }

    // Event handler for UserCreatedEvent
    @EventHandler
    public void on(UserCreatedEvent event) {
        System.out.println("UserCreatedEvent received: " + event);
        // Convert event to read model and save it in MongoDB
        UserReadModel user = new UserReadModel(event.getUserId(), event.getName(), event.getEmail());
        userRepository.save(user);  // This inserts the document into MongoDB
    }

    // Query handler to return all users
    @QueryHandler
    public List<UserReadModel> handle(FindAllUsersQuery query) {
        return userRepository.findAll();
    }
}

AxonAMQPConfiguration.java

package com.jasper.ecommerce.query.config;

import org.axonframework.config.EventProcessingConfigurer;
import org.axonframework.extensions.amqp.eventhandling.DefaultAMQPMessageConverter;
import org.axonframework.extensions.amqp.eventhandling.spring.SpringAMQPMessageSource;
import org.axonframework.serialization.Serializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class AxonAMQPConfiguration {

    @Bean
    public SpringAMQPMessageSource rabbitMQMessageSource(Serializer serializer) {
        return new SpringAMQPMessageSource(
                DefaultAMQPMessageConverter.builder()
                        .serializer(serializer)
                        .build()
        );
    }

    public void configure(EventProcessingConfigurer configurer, SpringAMQPMessageSource rabbitMQMessageSource) {
        configurer.registerSubscribingEventProcessor("UserEventHandler", c -> rabbitMQMessageSource);
    }
}

application.yml

spring:
  application:
    name: ecommerce-cqrs-query
  data:
    mongodb:
      uri: mongodb://localhost:27017/ecommerce-query
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest

axon:
  amqp:
    exchange: AxonExchange
  eventhandling:
    processors:
      UserEventHandler:
        mode: subscribing
  axonserver:
    enabled: false
  distributed:
    enabled: false

logging:
  level:
    org.axonframework: DEBUG
    org.springframework.amqp: DEBUG
    com.rabbitmq: DEBUG

server:
  port: 8080

问题排查与修复步骤

  1. 修复配置类回调触发问题
    你的AxonAMQPConfiguration中的configure方法未被Spring识别为配置回调,导致事件处理器没有绑定到RabbitMQ消息源。可以通过两种方式修复:

    • 实现AxonConfigurationConfigurer接口:
      @Configuration
      public class AxonAMQPConfiguration implements AxonConfigurationConfigurer {
      
          @Bean
          public SpringAMQPMessageSource rabbitMQMessageSource(Serializer serializer) {
              return new SpringAMQPMessageSource(
                      DefaultAMQPMessageConverter.builder()
                              .serializer(serializer)
                              .build()
              );
          }
      
          @Override
          public void configure(AxonConfiguration axonConfiguration) {
              EventProcessingConfigurer configurer = axonConfiguration.eventProcessingConfigurer();
              configurer.registerSubscribingEventProcessor("UserEventHandler", 
                  c -> rabbitMQMessageSource(axonConfiguration.serializer()));
          }
      }
      
    • 给configure方法添加@Autowired注解:
      @Autowired
      public void configure(EventProcessingConfigurer configurer, SpringAMQPMessageSource rabbitMQMessageSource) {
          configurer.registerSubscribingEventProcessor("UserEventHandler", c -> rabbitMQMessageSource);
      }
      
  2. 绑定事件处理器到指定分组
    当前事件处理器未指定分组,Axon会默认将其归到default组,但你配置的处理器组是UserEventHandler。给UserEventHandler类添加@ProcessingGroup注解:

    @Component
    @ProcessingGroup("UserEventHandler")
    public class UserEventHandler {
        // 现有代码不变
    }
    
  3. 确认RabbitMQ队列与绑定关系
    登录RabbitMQ管理后台检查:

    • AxonExchange交换机是否存在
    • UserEventHandler队列是否存在,且与AxonExchange正确绑定
    • 消息的路由键是否与绑定规则匹配(Axon默认用事件全类名作为路由键)
  4. 统一事件序列化方式
    确保消息发送端和消费端使用相同的序列化器,避免因序列化不一致导致消息无法解析。可以在application.yml中明确指定:

    axon:
      serializer:
        events: jackson
        messages: jackson
    
  5. 验证日志输出配置
    确认日志框架(如Logback)正确加载了配置,org.axonframework和org.springframework.amqp的DEBUG日志确实能输出到控制台或日志文件中,方便排查后续问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 06:45:21