微服务通信:同步不够用了之后

最近在项目里又碰了一次微服务通信,有几处判断值得留下来。

项目背景

项目最开始是单体架构,一个 Spring Boot 应用,所有功能都在一个进程里。后来业务增长,团队规模扩大,决定拆分成微服务:

  • 订单服务:处理订单创建、支付、发货
  • 库存服务:管理商品库存、库存锁定
  • 用户服务:用户信息、积分、优惠券
  • 通知服务:短信、邮件、推送通知

拆分是好事,但问题来了:订单创建要扣库存、扣积分、发通知。这些操作谁先谁后?失败了怎么办?接口调不通怎么办?

最初用的是最简单的 HTTP 同步调用:

// 订单服务创建订单
@PostMapping("/orders")
public Order createOrder(@RequestBody CreateOrderRequest request) {
    // 1. 调用库存服务锁定库存
    LockStockRequest lockRequest = new LockStockRequest(
        request.getProductId(), request.getQuantity()
    );
    ResponseEntity<LockStockResponse> lockResponse = restTemplate.postForEntity(
        "http://inventory-service/api/inventory/lock",
        lockRequest,
        LockStockResponse.class
    );

    if (!lockResponse.getStatusCode().is2xxSuccessful()) {
        throw new RuntimeException("锁定库存失败");
    }

    // 2. 创建订单
    Order order = orderRepository.save(new Order(request));

    // 3. 调用用户服务扣减积分
    DeductPointsRequest pointsRequest = new DeductPointsRequest(
        request.getUserId(), order.getPoints()
    );
    ResponseEntity<Void> pointsResponse = restTemplate.postForEntity(
        "http://user-service/api/users/points/deduct",
        pointsRequest,
        Void.class
    );

    if (!pointsResponse.getStatusCode().is2xxSuccessful()) {
        // 上面订单已经创建了,这里失败怎么办?
        throw new RuntimeException("扣减积分失败");
    }

    // 4. 调用通知服务发送短信
    SendSmsRequest smsRequest = new SendSmsRequest(
        request.getMobile(), "订单创建成功"
    );
    restTemplate.postForEntity(
        "http://notification-service/api/notifications/sms",
        smsRequest,
        Void.class
    );

    return order;
}

这段代码看起来没什么问题,但线上运行一周就暴露了一堆毛病。

为了更直观地看清楚同步调用的问题,我用时序图梳理一下调用链路:

sequenceDiagram participant Client participant OrderService participant InventoryService participant UserService participant NotificationService Client->>OrderService: 创建订单 OrderService->>InventoryService: 锁定库存 (200ms) InventoryService-->>OrderService: 返回 OrderService->>OrderService: 创建订单 (50ms) OrderService->>UserService: 扣减积分 (150ms) UserService-->>OrderService: 返回 OrderService->>NotificationService: 发送短信 (100ms) NotificationService-->>OrderService: 返回 OrderService-->>Client: 返回订单 Note over Client,NotificationService: 总耗时 500ms,任意环节失败整个请求失败

HTTP 同步调用的坑

坑一:调用链路太长,性能差

一个订单创建请求,要串行调用三个服务,每个服务响应时间加起来就是总耗时:

  • 库存锁定:200ms
  • 订单创建:50ms
  • 积分扣减:150ms
  • 短信通知:100ms

总共 500ms,这还是理想情况。一旦某个服务慢了,整个请求就被拖慢。

坑二:部分失败难以处理

最头疼的是部分失败。比如订单创建成功了,但积分扣减失败。这时候怎么办?

  • 回滚订单?订单号已经返回给用户了
  • 重试积分扣减?但订单状态已经变了
  • 记录失败日志人工处理?成本太高

最后搞了个复杂的补偿逻辑,代码越写越乱:

@PostMapping("/orders")
public Order createOrder(@RequestBody CreateOrderRequest request) {
    Order order = null;

    try {
        // 锁定库存
        lockStock(request);
        // 创建订单
        order = createOrder(request);
        // 扣减积分
        deductPoints(request);
    } catch (Exception e) {
        // 补偿:如果订单创建了但积分扣减失败,回滚订单
        if (order != null) {
            cancelOrder(order.getId());
            unlockStock(request);
        } else {
            unlockStock(request);
        }
        throw e;
    }

    // 短信通知异步发,失败了不影响主流程
    CompletableFuture.runAsync(() -> {
        try {
            sendSms(request);
        } catch (Exception e) {
            log.error("发送短信失败", e);
        }
    });

    return order;
}

坑三:服务不稳定影响主流程

通知服务经常因为短信通道问题不稳定,拖累了整个订单创建。把通知改成异步后好一些,但库存服务一挂,订单就创建不了。

测试环境故意停掉库存服务,订单接口直接报错,用户体验很差。


引入 gRPC

HTTP JSON 这种方式虽然简单,但性能确实有问题。于是决定试试 gRPC。

gRPC 的优势

gRPC 使用 Protocol Buffers 序列化,比 JSON 更紧凑,性能更好。而且支持流式调用、双向流。

定义 proto 文件

syntax = "proto3";

package inventory;

service InventoryService {
    rpc LockStock(LockStockRequest) returns (LockStockResponse);
    rpc UnlockStock(UnlockStockRequest) returns (UnlockStockResponse);
}

message LockStockRequest {
    string product_id = 1;
    int32 quantity = 2;
    string order_id = 3;
}

message LockStockResponse {
    bool success = 1;
    string message = 2;
}

message UnlockStockRequest {
    string order_id = 1;
}

message UnlockStockResponse {
    bool success = 1;
}

Java 客户端实现

@Service
public class InventoryGrpcClient {

    private final InventoryServiceGrpc.InventoryServiceBlockingStub blockingStub;

    public InventoryGrpcClient(ManagedChannel channel) {
        this.blockingStub = InventoryServiceGrpc.newBlockingStub(channel);
    }

    public LockStockResponse lockStock(String productId, int quantity, String orderId) {
        LockStockRequest request = LockStockRequest.newBuilder()
            .setProductId(productId)
            .setQuantity(quantity)
            .setOrderId(orderId)
            .build();

        inventory.LockStockResponse response = blockingStub.lockStock(request);
        return new LockStockResponse(response.getSuccess(), response.getMessage());
    }
}

性能提升

换成 gRPC 后,接口响应时间明显下降:

  • HTTP JSON:500ms
  • gRPC:320ms

提升了 35%,主要是序列化更快、连接复用更好。


gRPC 踩过的坑

坑一:服务发现麻烦

gRPC 客户端需要知道服务端的 IP 和端口,这在容器环境下很麻烦。最后用了 gRPC Name Resolver + Kubernetes Service:

ManagedChannel channel = ManagedChannelBuilder.forTarget(
    "dns:///inventory-service:9090"
)
    .usePlaintext()
    .build();

坑二:负载均衡不均衡

gRPC 的 HTTP/2 连接是长连接,默认负载均衡策略可能导致连接不均衡。得在客户端配置 round_robin 策略:

ManagedChannel channel = ManagedChannelBuilder.forTarget(
    "dns:///inventory-service:9090"
)
    .defaultLoadBalancingPolicy("round_robin")
    .usePlaintext()
    .build();

坑三:调试困难

HTTP 可以用 curl、浏览器直接调,gRPC 需要专门的工具。最后用了 grpcurl 命令行工具:

grpcurl -plaintext inventory-service:9090 list
grpcurl -plaintext inventory-service:9090 inventory.InventoryService/LockStock

坑四:跨语言协作

gRPC 确实跨语言友好,但团队有人写 Python,有人写 Java,proto 文件版本管理成了问题。最后建了个专门的 proto repo,统一管理版本。


引入消息队列

gRPC 解决了性能问题,但同步调用的本质问题没解决:一个服务挂了,整个链路就断了。

于是决定引入消息队列,把强依赖改成弱依赖。

为什么选 RabbitMQ

选 RabbitMQ 主要考虑:

  1. 延迟低:内部系统,不追求分布式一致性,延迟比 Kafka 低
  2. 功能全:死信队列、延迟消息、消息确认都有
  3. 运维简单:团队已经用过,比较熟悉
  4. Spring Boot 支持好:Spring AMQP 封装完善

架构改造

改造思路:订单创建后发消息到 RabbitMQ,库存服务、用户服务、通知服务各自消费消息:

graph TD A[订单服务] -->|订单创建事件| B[RabbitMQ] B -->|消费| C[库存服务] B -->|消费| D[用户服务] B -->|消费| E[通知服务] C -->|锁定库存成功| B D -->|扣减积分成功| B E -->|发送通知成功| B

消息定义

public class OrderCreatedEvent {
    private String orderId;
    private String userId;
    private String productId;
    private int quantity;
    private BigDecimal amount;
    private LocalDateTime createdAt;

    // getters and setters
}

订单服务发送消息

@Service
public class OrderService {

    private final RabbitTemplate rabbitTemplate;

    public Order createOrder(CreateOrderRequest request) {
        // 创建订单
        Order order = orderRepository.save(new Order(request));

        // 发送订单创建事件
        OrderCreatedEvent event = new OrderCreatedEvent(
            order.getId(),
            request.getUserId(),
            request.getProductId(),
            request.getQuantity(),
            order.getAmount(),
            LocalDateTime.now()
        );

        rabbitTemplate.convertAndSend(
            "order-exchange",
            "order.created",
            event
        );

        return order;
    }
}

库存服务消费消息

@Service
public class InventoryMessageConsumer {

    private final InventoryService inventoryService;

    @RabbitListener(queues = "inventory.lock-stock-queue")
    public void handleOrderCreated(OrderCreatedEvent event) {
        try {
            inventoryService.lockStock(
                event.getProductId(),
                event.getQuantity(),
                event.getOrderId()
            );
        } catch (InsufficientStockException e) {
            // 库存不足,发送到死信队列人工处理
            throw new AmqpRejectAndDontRequeueException(e);
        } catch (Exception e) {
            // 其他异常,稍后重试
            throw e;
        }
    }
}

配置死信队列

spring:
  rabbitmq:
    host: rabbitmq
    port: 5672
    username: guest
    password: guest
    listener:
      simple:
        acknowledge-mode: manual
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1000
          multiplier: 2.0

消息队列踩过的坑

坑一:消息丢失

RabbitMQ 默认不持久化消息,节点重启就丢了。得配置消息持久化:

@Bean
public Queue lockStockQueue() {
    return QueueBuilder.durable("inventory.lock-stock-queue")
        .withArgument("x-dead-letter-exchange", "order-dlx")
        .withArgument("x-dead-letter-routing-key", "order.created.dlq")
        .build();
}

还有消息确认,消费成功才 ack:

@RabbitListener(queues = "inventory.lock-stock-queue")
public void handleOrderCreated(
    OrderCreatedEvent event,
    Channel channel,
    @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag
) {
    try {
        inventoryService.lockStock(event.getProductId(), event.getQuantity(), event.getOrderId());
        channel.basicAck(deliveryTag, false);
    } catch (Exception e) {
        channel.basicNack(deliveryTag, false, true);
        throw e;
    }
}

坑二:重复消费

消息队列不保证 exactly-once,可能重复消费。消费者要幂等:

@Service
public class InventoryService {

    @Transactional
    public void lockStock(String productId, int quantity, String orderId) {
        // 先查是否已经锁定过
        StockLock existingLock = stockLockRepository.findByOrderId(orderId);
        if (existingLock != null) {
            return; // 幂等:已经锁定过就直接返回
        }

        // 锁定库存
        StockLock lock = new StockLock(productId, quantity, orderId);
        stockLockRepository.save(lock);
    }
}

或者用 Redis 分布式锁:

public void lockStock(String productId, int quantity, String orderId) {
    String lockKey = "stock:lock:" + orderId;
    Boolean locked = redisTemplate.opsForValue().setIfAbsent(
        lockKey, "1", 30, TimeUnit.SECONDS
    );

    if (!locked) {
        log.info("库存锁定已存在,orderId: {}", orderId);
        return;
    }

    try {
        // 实际锁定逻辑
        doLockStock(productId, quantity, orderId);
    } finally {
        redisTemplate.delete(lockKey);
    }
}

坑三:消息积压

活动期间订单激增,消息堆积在队列里,消费者处理不过来。最后加了水平扩容和优先级队列:

@Bean
public Queue lockStockQueue() {
    return QueueBuilder.durable("inventory.lock-stock-queue")
        .withArgument("x-max-priority", 10)
        .build();
}
rabbitTemplate.convertAndSend(
    "order-exchange",
    "order.created",
    event,
    message -> {
        message.getMessageProperties().setPriority(5);
        return message;
    }
);

坑四:事务问题

订单创建和发消息不在同一个事务里,可能订单创建了但消息没发出去。最后用 Spring 事务管理:

@Transactional
public Order createOrder(CreateOrderRequest request) {
    Order order = orderRepository.save(new Order(request));

    // 事务提交后发消息
    TransactionSynchronizationManager.registerSynchronization(
        new TransactionSynchronization() {
            @Override
            public void afterCommit() {
                rabbitTemplate.convertAndSend(
                    "order-exchange",
                    "order.created",
                    new OrderCreatedEvent(order)
                );
            }
        }
    );

    return order;
}

最终架构

折腾了几个月,最终的通信架构是这样的:

graph TD A[客户端] -->|HTTP/JSON| B[API 网关] B -->|gRPC| C[订单服务] C -->|发送事件| D[RabbitMQ] D -->|消费| E[库存服务] D -->|消费| F[用户服务] D -->|消费| G[通知服务] C -->|gRPC 同步调用| H[认证服务] C -->|gRPC 同步调用| I[风控服务]

核心链路用消息队列异步处理,实时性要求高的用 gRPC 同步调用:

  • 异步处理:库存锁定、积分扣减、通知发送
  • 同步调用:用户认证、风控检查

性能对比

不同方案的性能数据(订单创建接口):

方案P95 延迟P99 延迟可用性
HTTP JSON 同步500ms800ms99.5%
gRPC 同步320ms500ms99.7%
消息队列异步50ms80ms99.9%

下图把三种方案的 P95/P99 延迟和可用性放在一起,能更直观地看到从同步到异步的跃迁幅度。

订单创建接口:HTTP JSON、gRPC 与消息队列异步方案的 P95/P99 延迟及可用性对比

换成消息队列后,订单创建的 P95 延迟从 500ms 降到 50ms,可用性从 99.5% 提升到 99.9%。库存锁定、积分扣减这些操作异步处理,不影响主流程。


实践经验

这次通信模式改造,有几个经验值得总结:

  1. 不要一刀切:同步、异步各有适用场景,实时性要求高的用同步,可以延迟的用异步
  2. 消息队列不是万能的:带来了新的复杂性,要处理消息丢失、重复消费、积压问题
  3. 监控很重要:消息队列的监控比 HTTP 复杂,要盯住队列积压、消费延迟、错误率
  4. 幂等设计要提前:重复消费是常态,幂等设计要提前想清楚

写在最后

微服务通信这条路,从 HTTP 到 gRPC,再到消息队列,走得不容易。但每一步都是问题驱动,不是跟风。

解决了

  • 性能问题,接口响应时间大幅下降
  • 可用性问题,一个服务挂了不影响其他服务
  • 扩展性问题,可以独立扩容消费者

带来了

  • 复杂性提升,要处理消息丢失、重复消费
  • 调试难度增加,消息链路不如调用链路直观
  • 运维成本增加,要监控消息队列的指标

如果你的微服务还都是同步调用,而且已经感受到性能和可用性问题,可以考虑引入消息队列。但引入之前先想清楚:哪些操作可以异步?怎么保证最终一致性?团队能力能不能跟上?

毕竟,技术是为了解决问题,不是为了堆砌架构。


这次通信模式改造花了三个月,从 HTTP 到 gRPC,再到消息队列。改造完成后,订单创建的 P95 延迟从 500ms 降到 50ms,可用性从 99.5% 提升到 99.9%。

版权声明: 本文首发于 指尖魔法屋-微服务通信:同步不够用了之后https://blog.thinkmoon.cn/post/130-microservices-communication-sync-async-practice/) 转载或引用必须申明原指尖魔法屋来源及源地址!