Stream 实战:绑定 RabbitMQ 收发消息

本章完成一次真实的"订单事件"异步流转:order-service 发布消息,inventory-service 订阅并处理,中间件用 RabbitMQ。沿用第 01 章版本锚点与第 19 章的函数式模型,全程不出现 RabbitMQ 厂商 API。

前置:启动 RabbitMQ

本机安装 RabbitMQ(官方安装包或 Docker 镜像均可),默认端口 5672 提供 AMQP 协议,15672 是管理控制台。本地默认账号 guest/guest,仅允许 localhost 访问;生产环境务必改密并限制来源,细节以官方文档为准。控制台里可以查看交换机、队列与消息,是验证本实战最直观的工具。

依赖:加 Rabbit Binder

收发双方工程都引入 Stream 核心与 Rabbit 的 Binder(当前 Stream 主线使用 spring-cloud-stream-binder-rabbit 系列坐标,若你所用的版本命名不同,以官方文档为准):

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

生产者:order-service 发布事件

定义 Supplier Bean 产出订单事件,并在配置里声明启用该函数:

@Bean
public Supplier<OrderCreatedEvent> orderSupplier() {
    return () -> {
        log.info("publish order created");
        return new OrderCreatedEvent(1L, "NO-0001");
    };
}
spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
  cloud:
    stream:
      bindings:
        orderSupplier-out-0:
          destination: order-events
          content-type: application/json
      function:
        definition: orderSupplier

启动后框架会周期性调用该 Supplier 并发送消息(默认轮询间隔以官方文档为准,这里用于演示)。实际业务里更多是在"下单成功"的时机动态发送,那需要用 StreamBridge:注入 StreamBridge 后调用 send 方法,把事件发到声明好的输出绑定(此处即 orderSupplier-out-0)上,业务触发、即时发送。

消费者:inventory-service 订阅处理

定义 Consumer Bean 处理订单事件,绑定到同一 destination:

@Bean
public Consumer<OrderCreatedEvent> onOrderCreated() {
    return event -> log.info("inventory handling order: {}", event.getOrderNo());
}
spring:
  cloud:
    stream:
      bindings:
        onOrderCreated-in-0:
          destination: order-events
          group: inventory-group      # 消费组:同组多实例分摊消息
      function:
        definition: onOrderCreated

启动后观察日志:生产者每发一条,消费者就处理一条。消息体默认走 JSON 序列化,两端 POJO 字段一致即可互通。

消费组验证:分摊与广播

再启动一个相同配置的 inventory-service 实例(换端口):两个实例同属 inventory-group,新消息会在两者间分摊,不会重复处理——这是水平扩容消费能力的基础。若把 group 改掉或去掉,每个实例各自收到全部消息(广播),按需选择。RabbitMQ 控制台里能看到为 order-events + inventory-group 自动创建的队列,消息积压量一目了然。

可靠性要点

消费处理抛异常时,Binder 默认会按配置重试,多次失败的消息需要去处:Rabbit Binder 支持自动绑定死信队列(auto-bind-dlq 类配置),把"毒消息"隔离出来人工排查,具体配置项与默认重试次数以官方文档为准。生产者侧,"发出去即成功"与"确认落盘"之间也有可靠性档位可选。消息场景的可靠性设计(确认、重试、死信、幂等消费)是工程落地重点,超出本章范围的部分请按官方文档逐项确认。

小结

Stream + RabbitMQ 的实战路径:起 RabbitMQ → 双方加 binder 依赖 → 生产者声明 Supplier 并配置 out 绑定 → 消费者声明 Consumer 并配置 in 绑定与 group → 启动验证分摊与广播。收尾三问:消息会不会丢(确认机制)、会不会重复(消费幂等)、失败去哪(死信队列)。从下一章开始,回到网关与调用链,把安全与上下文串起来。

笔记加载中…