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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 01:31:02