Axon Saga中ProductReservationService订单总价变量异常求助
问题描述
我在Saga类中拆分出多个服务分别处理不同逻辑,其他功能均正常,但ProductReservationService中的orderTotalPrice变量存在异常:
- 调试时,
handle(ProductReservedEvent)方法中该变量能正确累加订单总价,但在OrderCreationHandler中获取该变量用于发送订阅查询更新时,得到的值为0; - 若移除
processOrder方法中的productReservationService.resetOrderTotalPrice()语句,该变量会累计应用中所有订单的价格(即使重启应用仍会保留)。
相关代码
ProductReservationService.java
import com.jchaaban.common.command.ReserveProductCommand; import com.jchaaban.common.command.UnreserveProductCommand; import com.jchaaban.common.event.ProductReservedEvent; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.axonframework.commandhandling.CommandExecutionException; import org.axonframework.commandhandling.gateway.CommandGateway; import org.axonframework.eventhandling.EventHandler; import org.springframework.stereotype.Service; import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.stream.Collectors; @Slf4j @Service public class ProductReservationService { private final transient CommandGateway commandGateway; private final Map<String, Integer> successfullyReservedProducts = new HashMap<>(); @Getter private double orderTotalPrice = 0; @EventHandler public void handle(ProductReservedEvent event) { log.info("Handling ProductReservedEvent for productId: {}, productPrice: {}, quantity: {}", event.getProductId(), event.getProductPrice(), event.getQuantity()); orderTotalPrice += event.getProductPrice().doubleValue() * event.getQuantity();; log.info("Updated orderPrice: {}", orderTotalPrice); } public ProductReservationService(CommandGateway commandGateway) { this.commandGateway = commandGateway; } public boolean reserveProducts(List<String> productIds) { Map<String, Integer> productOccurrences = countProductOccurrences(productIds); return productOccurrences.entrySet().stream().allMatch(entry -> reserveProduct(entry.getKey(), entry.getValue())); } public void rollbackReservations() { log.info("Rolling back product reservations."); successfullyReservedProducts.forEach((productId, quantity) -> commandGateway.send(new UnreserveProductCommand(productId, quantity))); successfullyReservedProducts.clear(); orderTotalPrice = 0; } public void resetOrderTotalPrice(){ orderTotalPrice = 0; } private Map<String, Integer> countProductOccurrences(List<String> productIds) { return productIds.stream() .collect(Collectors.groupingBy(productId -> productId, Collectors.summingInt(productId -> 1))); } private boolean reserveProduct(String productId, int quantity) { try { commandGateway.sendAndWait(new ReserveProductCommand(productId, quantity)); successfullyReservedProducts.put(productId, quantity); } catch (CommandExecutionException exception) { rollbackReservations(); return false; } return true; } }
OrderCreationHandler.java
import com.jchaaban.common.dto.ReadPaymentDetailsDto; import com.jchaaban.ordersservice.saga.service.*; import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import java.util.List; @Slf4j @Component @Data public class OrderCreationHandler { private static final String USER_NOT_FOUND_TEMPLATE = "User with ID: %s not found."; private static final String PAYMENT_DETAILS_NOT_FOUND_TEMPLATE = "Payment details for user with ID: %s not found."; private static final String INSUFFICIENT_BALANCE = "Insufficient balance."; private static final String ORDER_SUCCESS_TEMPLATE = "Order processed successfully. Total cost: %s."; private final transient OrderSummaryEmitter orderSummaryEmitter; private final transient ProductReservationService productReservationService; private final transient UserQueryService userQueryService; private final transient OrderCancellationService orderCancellationService; private final transient PaymentService paymentService; public OrderCreationHandler( OrderSummaryEmitter orderSummaryEmitter, ProductReservationService productReservationService, UserQueryService userQueryService, OrderCancellationService orderCancellationService, PaymentService paymentService ) { this.orderSummaryEmitter = orderSummaryEmitter; this.productReservationService = productReservationService; this.userQueryService = userQueryService; this.orderCancellationService = orderCancellationService; this.paymentService = paymentService; } public void processOrder(String orderId, String userId, List<String> productIds) { userQueryService.fetchUserInformation(userId).ifPresentOrElse( user -> processPaymentDetails(orderId, userId, productIds), () -> orderSummaryEmitter.emitOrderSummary(orderId, String.format(USER_NOT_FOUND_TEMPLATE, userId)) ); } private void processPaymentDetails(String orderId, String userId, List<String> productIds) { userQueryService.fetchUserPaymentDetails(userId).ifPresentOrElse( paymentDetails -> processOrder(orderId, productIds, paymentDetails), () -> orderSummaryEmitter.emitOrderSummary(orderId, String.format(PAYMENT_DETAILS_NOT_FOUND_TEMPLATE, userId)) ); } private void processOrder(String orderId, List<String> productIds, ReadPaymentDetailsDto paymentDetails) { if (!productReservationService.reserveProducts(productIds)) { orderCancellationService.cancelOrder(orderId, "Product reservation failed"); return; } double orderPrice = productReservationService.getOrderTotalPrice(); if (paymentDetails.getBalance() < orderPrice) { orderCancellationService.cancelOrder(orderId, INSUFFICIENT_BALANCE); return; } orderSummaryEmitter.emitOrderSummary(orderId, String.format(ORDER_SUCCESS_TEMPLATE, orderPrice)); paymentService.updatePaymentDetailsBalance(paymentDetails.getPaymentDetailsId(), orderPrice); productReservationService.resetOrderTotalPrice(); } }
OrderSaga.java
import com.jchaaban.common.event.ProductReservedEvent; import com.jchaaban.ordersservice.event.OrderCreatedEvent; import lombok.NoArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.axonframework.eventhandling.EventHandler; import org.axonframework.modelling.saga.SagaEventHandler; import org.axonframework.modelling.saga.StartSaga; import org.axonframework.spring.stereotype.Saga; import org.springframework.beans.factory.annotation.Autowired; @Saga // is already a spring component @Slf4j @NoArgsConstructor public class OrderSaga { private String orderId; @Autowired private transient OrderCreationHandler orderCreationHandler; @StartSaga @SagaEventHandler(associationProperty = "orderId") public void handleOrderCreated(OrderCreatedEvent event) { log.info("Handling OrderCreatedEvent for orderId: {}", event.getOrderId()); this.orderId = event.getOrderId(); orderCreationHandler.processOrder(orderId, event.getUserId(), event.getProductIds()); } }
日志信息
2023-10-30T16:15:47.496+01:00 INFO 19832 --- [agaProcessor]-0] c.jchaaban.ordersservice.saga.OrderSaga : Handling OrderCreatedEvent for orderId: 4ea2f30f-6a7a-4886-a99d-58558ef1f769 Hibernate: update token_entry set token=?,token_type=?,timestamp=? where owner=? and processor_name=? and segment=? 2023-10-30T16:15:47.550+01:00 INFO 19832 --- [saga.service]-0] c.j.o.s.s.ProductReservationService : Handling ProductReservedEvent for productId: cb70a736-05fd-4091-8441-d83eaa117673, productPrice: 33, quantity: 2 2023-10-30T16:15:47.550+01:00 INFO 19832 --- [saga.service]-0] c.j.o.s.s.ProductReservationService : Updated orderPrice: 66.0 Hibernate: select nextval('association_value_entry_seq')
问题分析
- 单例服务的共享状态冲突:
ProductReservationService是Spring默认的单例组件,所有Saga实例共享同一个服务实例。这导致orderTotalPrice成为全局变量,多个订单请求会互相覆盖该值;同时Axon的Saga持久化机制会保留服务状态,重启后依然会累计历史订单价格。 - 事件处理与命令执行的线程异步性:
reserveProducts中使用commandGateway.sendAndWait()同步发送命令,但ProductReservedEvent的处理是在Axon的事件处理线程中执行的。当命令调用返回时,对应的事件可能还未被处理,此时orderTotalPrice还未更新,因此获取到的值为0。
解决方案
1. 将状态与Saga实例绑定
每个Saga实例应该持有自己的订单状态,避免全局共享。修改OrderSaga,添加专属的状态变量和事件处理方法:
@Saga @Slf4j @NoArgsConstructor public class OrderSaga { private String orderId; private double orderTotalPrice = 0; private Map<String, Integer> successfullyReservedProducts = new HashMap<>(); @Autowired private transient OrderCreationHandler orderCreationHandler; @Autowired private transient CommandGateway commandGateway; @StartSaga @SagaEventHandler(associationProperty = "orderId") public void handleOrderCreated(OrderCreatedEvent event) { log.info("Handling OrderCreatedEvent for orderId: {}", event.getOrderId()); this.orderId = event.getOrderId(); orderCreationHandler.processOrder(this, orderId, event.getUserId(), event.getProductIds()); } @SagaEventHandler(associationProperty = "orderId") public void handleProductReserved(ProductReservedEvent event) { log.info("Handling ProductReservedEvent for productId: {}, productPrice: {}, quantity: {}", event.getProductId(), event.getProductPrice(), event.getQuantity()); this.orderTotalPrice += event.getProductPrice().doubleValue() * event.getQuantity(); log.info("Updated orderPrice for order {}: {}", this.orderId, this.orderTotalPrice); this.successfullyReservedProducts.put(event.getProductId(), event.getQuantity()); } public double getOrderTotalPrice() { return orderTotalPrice; } public void rollbackReservations() { log.info("Rolling back product reservations for order {}", this.orderId); successfullyReservedProducts.forEach((productId, quantity) -> commandGateway.send(new UnreserveProductCommand(productId, quantity))); successfullyReservedProducts.clear(); this.orderTotalPrice = 0; } public void resetOrderTotalPrice() { this.orderTotalPrice = 0; this.successfullyReservedProducts.clear(); } public String getOrderId() { return orderId; } }
2. 改造ProductReservationService为无状态服务
移除服务中的全局状态,只保留命令发送的核心逻辑,状态管理交给Saga:
@Slf4j @Service public class ProductReservationService { private final transient CommandGateway commandGateway; public ProductReservationService(CommandGateway commandGateway) { this.commandGateway = commandGateway; } public boolean reserveProducts(OrderSaga saga, List<String> productIds) { Map<String, Integer> productOccurrences = countProductOccurrences(productIds); return productOccurrences.entrySet().stream().allMatch(entry -> reserveProduct(saga, entry.getKey(), entry.getValue())); } private boolean reserveProduct(OrderSaga saga, String productId, int quantity) { try { // 发送命令时携带orderId,确保事件能关联到对应Saga commandGateway.sendAndWait(new ReserveProductCommand(productId, quantity, saga.getOrderId())); return true; } catch (CommandExecutionException exception) { saga.rollbackReservations(); return false; } } private Map<String, Integer> countProductOccurrences(List<String> productIds) { return productIds.stream() .collect(Collectors.groupingBy(productId -> productId, Collectors.summingInt(productId -> 1))); } }
3. 修改OrderCreationHandler,使用Saga实例状态
调整方法参数,传入Saga实例,从Saga中获取订单总价而非全局服务:
@Slf4j @Component @Data public class OrderCreationHandler { // 常量定义保持不变 private final transient OrderSummaryEmitter orderSummaryEmitter; private final transient ProductReservationService productReservationService; private final transient UserQueryService userQueryService; private final transient OrderCancellationService orderCancellationService; private final transient PaymentService paymentService; // 构造方法保持不变 public void processOrder(OrderSaga saga, String orderId, String userId, List<String> productIds) { userQueryService.fetchUserInformation(userId).ifPresentOrElse( user -> processPaymentDetails(saga, orderId, userId, productIds), () -> orderSummaryEmitter.emitOrderSummary(orderId, String.format(USER_NOT_FOUND_TEMPLATE, userId)) ); } private void processPaymentDetails(OrderSaga saga, String orderId, String userId, List<String> productIds) { userQueryService.fetchUserPaymentDetails(userId).ifPresentOrElse( paymentDetails -> processOrder(saga, orderId, productIds, paymentDetails), () -> orderSummaryEmitter.emitOrderSummary(orderId, String.format(PAYMENT_DETAILS_NOT_FOUND_TEMPLATE, userId)) ); } private void processOrder(OrderSaga saga, String orderId, List<String> productIds, ReadPaymentDetailsDto paymentDetails) { if (!productReservationService.reserveProducts(saga, productIds)) { orderCancellationService.cancelOrder(orderId, "Product reservation failed"); return; } double orderPrice = saga.getOrderTotalPrice(); if (paymentDetails.getBalance() < orderPrice) { orderCancellationService.cancelOrder(orderId, INSUFFICIENT_BALANCE); saga.rollbackReservations(); return; } orderSummaryEmitter.emitOrderSummary(orderId, String.format(ORDER_SUCCESS_TEMPLATE, orderPrice)); paymentService.updatePaymentDetailsBalance(paymentDetails.getPaymentDetailsId(), orderPrice); saga.resetOrderTotalPrice(); } }
4. 完善事件与命令的关联
修改ReserveProductCommand和ProductReservedEvent,添加orderId字段,确保事件能正确路由到对应的Saga实例:
// ReserveProductCommand public class ReserveProductCommand { private final String productId; private final int quantity; private final String orderId; public ReserveProductCommand(String productId, int quantity, String orderId) { this.productId = productId; this.quantity = quantity; this.orderId = orderId; } // getter方法 } // ProductReservedEvent public class ProductReservedEvent { private final String productId; private final int quantity; private final BigDecimal productPrice; private final String orderId; public ProductReservedEvent(String productId, int quantity, BigDecimal productPrice, String orderId) { this.productId = productId; this.quantity = quantity; this.productPrice = productPrice; this.orderId = orderId; } // getter方法 }
内容的提问来源于stack exchange,提问作者Hadi Rifaii
相关产品推荐
相关产品推荐

