Axon+Spring Boot实现CQRS:命令与查询服务事件传递异常排查
问题解决步骤与示例实现
你的核心问题是跨微服务的事件传递机制未配置——在排除Axon Server后,默认的Axon事件总线是本地内存模式,无法实现命令端与查询端的事件互通。以下是具体解决方案:
1. 核心依赖调整(命令/查询端均需配置)
更新build.gradle,添加分布式事件传递所需的AMQP(以RabbitMQ为例)和JPA事件存储依赖:
// 公共Axon依赖 implementation 'org.axon-framework:axon-spring-boot-starter:4.8.2' implementation 'org.axon-framework:axon-amqp:4.8.2' implementation 'org.springframework.boot:spring-boot-starter-amqp' // 排除Axon Server连接器 exclude group: 'org.axon-framework', module: 'axon-server-connector' // 命令端额外添加JPA事件存储依赖 implementation 'org.springframework.boot:spring-boot-starter-data-jpa' implementation 'org.postgresql:postgresql' // 查询端额外添加JPA查询存储依赖 implementation 'org.springframework.boot:spring-boot-starter-data-jpa' implementation 'org.postgresql:postgresql'
2. 命令端配置(事件发布+事件溯源存储)
2.1 配置文件application.yml
spring: # 命令端数据库配置(独立端口) datasource: url: jdbc:postgresql://localhost:5432/customer-command-db username: postgres password: postgres jpa: hibernate: ddl-auto: update show-sql: true # RabbitMQ连接配置 rabbitmq: host: localhost port: 5672 username: guest password: guest axon: # 事件序列化配置(确保跨服务兼容) serializer: general: jackson events: jackson messages: jackson # AMQP事件发布配置 amqp: exchange: axon-customer-events event-exchange: axon-customer-events publisher: confirm-enabled: true # JPA事件存储配置(替换默认内存存储) eventstore: jpa: entity-manager-factory: default eventhandling: processors: default: mode: subscribing
2.2 聚合类修正
确保聚合类有无参构造函数(Axon重建聚合必需),并完善事件发布逻辑:
@Aggregate public class CustomerAggregate { @AggregateIdentifier private String customerId; private String firstname; private String lastname; // 其他字段... // 必须保留无参构造函数 public CustomerAggregate() {} @CommandHandler public CustomerAggregate(CreateCustomerCommand command) { // 参数校验 if (Objects.isNull(command.getCustomerId()) || command.getCustomerId().isEmpty()) { throw new IllegalArgumentException("Customer ID cannot be empty"); } // 映射命令到事件(简化示例,可保留原Mapper逻辑) CustomerCreatedEvent event = CustomerCreatedEvent.builder() .customerId(command.getCustomerId()) .firstname(command.getFirstname()) .lastname(command.getLastname()) .fullName(command.getFirstname() + " " + command.getLastname()) // 其他字段赋值 .build(); // 发布事件 AggregateLifecycle.apply(event); } @EventSourcingHandler public void on(CustomerCreatedEvent event) { this.customerId = event.getCustomerId(); this.firstname = event.getFirstname(); this.lastname = event.getLastname(); // 其他字段赋值 } }
3. 查询端配置(事件订阅+查询存储)
3.1 配置文件application.yml
spring: # 查询端数据库配置(独立端口) datasource: url: jdbc:postgresql://localhost:5433/customer-query-db username: postgres password: postgres jpa: hibernate: ddl-auto: update show-sql: true # RabbitMQ连接配置(与命令端一致) rabbitmq: host: localhost port: 5672 username: guest password: guest axon: # 事件序列化配置(与命令端一致) serializer: general: jackson events: jackson messages: jackson # AMQP事件订阅配置 amqp: exchange: axon-customer-events event-exchange: axon-customer-events listener: auto-start: true # 查询端无需事件溯源,禁用JPA事件存储 eventstore: jpa: enabled: false eventhandling: processors: default: mode: subscribing
3.2 事件处理器修正
确保处理器是Spring管理的Bean,完善事件持久化逻辑:
@Component public class CustomerProjection { private final CustomerRepository customerRepository; // 构造注入(避免字段注入) public CustomerProjection(CustomerRepository customerRepository) { this.customerRepository = customerRepository; } @EventHandler public void handle(CustomerCreatedEvent event) { CustomerEntity entity = new CustomerEntity(); BeanUtils.copyProperties(event, entity); // 可添加额外字段处理逻辑 customerRepository.save(entity); } }
4. 关键注意事项
- 事件类一致性:
CustomerCreatedEvent必须在命令端和查询端保持完全一致(包名、字段、序列化方式),建议抽为公共模块依赖。 - 消息中间件状态:确保RabbitMQ已启动,且配置的主机、端口、账号正确。
- 事件存储验证:命令端执行命令后,可检查
domain_event_entry表是否生成事件记录。 - 日志排查:查看查询端日志,确认是否有RabbitMQ连接失败、事件序列化错误等异常。
内容的提问来源于stack exchange,提问作者hamed
相关产品推荐
相关产品推荐

