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

事件驱动架构中如何判断请求已完成?求Spring Boot示例

事件驱动架构下同步订单创建的解决方案

核心思路修正

你的初始流程存在一个关键误区:事件驱动架构并非要求所有流程完全异步到底。当客户端需要同步响应时,必须结合请求ID绑定+异步状态聚合的模式,同时利用消息中间件的请求-响应能力,解决多实例、事件匹配等问题。

最优方案:请求ID绑定+共享状态缓存+异步结果聚合

关键问题解决

  • 事件匹配:为每个客户端请求生成唯一requestId,所有依赖服务的检查结果事件都携带该ID,CreateOrder服务仅处理与当前请求ID匹配的事件。
  • 多实例部署:用Redis作为共享状态存储,多个CreateOrder实例均可读取/更新同一requestId的检查状态;结合Kafka消费者组特性,确保结果消息被高效分配处理。
  • RestController等待事件:通过CompletableFuture+超时机制,阻塞等待两个检查结果聚合完成,超时则返回错误响应。

具体流程

  1. 客户端调用CreateOrder的REST接口,传入订单商品、配送地址等信息。
  2. CreateOrder生成唯一requestId,将订单信息+requestId分别发送到inventory-check-topic和location-check-topic。
  3. Inventory服务消费消息完成库存检查,将携带requestId的结果发送到inventory-check-result-topic。
  4. Location服务同理,将携带requestId的地址检查结果发送到location-check-result-topic。
  5. CreateOrder服务订阅两个结果主题,通过requestId匹配对应请求,更新Redis中的检查状态。
  6. 初始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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 16:03:17