乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • 学习路径
  • 第1章 微服务与 Spring Cloud 概述

    • 微服务与 Spring Cloud 概述
  • 第2章 服务注册与发现

    • Eureka 服务注册与发现
    • Consul 服务注册与发现
  • 第3章 客户端负载均衡

    • Ribbon 客户端负载均衡
    • Spring Cloud LoadBalancer
  • 第4章 声明式服务调用

    • Feign 声明式 HTTP 客户端
    • OpenFeign 高级特性与 Spring Cloud 集成
  • 第5章 服务容错与熔断

    • Resilience4j 服务容错与熔断
    • Sentinel 流量控制与熔断降级
  • 第6章 配置中心

    • Spring Cloud Config 配置中心
  • 第7章 API 网关

    • Spring Cloud Gateway 现代化 API 网关
    • Zuul 网关(已进入维护模式)
  • 第8章 消息驱动的微服务

    • Spring Cloud Stream 消息驱动的微服务
  • 第9章 分布式追踪与监控

    • Sleuth + Zipkin 分布式链路追踪
  • 第10章 最佳实践与面试考点

    • Spring Cloud 最佳实践与面试考点

Spring Cloud Stream 消息驱动的微服务

导学

飞翔科技的订单系统最近遇到了一个问题。孔蓝提了需求:"用户下单后,要同时做三件事——扣库存、发短信通知、写审计日志。现在这三个操作串行执行,用户等 3 秒才能看到'下单成功'。"

白歌画了一张对比图:

"用消息队列解耦——这就是 Spring Cloud Stream 的价值。"白歌说。


定位与问题场景

Spring Cloud Stream 是一个消息中间件抽象层,让开发者面向统一的消息编程模型编写代码,无需关心底层是 Kafka、RabbitMQ 还是 RocketMQ。

核心概念

概念说明
Binder与消息中间件的适配层,屏蔽 RabbitMQ/Kafka 差异
Destination消息目标(Kafka 中的 Topic,RabbitMQ 中的 Exchange)
Supplier消息生产者(函数式),输出消息
Function消息处理器,接收消息 + 处理后输出
Consumer消息消费者(函数式),接收消息
StreamBridge程序化发送消息的工具类

Stream 3.x 函数式编程模型

Spring Cloud Stream 3.x 废弃了 @EnableBinding / @Input / @Output 注解式编程,全面转向函数式编程:


完整示例:飞翔科技订单异步处理

场景描述

用户下单后:

  1. 订单服务发送 OrderCreatedEvent 到消息队列
  2. 库存服务消费事件扣减库存
  3. 通知服务消费事件发送短信
  4. 审计服务消费事件记录日志

操作前后对比:

维度引入 Stream 前引入 Stream 后
下单响应时间3 秒(串行)0.5 秒(异步)
耦合度订单服务直接调用库存/通知/审计订单服务只发消息,完全解耦
削峰填谷无消息队列缓冲

示例一:订单服务——发送消息(Supplier + StreamBridge)

依赖:

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

配置:

spring:
  cloud:
    stream:
      bindings:
        orderCreated-out-0:              # 绑定名:<functionName>-<in|out>-<index>
          destination: order.created      # 对应 RabbitMQ Exchange / Kafka Topic
          content-type: application/json
        orderRefunded-out-0:
          destination: order.refunded
          content-type: application/json
  rabbitmq:
    host: localhost
    port: 5672
    username: admin
    password: admin

订单服务代码:

@RestController
public class OrderController {
    @Autowired
    private StreamBridge streamBridge;

    @PostMapping("/orders")
    public Order createOrder(@RequestBody OrderRequest request) {
        // 1. 保存订单到数据库(核心业务,同步)
        Order order = orderService.save(request);

        // 2. 发送事件(通知+审计,异步)
        OrderCreatedEvent event = new OrderCreatedEvent(
            order.getId(), order.getUserId(), order.getTotalAmount()
        );
        streamBridge.send("orderCreated-out-0", event);

        // 3. 立即返回(不等待库存扣减、短信、日志)
        return order;
    }
}

示例二:库存服务——消费消息(Consumer)

@Configuration
public class InventoryConsumer {

    @Bean
    public Consumer<OrderCreatedEvent> orderCreated() {
        return event -> {
            log.info("收到订单创建事件,扣减库存:orderId={}", event.getOrderId());
            inventoryService.deduct(event);
        };
    }
}
spring:
  cloud:
    stream:
      bindings:
        orderCreated-in-0:
          destination: order.created     # 必须与生产者 destination 一致
          group: inventory-consumer-group # 消费者组(持久化 + 负载均衡)

示例三:通知服务——消费消息(Consumer)

@Configuration
public class NotificationConsumer {

    @Bean
    public Consumer<OrderCreatedEvent> orderCreated() {
        return event -> {
            log.info("收到订单创建事件,发送短信:orderId={}", event.getOrderId());
            smsService.sendOrderConfirmation(event.getUserId(), event.getOrderId());
        };
    }
}

示例四:消息处理——Function(接收 + 转换 + 输出)

@Configuration
public class OrderProcessor {

    @Bean
    public Function<OrderCreatedEvent, AuditLog> orderAudit() {
        return event -> {
            log.info("审计订单:orderId={}", event.getOrderId());
            return new AuditLog("ORDER_CREATED", event.getOrderId(), Instant.now());
        };
    }
}
spring:
  cloud:
    stream:
      bindings:
        orderAudit-in-0:
          destination: order.created
        orderAudit-out-0:
          destination: audit.log

消费组与分区

消费组(Consumer Group):

# 库存服务部署 3 个实例,同一 group 内只有一个实例消费
spring.cloud.stream.bindings.orderCreated-in-0.group: inventory-consumer-group
# 通知服务部署 2 个实例
spring.cloud.stream.bindings.orderCreated-in-0.group: notification-consumer-group

效果:消息只会被每个 group 中的一个实例消费,不同 group 各自独立消费(类似 Kafka 的消费者组)。

分区(Partitioning):

spring:
  cloud:
    stream:
      bindings:
        orderCreated-out-0:
          producer:
            partition-key-expression: payload.orderId
            partition-count: 3

效果:相同 orderId 的消息总是路由到同一分区,保证有序处理。


易错场景

1. 消费者组缺失导致消息丢失

# ❌ 错误:未配置 group,默认匿名消费者,重启后会丢失所有未消费消息
spring.cloud.stream.bindings.orderCreated-in-0:
  destination: order.created

# ✅ 正确:配置 group,消息持久化
spring.cloud.stream.bindings.orderCreated-in-0:
  destination: order.created
  group: my-group

2. 函数命名与 binding 命名不匹配

@Bean
public Consumer<OrderCreatedEvent> processOrder() { ... }
// 期望的 binding 名:processOrder-in-0
// 如果配置文件写成了 orderCreated-in-0 → 绑定失败

3. 忘记配置 content-type

# 默认 content-type 为 application/json
# 但发送字节数组或纯文本时需要显式指定
spring.cloud.stream.bindings.myEvent-out-0.content-type: text/plain

4. DLQ(死信队列)未配置导致异常消息丢弃

# RabbitMQ 死信队列
spring:
  cloud:
    stream:
      rabbit:
        bindings:
          orderCreated-in-0:
            consumer:
              auto-bind-dlq: true
              dlq-ttl: 5000

面试考点

Spring Cloud Stream 的 Binder 是什么?

Binder 是 Spring Cloud Stream 与底层消息中间件的适配层,负责将消息通道(Destination)映射为具体的中间件概念(Kafka 的 Topic、RabbitMQ 的 Exchange/Queue)。开发者面向统一的 Supplier/Function/Consumer 编程,切换中间件只需更换 Binder 依赖,无需修改业务代码。

Consumer Group 的原理和作用?

相同的 group 内的多个实例共享消息(竞争消费,每条消息只被一个实例处理),实现负载均衡。不同的 group 各自独立消费全量消息(广播模式)。group 还提供持久化功能:消息中间件会记录每个 group 的消费偏移量,服务重启后从上次位置继续消费。不配置 group 的消费者是匿名消费者,重启后消息会丢失。

Stream 3.x 为什么废弃 @EnableBinding?

Spring Cloud Stream 3.x 全面转向函数式编程模型,原因:

  1. 注解式编程(@Input/@Output/@StreamListener)需要维护繁琐的接口声明
  2. 函数式 Bean(Supplier/Function/Consumer)更符合 Spring Boot 的 Auto-Configuration 理念
  3. 框架可以根据函数的输入输出类型自动推断消息转换逻辑
  4. 更易于单元测试

StreamBridge 与 Supplier 的区别?

  • Supplier:声明式,框架按配置的轮询间隔自动调用,适合轮询数据库、定时任务等场景
  • StreamBridge:程序化,业务代码主动调用 .send() 发送消息,适合用户操作触发的场景(如 HTTP 请求 → 发消息)
  • 大多数业务场景使用 StreamBridge,定时或流式数据源使用 Supplier

小结

Spring Cloud Stream 通过 Binder 抽象层屏蔽了消息中间件的差异,函数式编程模型让消息处理代码更简洁。消费组实现了消费端的负载均衡和持久化,StreamBridge 让业务代码可以灵活地发送消息。至此,微服务的同步调用链路和异步消息通道都已打通。下一章进入可观测性领域——分布式追踪。