事件驱动架构中如何判断请求已完成?求Spring Boot示例
事件驱动架构下同步订单创建的解决方案
核心思路修正
你的初始流程存在一个关键误区:事件驱动架构并非要求所有流程完全异步到底。当客户端需要同步响应时,必须结合请求ID绑定+异步状态聚合的模式,同时利用消息中间件的请求-响应能力,解决多实例、事件匹配等问题。
最优方案:请求ID绑定+共享状态缓存+异步结果聚合
关键问题解决
- 事件匹配:为每个客户端请求生成唯一
requestId,所有依赖服务的检查结果事件都携带该ID,CreateOrder服务仅处理与当前请求ID匹配的事件。 - 多实例部署:用Redis作为共享状态存储,多个CreateOrder实例均可读取/更新同一
requestId的检查状态;结合Kafka消费者组特性,确保结果消息被高效分配处理。 - RestController等待事件:通过
CompletableFuture+超时机制,阻塞等待两个检查结果聚合完成,超时则返回错误响应。
具体流程
- 客户端调用CreateOrder的REST接口,传入订单商品、配送地址等信息。
- CreateOrder生成唯一
requestId,将订单信息+requestId分别发送到inventory-check-topic和location-check-topic。 - Inventory服务消费消息完成库存检查,将携带
requestId的结果发送到inventory-check-result-topic。 - Location服务同理,将携带
requestId的地址检查结果发送到location-check-result-topic。 - CreateOrder服务订阅两个结果主题,通过
requestId匹配对应请求,更新Redis中的检查状态。 - 初始REST请求通过
CompletableFuture轮询Redis状态,聚合完成后返回订单ID或错误信息。
Spring Boot示例代码
1. 核心依赖(pom.xml)
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> </dependency> </dependencies>
2. 配置文件(application.yml)
spring: kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.springframework.kafka.support.serializer.JsonSerializer consumer: group-id: create-order-group key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer properties: spring.json.trusted.packages: "*" redis: host: localhost port: 6379
3. 消息实体类
// 订单检查请求 public class OrderCheckRequest { private String requestId; private String productId; private String destination; // 省略getter、setter、全参/无参构造方法 } // 检查结果 public class CheckResult { private String requestId; private boolean success; private String message; // 省略getter、setter、全参/无参构造方法 } // 客户端订单请求 public class OrderRequest { private String productId; private String destination; // 省略getter、setter }
4. CreateOrder控制器与服务
@RestController @RequestMapping("/orders") public class CreateOrderController { private final CreateOrderService orderService; public CreateOrderController(CreateOrderService orderService) { this.orderService = orderService; } @PostMapping public ResponseEntity<?> createOrder(@RequestBody OrderRequest orderRequest) throws Exception { String requestId = UUID.randomUUID().toString(); CompletableFuture<Boolean> checkFuture = orderService.initiateDualCheck(requestId, orderRequest); // 等待检查结果,设置10秒超时 Boolean allPass = checkFuture.get(10, TimeUnit.SECONDS); if (allPass) { String orderId = UUID.randomUUID().toString(); // 此处可添加订单持久化逻辑 return ResponseEntity.ok(Map.of("orderId", orderId, "msg", "订单创建成功")); } else { return ResponseEntity.badRequest().body(Map.of("msg", "库存不足或配送地址不支持")); } } } @Service public class CreateOrderService { private final KafkaTemplate<String, OrderCheckRequest> kafkaTemplate; private final StringRedisTemplate redisTemplate; private static final String INVENTORY_TOPIC = "inventory-check-topic"; private static final String LOCATION_TOPIC = "location-check-topic"; private static final String RESULT_KEY_PREFIX = "order-check:"; public CreateOrderService(KafkaTemplate<String, OrderCheckRequest> kafkaTemplate, StringRedisTemplate redisTemplate) { this.kafkaTemplate = kafkaTemplate; this.redisTemplate = redisTemplate; } public CompletableFuture<Boolean> initiateDualCheck(String requestId, OrderRequest orderRequest) { CompletableFuture<Boolean> future = new CompletableFuture<>(); // 初始化Redis状态:0=未开始,1=库存通过,2=地址通过,3=全部通过,-1=任意失败 redisTemplate.opsForValue().set(RESULT_KEY_PREFIX + requestId, "0", 15, TimeUnit.MINUTES); // 发送库存检查请求 OrderCheckRequest inventoryReq = new OrderCheckRequest(requestId, orderRequest.getProductId(), null); kafkaTemplate.send(INVENTORY_TOPIC, requestId, inventoryReq); // 发送地址检查请求 OrderCheckRequest locationReq = new OrderCheckRequest(requestId, null, orderRequest.getDestination()); kafkaTemplate.send(LOCATION_TOPIC, requestId, locationReq); // 启动线程轮询Redis状态(可替换为Redis Keyspace Notification优化) new Thread(() -> { try { for (int i = 0; i < 20; i++) { String status = redisTemplate.opsForValue().get(RESULT_KEY_PREFIX + requestId); if ("3".equals(status)) { future.complete(true); break; } else if ("-1".equals(status)) { future.complete(false); break; } Thread.sleep(500); } if (!future.isDone()) { future.completeExceptionally(new TimeoutException("检查超时")); } } catch (InterruptedException e) { future.completeExceptionally(e); } finally { redisTemplate.delete(RESULT_KEY_PREFIX + requestId); } }).start(); return future; } }
5. Inventory服务消费者
@Service public class InventoryConsumer { private final KafkaTemplate<String, CheckResult> kafkaTemplate; private static final String INVENTORY_RESULT_TOPIC = "inventory-check-result-topic"; public InventoryConsumer(KafkaTemplate<String, CheckResult> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } @KafkaListener(topics = "inventory-check-topic", groupId = "inventory-group") public void handleInventoryCheck(OrderCheckRequest request) { // 模拟库存检查逻辑 boolean available = "PROD001".equals(request.getProductId()); CheckResult result = new CheckResult(request.getRequestId(), available, available ? "库存充足" : "库存不足"); kafkaTemplate.send(INVENTORY_RESULT_TOPIC, request.getRequestId(), result); } }
6. Location服务消费者
@Service public class LocationConsumer { private final KafkaTemplate<String, CheckResult> kafkaTemplate; private static final String LOCATION_RESULT_TOPIC = "location-check-result-topic"; public LocationConsumer(KafkaTemplate<String, CheckResult> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } @KafkaListener(topics = "location-check-topic", groupId = "location-group") public void handleLocationCheck(OrderCheckRequest request) { // 模拟地址检查逻辑 boolean supported = "BEIJING".equals(request.getDestination()); CheckResult result = new CheckResult(request.getRequestId(), supported, supported ? "地址支持配送" : "地址不支持"); kafkaTemplate.send(LOCATION_RESULT_TOPIC, request.getRequestId(), result); } }
7. CreateOrder结果聚合消费者
@Service public class CheckResultConsumer { private final StringRedisTemplate redisTemplate; private static final String RESULT_KEY_PREFIX = "order-check:"; public CheckResultConsumer(StringRedisTemplate redisTemplate) { this.redisTemplate = redisTemplate; } @KafkaListener(topics = "inventory-check-result-topic", groupId = "create-order-group") public void handleInventoryResult(CheckResult result) { updateStatus(result.getRequestId(), result.isSuccess(), 1); } @KafkaListener(topics = "location-check-result-topic", groupId = "create-order-group") public void handleLocationResult(CheckResult result) { updateStatus(result.getRequestId(), result.isSuccess(), 2); } private void updateStatus(String requestId, boolean success, int type) { String key = RESULT_KEY_PREFIX + requestId; String currentStatus = redisTemplate.opsForValue().get(key); if (currentStatus == null) return; if (!success) { redisTemplate.opsForValue().set(key, "-1"); return; } int newStatus = Integer.parseInt(currentStatus) | type; redisTemplate.opsForValue().set(key, String.valueOf(newStatus)); } }
补充优化点
- 替代轮询:启用Redis的Keyspace Notification,当状态变化时主动通知CreateOrder服务,避免轮询消耗。
- 容错处理:为Kafka消息添加重试机制,设置消息过期时间,避免请求长期挂起。
- 分布式追踪:可引入Spring Cloud Sleuth,通过TraceId关联所有请求链路,便于问题排查。
内容的提问来源于stack exchange,提问作者Mycotina
相关产品推荐
相关产品推荐

