如何将Axon Saga抛出的异常传递到Spring Boot REST控制器
解决Axon Saga抛出异常无法传递到REST控制器的问题
核心原因很明确:Saga是异步处理事件的,当初始的CreateOrderCommand通过sendAndWait完成后,REST请求的响应已经返回给客户端了。Saga处理OrderCreatedEvent是后续的异步操作,抛出的异常不在原请求的上下文里,自然无法被全局异常处理器捕获。
下面是几种可行的解决方案,按推荐优先级排序:
1. 前置验证:把逻辑移到命令处理阶段(最优)
如果业务允许,将产品有效性验证提前到CreateOrderCommand的聚合处理逻辑中,而不是放到Saga里。这样验证失败时,异常会直接在sendAndWait调用时抛出,能被全局异常处理器捕获。
修改OrderAggregate的命令处理逻辑:
@Aggregate public class OrderAggregate { // ... @CommandHandler public OrderAggregate(CreateOrderCommand command) { // 直接在这里调用产品服务验证有效性 try { productService.validateProduct(command.getProductId()); } catch (InvalidProductException ex) { throw new InvalidOrderException("产品ID无效: " + command.getProductId()); } // 发布OrderCreatedEvent apply(new OrderCreatedEvent(command.getOrderId(), command.getProductId())); } // ... }
这样控制器的sendAndWait会直接抛出InvalidOrderException,全局异常处理器就能正常捕获并返回错误响应。
2. 异步反馈:通过领域事件+查询模型通知客户端
这是CQRS架构下的标准异步处理方式,适合不需要同步响应的场景:
- 在Saga捕获异常时,发布一个
OrderFailedEvent,包含订单ID和错误信息:public class OrderSaga { // ... @StartSaga @SagaEventHandler(associationProperty = "orderId") public void handle(OrderCreatedEvent orderCreatedEvent) { // ... try { commandGateway.sendAndWait(reserveProductCommand); } catch (CommandExecutionException ex) { rollbackReservations(); // 发布失败事件,替代直接抛异常 eventGateway.publish(new OrderFailedEvent(orderCreatedEvent.getOrderId(), "产品ID无效: " + productId)); // 结束Saga endSaga(); } // ... } // ... } - 创建一个订单状态的查询投影(Projection),监听
OrderCreatedEvent和OrderFailedEvent,维护订单的状态和错误信息:@Projection public class OrderStatusProjection { private final OrderStatusRepository repository; public OrderStatusProjection(OrderStatusRepository repository) { this.repository = repository; } @EventHandler public void on(OrderCreatedEvent event) { repository.save(new OrderStatus(event.getOrderId(), "CREATED", null)); } @EventHandler public void on(OrderFailedEvent event) { OrderStatus status = repository.findById(event.getOrderId()).orElseThrow(); status.setStatus("FAILED"); status.setErrorMessage(event.getErrorMessage()); repository.save(status); } } - 修改控制器,返回订单ID给客户端,让客户端通过轮询或WebSocket查询订单状态:
@PostMapping public ResponseEntity<String> createOrder(@RequestBody CreateOrderDto createOrderDto){ String orderId = UUID.randomUUID().toString(); CreateOrderCommand createOrderCommand = new CreateOrderCommand(orderId, createOrderDto.getProductId()); commandGateway.send(createOrderCommand); // 返回订单ID,让客户端查询状态 return new ResponseEntity<>(orderId, HttpStatus.ACCEPTED); } @GetMapping("/{orderId}/status") public ResponseEntity<OrderStatus> getOrderStatus(@PathVariable String orderId){ OrderStatus status = orderStatusRepository.findById(orderId).orElseThrow(); return ResponseEntity.ok(status); }
3. 同步等待:通过请求追踪关联Saga错误到初始请求
如果必须同步返回响应,可以通过请求ID追踪的方式实现:
- 发送初始命令时,添加唯一的
requestId作为元数据:@PostMapping public ResponseEntity<ReadOrderDto> createOrder(@RequestBody CreateOrderDto createOrderDto){ String requestId = UUID.randomUUID().toString(); CreateOrderCommand createOrderCommand = new CreateOrderCommand(/* 参数 */); // 添加元数据 MetaData metaData = MetaData.with("requestId", requestId); // 使用带回调的send方法,而不是sendAndWait CompletableFuture<ReadOrderDto> future = new CompletableFuture<>(); commandGateway.send(createOrderCommand, metaData, new CommandCallback<>() { @Override public void onResult(CommandMessage<?> commandMessage, CommandResultMessage<?> commandResultMessage) { if (commandResultMessage.isSuccess()) { future.complete(createOrderCommand.toToReadProductDto()); } else { future.completeExceptionally(commandResultMessage.exceptionResult()); } } }); // 注册一个监听器,监听Saga抛出的错误事件(需要自定义事件) errorEventListener.registerRequest(requestId, future); try { // 设置超时时间 ReadOrderDto result = future.get(10, TimeUnit.SECONDS); return ResponseEntity.ok(result); } catch (Exception e) { if (e.getCause() instanceof InvalidOrderException) { return ResponseEntity.badRequest().body(e.getCause().getMessage()); } return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).build(); } } - 在Saga中获取
requestId元数据,抛出异常时发布带requestId的OrderProcessingFailedEvent:public class OrderSaga { // ... @StartSaga @SagaEventHandler(associationProperty = "orderId", metaDataProperty = "requestId") public void handle(OrderCreatedEvent orderCreatedEvent, @MetaDataValue("requestId") String requestId) { // 关联requestId到Saga associateWith("requestId", requestId); // ... try { commandGateway.sendAndWait(reserveProductCommand); } catch (CommandExecutionException ex) { rollbackReservations(); eventGateway.publish(new OrderProcessingFailedEvent(requestId, "产品ID无效: " + productId)); endSaga(); } // ... } // ... } - 实现一个
ErrorEventListener,监听OrderProcessingFailedEvent,找到对应的CompletableFuture并触发异常:@Component public class ErrorEventListener { private final Map<String, CompletableFuture<?>> requestMap = new ConcurrentHashMap<>(); public void registerRequest(String requestId, CompletableFuture<?> future) { requestMap.put(requestId, future); } @EventHandler public void on(OrderProcessingFailedEvent event) { CompletableFuture<?> future = requestMap.remove(event.getRequestId()); if (future != null) { future.completeExceptionally(new InvalidOrderException(event.getErrorMessage())); } } }
这种方式需要处理超时和并发问题,适合必须同步响应的场景,但复杂度较高。
额外注意:Saga的异常重试配置
Axon默认会重试Saga的事件处理,如果不希望无限重试,需要配置重试策略或在异常后结束Saga:
- 在Saga的异常处理方法上添加
@EndSaga,或者手动调用endSaga() - 通过
SagaConfiguration配置重试次数:@Bean public SagaConfiguration<OrderSaga> orderSagaConfiguration() { return SagaConfiguration.trackingSagaManager(OrderSaga.class, config -> config.configureTrackingProcessor(processor -> processor.registerHandlerInterceptor(new RetryingErrorHandler( RetryConfig.custom().maxAttempts(3).build() )) ) ); }
内容的提问来源于stack exchange,提问作者Hadi Rifaii
相关产品推荐
相关产品推荐

