下载工作台
Spring Cloud 微服务

消息驱动与 Spring Cloud Stream

试读上半部分 · 解锁后可读全文

第 6 章 · Spring Cloud Stream 与 RabbitMQ / Kafka 概念

本章目标:理解事件驱动在微服务中的角色;掌握 Spring Cloud StreamBinderBindingChannel 抽象;对比 RabbitMQApache Kafka 适用场景;在 svc-spring-demo 实现「订单创建 → 异步通知库存/日志」教学链路;了解消息可靠性、幂等消费与死信队列概念。

学时建议:5~6 小时(含 2 小时 RabbitMQ 本地联调)

前置spring-cloud-web ch01~ch05spring-boot-web ch10 @Async;消息队列基本概念。


6.1 同步 Feign 的局限与事件驱动

ch04 中 order-svc 同步调用 product-svc 校验商品,适合强一致读场景。但以下需求更适合异步消息

场景同步 Feign 问题消息方案
发邮件/短信通知拉长订单接口 RT发布事件,通知服务订阅
扣减库存(最终一致)库存服务慢或不可用拖垮下单发布 OrderCreated,inventory 异步扣
搜索索引更新非核心路径异步同步 ES
削峰填谷秒杀瞬时流量压垮 DB队列缓冲
同步(ch04):
  order-svc ──Feign──► product-svc(查价)

异步(本章):
  order-svc ──publish──► order-created ──► inventory-svc
                                      └──► notify-svc
svc-spring-demo 为虚构练习;消息 Broker 地址使用 localhostmq.example.com 占位,勿连接生产集群。

6.2 RabbitMQ 与 Kafka 概念对照

维度RabbitMQ(AMQP)Kafka
模型Exchange → Queue → ConsumerTopic → Partition → Consumer Group
吞吐量万级 QPS 足够多数业务百万级日志流
路由灵活(direct/topic/fanout)按 key 分区
消息回溯消费即删(可 TTL/ DLX)按 offset 可回溯
运维Erlang,相对轻3.x 起默认 KRaft 免 ZooKeeper(4.0 已移除 ZK),集群仍较重
教学默认本章 RabbitMQ 为主概念对比 + 配置差异

选型建议

  • 业务事件、任务队列 → RabbitMQ
  • 行为日志、大数据管道 → Kafka
  • Spring Cloud Stream 通过换 Binder 降低迁移成本

6.3 Spring Cloud Stream 核心抽象

┌─────────────────────────────────────────┐
│  order-svc 应用代码                      │
│  @Bean Supplier / Consumer / Function   │
└─────────────────┬───────────────────────┘
                  │ Binding: order-created-out
                  ▼
┌─────────────────────────────────────────┐
│  Stream Bridge / Functional Model         │
└─────────────────┬───────────────────────┘
                  │ RabbitMQ Binder
                  ▼
┌─────────────────────────────────────────┐
│  RabbitMQ Exchange: order.events          │
│  Queue: inventory.order.created           │
└─────────────────────────────────────────┘
概念说明
Binder与 MQ 打交道的适配层(rabbitkafka
Binding连接应用与中间件目的地
destination逻辑通道名,如 order-created
group消费组,竞争消费 vs 广播

Spring Cloud Stream 3.x+ 推荐 函数式编程模型SupplierConsumerFunction)。


6.4 本地启动 RabbitMQ

docker-compose.dev.yml 追加:

  rabbitmq:
    image: rabbitmq:3.13-management
    ports:
      - "5672:5672"
      - "15672:15672"
    environment:
      RABBITMQ_DEFAULT_USER: guest
      RABBITMQ_DEFAULT_PASS: guest

管理台:http://127.0.0.1:15672(guest/guest,仅本地)。


6.5 order-svc:发布订单创建事件

6.5.1 依赖

<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-stream-rabbit</artifactId>
</dependency>

Kafka 则换:

<artifactId>spring-cloud-starter-stream-kafka</artifactId>

6.5.2 事件 DTO(svc-common)

@Data
@AllArgsConstructor
@NoArgsConstructor
public class OrderCreatedEvent {
    private Long orderId;
    private Long userId;
    private Long productId;
    private Integer quantity;
    private BigDecimal amount;
    private Instant createdAt;
}

6.5.3 配置

spring:
  cloud:
    stream:
      bindings:
        orderCreated-out-0:
          destination: order.created
          content-type: application/json
      rabbit:
        bindings:
          orderCreated-out-0:
            producer:
              routing-key-expression: "'order.created'"

6.5.4 使用 StreamBridge 发送

@Service
@RequiredArgsConstructor
@Slf4j
public class OrderEventPublisher {

    private final StreamBridge streamBridge;

    private static final String BINDING = "orderCreated-out-0";

    public void publishOrderCreated(OrderCreatedEvent event) {
        boolean sent = streamBridge.send(BINDING, MessageBuilder
            .withPayload(event)
            .setHeader("eventType", "OrderCreated")
            .setHeader("X-Trace-Id", MDC.get("traceId"))
            .build());
        if (!sent) {
            log.error("Failed to publish OrderCreatedEvent orderId={}", event.getOrderId());
            throw new BusinessException(ErrorCode.MESSAGE_PUBLISH_FAILED);
        }
        log.info("Published OrderCreatedEvent orderId={}", event.getOrderId());
    }
}

OrderCreateService.create 成功落库后调用 publishOrderCreated


6.6 inventory-svc(选修模块):消费事件

新建轻量 svc-inventory 或先在 order-svc 内写 Consumer 做教学:

6.6.1 Consumer 函数

@Configuration
public class OrderEventConsumerConfig {

    @Bean
    public Consumer<OrderCreatedEvent> orderCreated() {
        return event -> {
            log.info("Received OrderCreated orderId={} productId={} qty={}",
                event.getOrderId(), event.getProductId(), event.getQuantity());
            // 扣减库存(教学:打印日志;完整逻辑需 inventory 表)
        };
    }
}

以下内容需解锁后阅读

试读已结束。解锁本章 ¥5.00,或开通年度会员畅读全部教程。
年度会员 ¥199.00/年; 小紫 AI 工作台有效会员 ¥99.00/年

正文仅在服务端鉴权后下发,未付费无法获取下半部分内容。