Skip to content

Stream与Bus内部原理与生产治理

Spring Cloud Stream不是“换一套注解发送MQ”,Spring Cloud Bus也不是“配置中心本身”。Stream在业务函数与Kafka、RabbitMQ之间增加了Binding和Binder抽象;Bus则复用Stream传播Spring Cloud系统事件。只有把应用启动、绑定创建、消息收发、重试确认、配置广播和失败恢复连成一条链,才能判断一条消息到底丢在哪里、为什么重复、为什么扩容后仍积压,以及配置为什么只刷新了部分实例。

本章明确区分两条版本线:

版本线Java与Spring BootSpring CloudStream / Bus样本版本主要用途
存量线JDK 8、Spring Boot 2.7.182021.0.8Stream 3.2.9、Bus 3.1.2维护现有Java 8系统
现代线Java 17、Spring Boot 3.2.42023.0.1Stream 4.1.1、Bus 4.1.1新项目与本章源码主线

Spring Cloud必须按发布列车与Boot配套,不能只把spring-boot.version改成3.x。Boot 3以Java 17为最低基线,并迁移到Jakarta EE 9命名空间;Stream 4以函数式模型为主。后续3.x小版本可能升级客户端和属性,生产项目必须用mvn dependency:tree确认实际版本。

一、学习目标

学完后应能独立回答:

  1. Function、Binding、Binder、Broker分别解决哪一层问题。
  2. orderConsumer-in-0从哪里生成,应用启动时怎样变成Kafka消费者或Rabbit监听容器。
  3. 同组、异组、匿名组分别怎样投递,为什么Bus不能使用所有实例共享的普通消费组。
  4. 分区键、分区选择器、Broker分区和消费者并发是什么关系。
  5. Binding重试、框架消费重试、Broker重投和业务补偿为什么不能混为一谈。
  6. Kafka Binder与Rabbit Binder在确认、重试、死信和积压上的根本差异。
  7. StreamBridge怎样创建动态输出绑定,为什么无限动态目的地会耗尽资源。
  8. /actuator/busrefresh之后事件怎样发送、匹配、刷新和回传ACK。
  9. 为什么Bus广播不能代替配置中心,离线实例怎样最终收敛。
  10. 如何用指标、日志、Broker工具和Runbook定位消息不消费、重复、积压与部分刷新。

二、先建立四层心智模型

mermaid
flowchart TD
    A["业务函数:Supplier、Function、Consumer"] --> B["Binding:逻辑输入输出与配置"]
    B --> C["Binder:适配具体中间件"]
    C --> D["Broker:Kafka或RabbitMQ"]

四层职责不能互相替代:

负责什么不负责什么
Function业务输入、处理、输出不直接决定Kafka offset或Rabbit ACK
Binding函数端口到destination、group的映射不存储消息
Binder把统一配置翻译成客户端、监听容器和Broker资源不能抹平Broker语义差异
Broker持久化、路由、分区、副本、确认不知道数据库事务是否成功

“用了Stream就可以完全不懂Kafka/RabbitMQ”是危险结论。Stream能减少样板代码,却不能改变Kafka按分区保存日志、RabbitMQ按Exchange路由到Queue的事实。

三、核心源码对象职责

对象关键职责排查价值
BindableFunctionProxyFactory为函数创建输入、输出Binding名称及通道名称不对、函数没有绑定时先看这里
FunctionConfigurationFunctionCatalog查找函数并把通道连接到函数判断函数是否真正接入消息链
BindingService选择Binder,合并属性并调用绑定中间件不可用、Binder选择错误时的核心入口
BinderFactory / DefaultBinderFactory按名称和目标类型取得Binder,可维护Binder子上下文多Binder系统的选择与隔离
Binder定义bindConsumerbindProducer契约统一抽象与具体实现的分界线
AbstractMessageChannelBinder供应目的地、创建生产处理器或消费端点、错误基础设施生产与消费启动主链
ConsumerPropertiesgroup、并发、重试、分区等通用消费属性判断通用重试是否生效
ProducerProperties分区、requiredGroups、编码、错误通道等通用生产属性判断消息如何选分区
PartitionHandler提取分区键并计算目标分区同一业务键是否有序
StreamBridge在普通Service中发送,必要时动态创建输出Binding动态目的地与缓存治理
BusEnvironmentPostProcessor注入Bus函数定义、Binding映射和默认IDBus为什么无需手写Consumer
StreamBusBridge把远程应用事件发入Stream输出BindingBus发送链入口
BusConsumer接收、目标匹配、本地发布、ACK和traceBus部分实例不刷新时的关键对象

四、函数Binding名称怎样生成

函数式模型的默认规则是:

text
<functionName>-in-<index>
<functionName>-out-<index>

例如:

java
@Bean
public Function<OrderCreated, InventoryCommand> reserveInventory() {
    return event -> new InventoryCommand(event.orderId(), event.skuId(), event.quantity());
}

会生成:

text
reserveInventory-in-0
reserveInventory-out-0

源码中的BindableFunctionProxyFactory根据输入、输出数量建立代理Binding;FunctionConfiguration随后从FunctionCatalog找出reserveInventory并连接通道。多输入或多输出时索引递增,不要猜名称。

可以显式映射为稳定的业务名:

yaml
spring:
  cloud:
    stream:
      function:
        bindings:
          reserveInventory-in-0: orderCreatedInput
          reserveInventory-out-0: inventoryCommandOutput
      bindings:
        orderCreatedInput:
          destination: order.created.v1
          group: inventory-service
        inventoryCommandOutput:
          destination: inventory.reserve.v1

映射改变的是配置键,不会改变Java函数名。排查时应同时记录“函数隐式名、映射后Binding名、Broker destination”三者,否则很容易在配置里找错对象。

五、应用启动时怎样建立Binding

mermaid
flowchart TD
    A["解析函数定义"] --> B["创建输入输出通道"]
    B --> C["Input与Output生命周期启动"]
    C --> D["BindingService选择Binder"]
    D --> E["Binder供应Broker目的地"]
    E --> F["创建并启动消费端点或生产处理器"]
    F --> G["发布BindingCreatedEvent"]

5.1 消费端启动主链

BindingService.bindConsumer(input, inputName)大致完成:

  1. 读取Binding的destination、group和通用ConsumerProperties
  2. 按显式binder、默认Binder和目标通道类型选择具体Binder。
  3. 如果Binder支持扩展属性,再合并Kafka或Rabbit专属配置。
  4. 调用binder.bindConsumer(destination, group, input, properties)
  5. 具体Binder创建Topic/Queue等资源、监听容器和消息适配端点。
  6. 消费端点输出被连接到Stream输入通道,收到的Broker消息才能进入业务函数。
mermaid
flowchart TD
    A["Broker消息到达客户端"] --> B["Binder消费端点"]
    B --> C["转换为Spring Message"]
    C --> D["输入Binding通道"]
    D --> E["Function包装与调用"]
    E --> F["成功确认或失败处理"]

如果Broker里有消息但函数完全没日志,应先判断监听端点是否创建成功、Binding是否使用正确destination/group,而不是立刻怀疑业务代码。

5.2 生产端启动主链

BindingService.bindProducer(output, outputName)会选择Binder、合并生产属性并调用bindProducerAbstractMessageChannelBinder的典型过程是:

  1. provisionProducerDestination供应或校验目的地。
  2. createProducerMessageHandler创建具体发送处理器。
  3. 初始化并启动处理器。
  4. SendingHandler订阅到输出通道。
  5. 发布BindingCreatedEvent

之后业务函数把消息发入输出通道,真正发送由处理器完成。业务方法“返回了对象”不等于Broker已经持久化,更不等于下游业务已经成功。

5.3 Binding创建重试不等于消息重试

如果应用启动时Broker暂不可达,BindingService可以按bindingRetryInterval延迟重试创建Binding,并用LateBinding占位。这一重试解决的是“连接和监听容器还没有建起来”,不是“某条消息处理失败”。

重试层发生时间重试对象典型问题
Binding创建重试应用启动或重绑建立生产者/消费者BindingBroker暂不可达
Stream消费重试收到某条消息后再调用业务函数临时数据库异常
Broker重投消费失败或未确认后同一Broker消息Kafka offset未提交、Rabbit未ACK
业务补偿进入重试Topic、任务表后一次业务操作长时间外部依赖故障

把四层同时打开可能产生乘法效应。例如框架3次、Broker再投5次、补偿任务又3次,最坏可能执行45次,因此幂等不是可选项。

六、Binder怎样选择,多Binder为什么容易出错

单Binder时依赖通常足够决定实现;多Binder时必须明确配置。选择过程可理解为:

mermaid
flowchart TD
    A["读取Binding配置"] --> B["优先使用显式binder"]
    B --> C["否则使用defaultBinder"]
    C --> D["仍无则按目标类型推断"]
    D --> E["校验唯一候选并取得Binder"]

现代线多Binder示意:

yaml
spring:
  cloud:
    stream:
      default-binder: kafka
      bindings:
        auditConsumer-in-0:
          binder: rabbit
          destination: audit.event.v1
          group: audit-service
        orderConsumer-in-0:
          binder: kafka
          destination: order.created.v1
          group: order-analysis

真正隔离不同集群时还要配置Binder environments。DefaultBinderFactory可为Binder维护独立应用上下文,这意味着连接工厂、客户端配置和Bean可隔离;但子上下文也增加启动时间、内存和配置复杂度。

常见错误:

  • 同时引入Kafka和Rabbit Binder却没有明确默认值,启动时出现候选冲突。
  • destination相同但两个Binder指向不同Broker,排查人员误以为消息“神秘消失”。
  • 把Kafka专属属性写到通用Binding节点,属性没有生效。
  • 用统一抽象掩盖事务、顺序、死信等强依赖Broker的设计。

七、Group、广播和匿名订阅的真实语义

Binder契约规定:同一个非空group共享订阅;空group要按匿名、互不共享的订阅处理。

mermaid
flowchart TD
    A["同一业务事件"] --> B["库存组与积分组各得一份"]
    B --> C["每个组内由一个实例处理"]

因此:

  • 同一服务的多个实例通常使用相同group,实例间竞争,整组处理一份逻辑消息。
  • 不同业务服务使用不同group,每个group各得到一份逻辑消息。
  • 不写group会创建匿名非共享订阅,每个实例都可能收到一份;具体是否持久、离线是否保留取决于Binder和Broker配置,不能笼统承诺。

错误示例:库存服务每次发布都生成随机group。结果是旧group不断遗留、每个实例都消费全量消息,既重复扣库存又制造Broker资源。

Bus为了让每个在线实例都收到广播,正是利用非共享消费语义;这也带来“离线实例可能错过事件”的边界,后文会详细解释。

八、分区原理:同一订单为什么进入同一条处理通道

Stream通用分区链是:

mermaid
flowchart TD
    A["消息"] --> B["提取partition key"]
    B --> C["选择器计算原始值"]
    C --> D["对partitionCount取模并归一化"]
    D --> E["写入目标分区信息"]
    E --> F["Binder映射到Broker分区或队列"]

PartitionHandler先使用自定义PartitionKeyExtractorStrategypartitionKeyExpression取得键,再使用自定义选择器、表达式或默认哈希,最后归一化为0..partitionCount-1。默认实现还处理了Integer.MIN_VALUE绝对值异常边界。

yaml
spring:
  cloud:
    stream:
      bindings:
        publishOrder-out-0:
          destination: order.created.v1
          producer:
            partition-key-expression: headers['orderId']
            partition-count: 12

必须理解四个限制:

  1. 只保证同一键稳定映射,不保证全局顺序。
  2. Kafka实际Topic分区数不能小于设计值;修改分区数后哈希映射会变化。
  3. 同一分区通常同时只能被组内一个消费者实例占有,消费者数超过分区数会有空闲实例。
  4. 热门键会形成热分区,整体消费者数量再多也无法并行处理该键。

若业务要求同一订单严格有序,应以orderId为键,并在消费端把数据库状态机作为最终约束。仅依靠队列顺序不能防止重试、人工补偿和跨Topic造成的乱序。

九、消息转换、contentType与原生编码

函数收到Java对象之前通常经历:Broker字节、Spring Message、内容类型转换、函数参数适配。contentType描述载荷如何序列化/反序列化,默认转换和Binder原生编码是两条不同路径。

mermaid
flowchart TD
    A["Java事件对象"] --> B["Message转换"]
    B --> C["字节或原生客户端对象"]
    C --> D["Broker"]
    D --> E["消费端反序列化"]
    E --> F["Consumer参数"]

建议事件契约至少包含:eventIdeventTypeeventVersionoccurredAttraceId和业务主键。不要直接把数据库实体序列化为跨服务契约,因为字段删除、懒加载代理和内部状态会泄漏到消费者。

事件演进原则:

  • 新增可选字段通常向后兼容。
  • 删除、改名、改变类型要升版本或做双读。
  • 消费者要容忍未知字段。
  • 时间、金额、枚举必须明确格式与单位。
  • PII、token和密钥不得写入消息体或日志。

启用useNativeEncoding后序列化更多交给具体Broker客户端。它能使用Kafka Serializer等原生能力,但也使迁移Binder更困难。切换前必须让生产者和消费者对序列化协议达成一致。

十、StreamBridge为什么适合业务Service发送消息

业务消息通常由HTTP请求或数据库事务触发,不适合轮询式SupplierStreamBridge允许普通Service主动发送:

java
@Service
public class OrderEventPublisher {
    private final StreamBridge streamBridge;

    public OrderEventPublisher(StreamBridge streamBridge) {
        this.streamBridge = streamBridge;
    }

    public boolean publish(OrderCreated event) {
        return streamBridge.send("orderCreatedOutput", event);
    }
}

内部主链是:

mermaid
flowchart TD
    A["send绑定名与数据"] --> B["查找或创建MessageChannel"]
    B --> C["解析Binder与contentType"]
    C --> D["转换消息并处理分区"]
    D --> E["必要时动态创建生产Binding"]
    E --> F["发送到输出通道"]

10.1 send()返回true代表什么

它只代表消息被当前Spring消息通道接受,具体发送确认还受通道类型、Binder和客户端配置影响;它绝不代表消费者已处理成功,更不代表消费者数据库事务已提交。

10.2 动态目的地的风险

如果Binding不存在,StreamBridge可以动态创建并缓存生产Binding。动态租户Topic看似方便,但若租户ID、用户输入或随机值直接作为destination,会导致:

  • 应用内Binding缓存膨胀;
  • Kafka Topic或Rabbit Exchange数量失控;
  • 连接、线程、指标标签和权限规则爆炸;
  • 缓存淘汰时频繁解绑和重建,发送抖动。

生产规则应是固定目的地优先;确需动态目的地时使用白名单、限制基数、统一命名、设置生命周期和Broker配额。异步发送还要验证trace上下文是否传播。

十一、消费重试、错误通道与死信边界

Stream 4.1.1样本的通用消费默认值包括:concurrency=1maxAttempts=3、初始退避1000ms、最大10000ms、倍数2.0。maxAttempts包含第一次调用,因此设为1才是关闭框架重试。

yaml
spring:
  cloud:
    stream:
      bindings:
        orderConsumer-in-0:
          destination: order.created.v1
          group: inventory-service
          consumer:
            max-attempts: 3
            back-off-initial-interval: 1000
            back-off-max-interval: 10000
            back-off-multiplier: 2.0

一次失败的正确决策:

mermaid
flowchart TD
    A["业务处理抛异常"] --> B["判断是否短暂可恢复"]
    B --> C["可恢复则有限指数退避"]
    C --> D["仍失败则进入DLQ或重试Topic"]
    D --> E["告警、修复、受控回放"]

错误通道承载的是消息处理异常信息,不等于持久化死信队列。若错误处理器本身失败、实例重启或错误通道没有持久消费者,异常记录仍可能丢失。需要长期保留和回放时,应使用Broker DLQ、重试Topic或数据库补偿表。

以下异常通常不应盲目原地重试:参数校验失败、事件版本不支持、业务唯一约束明确冲突、永久权限错误。原地重试会占用消费线程并阻塞同分区后续消息。

十二、至少一次投递下消费幂等怎样落地

可靠消息系统通常追求at-least-once,这意味着“可能重复,但不能静默丢失”。消费者必须让同一eventId重复到达时得到相同业务结果。

数据库去重与业务操作应在同一个本地事务内:

sql
CREATE TABLE consumed_event (
    consumer_group VARCHAR(100) NOT NULL,
    event_id VARCHAR(64) NOT NULL,
    consumed_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (consumer_group, event_id)
);
java
@Transactional
public void consume(OrderCreated event) {
    int inserted = consumedEventRepository.insertIgnore(
            "inventory-service", event.eventId());
    if (inserted == 0) {
        return;
    }
    inventoryRepository.reserve(event.skuId(), event.quantity());
}

为什么不能先在Redis写“已消费”再更新数据库:如果Redis成功而数据库失败,重投消息会被误判为已处理,形成真实数据丢失。Redis可做快速前置过滤,但最终幂等状态要与核心业务结果具备同一事务边界或可验证补偿链。

更多实现见请求、消息与任务幂等

12.1 消费确认、Offset/ACK与业务事务顺序

Stream统一的是编程模型,不是把Kafka和RabbitMQ的确认语义抹平。判断“消息是否消费成功”时,不能只看Consumer方法有没有执行完,还要同时看三件事:

事情谁负责成功含义
业务函数执行Spring Cloud Stream和业务代码方法没有抛异常,或异常被框架处理
本地业务事务Spring事务和数据库订单、库存、消费日志等数据真正提交
Broker确认Kafka offset提交或Rabbit ACKBroker认为这条消息或这个位置已经处理

这三件事没有天然原子性。数据库事务提交不了Kafka offset,Rabbit ACK也回滚不了MySQL事务,所以必须人为设计顺序和幂等。

mermaid
flowchart TD
    A["Broker投递消息"] --> B["Stream调用业务函数"]
    B --> C["开启本地事务"]
    C --> D["写消费日志和业务数据"]
    D --> E["提交数据库事务"]
    E --> F["提交Kafka offset或发送Rabbit ACK"]
    F --> G["Broker不再正常重投"]

最危险的顺序是“先确认Broker,再执行业务”:

mermaid
flowchart TD
    A["收到消息"] --> B["先提交offset或ACK"]
    B --> C["执行业务数据库更新"]
    C --> D["业务失败或进程崩溃"]
    D --> E["Broker认为已消费"]
    E --> F["业务数据永久缺失"]

这种顺序会把“至少一次投递”变成“最多一次投递”:消息可能不重复,但业务可能丢。订单、扣库存、支付通知、资产采集入库这类场景不能接受。

更常见、也更合理的顺序是“业务事务先提交,再确认Broker”:

窗口结果为什么必须幂等
业务事务提交成功,offset/ACK也成功正常完成无重复
业务事务提交成功,offset/ACK失败或进程崩溃Broker会再次投递第二次必须识别已处理并直接确认
业务事务回滚,offset/ACK未提交Broker会再次投递下次可重新处理
业务事务回滚,但offset/ACK已提交消息丢失这是要避免的错误顺序

Kafka和RabbitMQ的差异要说清楚:

对比项KafkaRabbitMQ
确认对象consumer group在partition上的offset位置Queue中某条投递消息的delivery tag
确认后果下次从更后面的offset继续pollBroker可从队列移除这条已ACK消息
未确认后果offset未前进,重启或重平衡后可能重读连接断开或NACK后可重新入队或进死信
排查重点committed offset、log end offset、Lag、rebalanceready、unacked、prefetch、redelivered、DLQ

所以Stream消费端的生产原则是:

  1. 业务数据和消费日志放在同一个本地数据库事务里。
  2. 消费日志使用consumer_group + event_id唯一键。
  3. 事务提交成功后再允许Binder提交Kafka offset或Rabbit ACK。
  4. 如果ACK失败导致重投,第二次根据消费日志直接返回成功,让Broker完成确认。
  5. 如果业务失败,不要吞异常伪装成功,否则框架可能提交确认。
  6. 对不可恢复错误要进入DLQ或补偿表,不能无限占住主分区或主队列。

一个更贴近商业项目的消费日志表应至少包含状态,而不是只有一条“已消费”:

sql
CREATE TABLE consumer_event_log (
    consumer_group VARCHAR(100) NOT NULL,
    event_id VARCHAR(64) NOT NULL,
    status VARCHAR(32) NOT NULL,
    retry_count INT NOT NULL DEFAULT 0,
    last_error VARCHAR(500),
    created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
    updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (consumer_group, event_id)
);

状态含义:

状态含义处理方式
PROCESSING已抢到处理权,业务正在执行超时后由补偿任务判断是否可接管
SUCCESS业务已经成功提交重复消息直接返回成功并ACK
FAILED_RETRYABLE临时失败,例如DB超时、下游短暂不可用允许框架重试、重试Topic或补偿任务
FAILED_FINAL永久失败,例如事件版本不支持、参数非法进入DLQ或人工处理,不要原地热重试

JDK 8最小Demo如下,用ConcurrentHashMap模拟消费日志。第一次业务成功但ACK失败,Broker重投后,第二次命中SUCCESS,直接确认,不会重复扣库存。

java
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

public class ConsumerAckOrderDemo {

    enum Status {
        PROCESSING, SUCCESS
    }

    static class OrderCreated {
        private final String eventId;
        private final String orderId;
        private final String skuId;
        private final int quantity;

        OrderCreated(String eventId, String orderId, String skuId, int quantity) {
            this.eventId = eventId;
            this.orderId = orderId;
            this.skuId = skuId;
            this.quantity = quantity;
        }
    }

    static class InventoryConsumer {
        private final Map<String, Status> consumerLog = new ConcurrentHashMap<String, Status>();
        private int stock = 10;

        void onMessage(OrderCreated event, boolean ackWillFail) {
            String key = "inventory-service:" + event.eventId;
            Status old = consumerLog.putIfAbsent(key, Status.PROCESSING);
            if (old == Status.SUCCESS) {
                System.out.println("重复消息,业务已成功,直接ACK: " + key);
                ack(ackWillFail);
                return;
            }

            stock = stock - event.quantity;
            consumerLog.put(key, Status.SUCCESS);
            System.out.println("扣库存成功,orderId=" + event.orderId + ", stock=" + stock);

            ack(ackWillFail);
        }

        private void ack(boolean fail) {
            if (fail) {
                throw new RuntimeException("模拟ACK或offset提交失败,Broker稍后重投");
            }
            System.out.println("Broker确认成功");
        }
    }

    public static void main(String[] args) {
        InventoryConsumer consumer = new InventoryConsumer();
        OrderCreated event = new OrderCreated("evt-1001", "order-1001", "sku-7", 2);

        try {
            consumer.onMessage(event, true);
        } catch (Exception ex) {
            System.out.println(ex.getMessage());
        }

        consumer.onMessage(event, false);
    }
}

这个Demo刻意把ACK做成方法调用,实际项目里ACK/offset由Kafka Binder或Rabbit Binder监听容器根据配置完成。真正落地时要把consumerLog.putIfAbsent、业务更新、状态改为SUCCESS放入同一个数据库事务;内存Map只能说明流程,不能用于生产去重。

排查时按现象倒推:

现象优先怀疑证据
Broker认为消费了,但业务表没数据先ACK后业务、异常被吞、事务未提交offset/ACK日志、业务事务日志、异常处理器
业务表有数据,但消息又来了业务提交后ACK失败、rebalance、消费者宕机eventId、消费日志、重平衡日志、Broker重投标记
重复消费突然增多下游变慢导致处理超时、消费者频繁重平衡、ACK失败P99、poll间隔、心跳、连接池、rebalance次数
某个分区或队列卡住毒消息、无限原地重试、顺序键热点Lag分布、重试次数、DLQ、业务错误分类

面试里可以这样回答:

text
Stream消费成功不能只理解为Consumer方法执行完成。生产上应让业务数据和消费日志在同一个
本地事务中提交,提交成功后再让Binder提交Kafka offset或Rabbit ACK。业务成功后ACK失败会
导致重复投递,所以消费者必须用consumer group加eventId做幂等;如果先ACK再写库,写库失败
会造成消息丢失。Stream统一编程模型,但Kafka offset和Rabbit ACK的底层语义仍要分别分析。

十三、Kafka Binder到底映射了什么

Kafka Binder把destination映射到Topic,把group映射到Kafka consumer group,把Stream分区信息映射到Kafka partition,并通过监听容器管理poll、offset和并发。

mermaid
flowchart TD
    A["Stream输出Binding"] --> B["Kafka生产处理器"]
    B --> C["Topic与Partition"]
    C --> D["Consumer Group"]
    D --> E["Kafka监听容器"]
    E --> F["Stream输入Binding"]

关键配置边界:

  • ackMode决定监听容器何时提交offset;旧autoCommitOffset属性已逐步废弃,不要用旧文章直接套现代线。
  • enableDlqdlqName可启用Kafka Binder死信发布,但要验证失败发布、offset提交和事务配置的组合。
  • startOffset决定新group从哪里开始;resetOffsets会主动重置,误用可能大量重放或跳过数据。
  • 消费并发的有效上限受分区数限制。
  • Kafka顺序只在单分区内成立。

Kafka消费成功的业务边界通常是“数据库事务提交后再允许offset确认”。如果先提交offset再写库,宕机会丢业务;如果先写库后offset提交失败,会重复,所以必须幂等。

完整底层原理见Kafka专栏

十四、Rabbit Binder到底映射了什么

Rabbit Binder需要把destination、group映射为Exchange、Queue和Binding关系,再通过监听容器消费并ACK/NACK。

mermaid
flowchart TD
    A["Stream输出Binding"] --> B["Rabbit Exchange"]
    B --> C["Routing与Queue"]
    C --> D["Rabbit监听容器"]
    D --> E["Stream输入Binding"]
    E --> F["ACK、拒绝或死信"]

现代样本中的重要属性:

属性作用风险
autoBindDlq让Binder自动供应DLX/DLQ拓扑生产权限不足时启动失败
republishToDlq失败后由Binder重新发布,并附加异常头消息变大、DLQ发布也可能失败
requeueRejected拒绝后是否重新入原队列true可能形成无退避热循环

Rabbit的Queue积压、unacked数量、prefetch和消费者并发共同决定吞吐。只看Queue ready不看unacked,会把“消费者拿走但处理很慢”误判为“没有积压”。

完整ACK、重试和DLX原理见RabbitMQ专栏

十五、Kafka Binder与Rabbit Binder怎么选

维度KafkaRabbitMQ
核心模型分区追加日志,消费者维护offsetExchange路由到Queue,Broker跟踪投递与确认
广播方式不同consumer group各读一份不同Queue绑定同一Exchange
顺序分区内有序单Queue单消费者较直观,多消费者会并发
重放保留期内调整offset重放消息ACK后通常从队列删除,需另建归档/重试机制
高吞吐日志流更擅长可用但模型不同
灵活路由依赖Topic、key和应用设计Exchange、routing key能力强
积压观察consumer group lagready、unacked、publish/deliver/ack速率
扩容上限受Topic分区数约束受Queue、prefetch、消费者和下游约束

选择不能只看“哪个吞吐高”。订单事件需要长时间重放、数据管道和大吞吐时常偏Kafka;复杂路由、工作队列、低延迟命令分发常偏Rabbit。若项目强依赖某一方高级语义,直接使用Spring Kafka或Spring AMQP有时比强行抽象更清晰。

十六、消息堆积为什么扩容后只短暂有效

积压变化可用一个简单关系理解:

text
积压增长速度 = 生产速率 - 实际确认速率

扩容后短暂下降、随后继续增长,常见原因不是“扩容没生效”,而是实际确认速率再次低于生产速率:

  1. Kafka消费者数已超过分区数,多出的实例空闲。
  2. Lag集中在少数热分区,新增消费者无法拆分单分区。
  3. 数据库、Redis、HTTP下游到达瓶颈,消费者越多争用越严重。
  4. 连接池小于消费并发,大量线程等待连接。
  5. 失败消息反复重试,表面消费很快但有效成功率低。
  6. 单条消息体积或批次变大,平均处理时间上升。
  7. GC停顿、CPU限流、容器内存压力使实例实际吞吐下降。
  8. 扩容期间重平衡暂停消费,频繁自动伸缩造成抖动。
  9. Rabbit prefetch过大,消息堆在少数实例的unacked区。
  10. 生产侧流量也同步上涨,扩容增加量小于新增流量。

“Lag分布”不是只看总Lag,而是查看每个Topic、group、partition的log-end-offset - committed-offset及其变化率。总Lag 100万可能均匀分在100个分区,也可能99万集中在一个热分区,两个场景的扩容策略完全不同。

16.1 积压Runbook

第一步,冻结证据:记录时间窗、发布速率、确认速率、总Lag、分区Lag分布、消费者实例数、分区数、重平衡次数、失败率和下游延迟。

第二步,先算理论吞吐:

text
单实例理论吞吐约等于 并发数 / 平均处理秒数
有效吞吐还要乘以成功率,并受分区数和下游容量限制

第三步,区分瓶颈:

证据更可能的根因合理动作
少数分区Lag很高热key或分区倾斜调整key、拆Topic、业务键分片
所有分区均匀增长总容量不足在下游允许范围内扩消费者
DB连接池等待高数据库/连接池瓶颈优化SQL、批量写、限并发,不能盲目扩容
重试次数激增毒消息或依赖故障暂停热循环、隔离DLQ、修复依赖
CPU低但消费慢I/O等待或锁竞争查线程栈、连接池、锁与远程调用
扩容伴随频繁rebalance伸缩和消费参数不稳稳定实例、调整poll与处理边界

第四步,恢复后受控回放。不要一次把消费并发拉满冲垮数据库;先限速消化,观察业务成功率、DB负载和Lag斜率,再逐步提高。

更完整专题见消息堆积、扩容与背压

十七、命令式函数与响应式函数不能用同一套直觉

命令式Consumer<T>通常由可订阅通道触发一次消息调用;响应式Function<Flux<T>, Flux<R>>连接的是Publisher链。响应式函数不是“每条消息反射调用一次普通方法”,订阅、背压、线程切换和错误传播边界不同。

java
@Bean
public Function<Flux<OrderCreated>, Flux<OrderRiskResult>> riskCheck() {
    return input -> input.flatMap(this::checkRisk, 8);
}

flatMap(..., 8)允许最多8路并发,可能改变输出顺序;若业务要求同一订单有序,应按键分组或使用顺序算子,并评估内存。响应式链里阻塞JDBC调用仍会阻塞线程,不能因为返回Flux就自动获得非阻塞能力。

十八、Java 17与Spring Boot 3可运行订单事件Demo

本Demo使用Java 17、Boot 3.2.4、Cloud 2023.0.1和Kafka Binder。事件使用Record;这段代码不能复制到JDK 8。

18.1 Maven配置

xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>

<parent>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-parent</artifactId>
    <version>3.2.4</version>
    <relativePath/>
</parent>

<groupId>com.example</groupId>
<artifactId>stream-order-demo</artifactId>
<version>1.0.0</version>

<properties>
    <java.version>17</java.version>
    <spring-cloud.version>2023.0.1</spring-cloud.version>
</properties>

<dependencyManagement>
    <dependencies>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-dependencies</artifactId>
            <version>${spring-cloud.version}</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>
    </dependencies>
</dependencyManagement>

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.cloud</groupId>
        <artifactId>spring-cloud-starter-stream-kafka</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-actuator</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-jdbc</artifactId>
    </dependency>
    <dependency>
        <groupId>com.h2database</groupId>
        <artifactId>h2</artifactId>
        <scope>runtime</scope>
    </dependency>
</dependencies>

<build>
    <plugins>
        <plugin>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-maven-plugin</artifactId>
        </plugin>
    </plugins>
</build>
</project>

18.2 事件与发送接口

启动类:

java
@SpringBootApplication
public class StreamOrderApplication {
    public static void main(String[] args) {
        SpringApplication.run(StreamOrderApplication.class, args);
    }
}
java
public record OrderCreated(
        String eventId,
        long orderId,
        long skuId,
        int quantity,
        Instant occurredAt,
        int eventVersion) {
}
java
@RestController
@RequestMapping("/orders")
public class OrderController {
    private final StreamBridge streamBridge;

    public OrderController(StreamBridge streamBridge) {
        this.streamBridge = streamBridge;
    }

    @PostMapping("/{orderId}/events")
    public ResponseEntity<Void> publish(@PathVariable long orderId) {
        OrderCreated event = new OrderCreated(
                UUID.randomUUID().toString(), orderId, 10086L,
                2, Instant.now(), 1);
        boolean accepted = streamBridge.send("orderCreatedOutput", event);
        return accepted
                ? ResponseEntity.accepted().build()
                : ResponseEntity.internalServerError().build();
    }
}

生产项目不能把HTTP成功、数据库提交与消息发送当成一个原子动作。若订单先提交而进程在发送前崩溃,事件丢失;应使用本地消息表/Outbox或Broker事务消息建立一致性闭环。

18.3 幂等消费者

java
@Configuration
public class OrderConsumerConfiguration {
    @Bean
    Consumer<OrderCreated> reserveInventory(InventoryApplicationService service) {
        return service::reserveOnce;
    }
}
java
@Service
public class InventoryApplicationService {
    private final ConsumedEventRepository consumedEvents;
    private final InventoryRepository inventory;

    public InventoryApplicationService(ConsumedEventRepository consumedEvents,
                                       InventoryRepository inventory) {
        this.consumedEvents = consumedEvents;
        this.inventory = inventory;
    }

    @Transactional
    public void reserveOnce(OrderCreated event) {
        if (!consumedEvents.tryInsert("inventory-service", event.eventId())) {
            return;
        }
        int changed = inventory.reserveIfEnough(
                event.skuId(), event.quantity());
        if (changed != 1) {
            throw new IllegalStateException("库存不足或SKU不存在");
        }
    }
}

Repository契约与最小JDBC实现:

java
public interface ConsumedEventRepository {
    boolean tryInsert(String consumerGroup, String eventId);
}

public interface InventoryRepository {
    int reserveIfEnough(long skuId, int quantity);
}
java
@Repository
public class JdbcConsumedEventRepository implements ConsumedEventRepository {
    private final JdbcTemplate jdbcTemplate;

    public JdbcConsumedEventRepository(JdbcTemplate jdbcTemplate) {
        this.jdbcTemplate = jdbcTemplate;
    }

    @Override
    public boolean tryInsert(String consumerGroup, String eventId) {
        try {
            return jdbcTemplate.update(
                    "INSERT INTO consumed_event(consumer_group, event_id) VALUES (?, ?)",
                    consumerGroup, eventId) == 1;
        }
        catch (DuplicateKeyException duplicate) {
            return false;
        }
    }
}
java
@Repository
public class JdbcInventoryRepository implements InventoryRepository {
    private final JdbcTemplate jdbcTemplate;

    public JdbcInventoryRepository(JdbcTemplate jdbcTemplate) {
        this.jdbcTemplate = jdbcTemplate;
    }

    @Override
    public int reserveIfEnough(long skuId, int quantity) {
        return jdbcTemplate.update("""
                UPDATE inventory
                   SET available = available - ?
                 WHERE sku_id = ?
                   AND available >= ?
                """, quantity, skuId, quantity);
    }
}

这里的文本块是Java 15+语法,现代线Java 17可以使用;JDK 8必须改成普通字符串拼接。示例H2捕获唯一键异常后事务仍可继续,生产数据库要按方言实现原子“插入并判断”:MySQL可用INSERT IGNORE,PostgreSQL应使用INSERT ... ON CONFLICT DO NOTHING并根据影响行数判断,不能不加验证地照搬异常捕获。

src/main/resources/schema.sql

sql
CREATE TABLE consumed_event (
    consumer_group VARCHAR(100) NOT NULL,
    event_id VARCHAR(64) NOT NULL,
    consumed_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (consumer_group, event_id)
);

CREATE TABLE inventory (
    sku_id BIGINT PRIMARY KEY,
    available INT NOT NULL
);

INSERT INTO inventory(sku_id, available) VALUES (10086, 100);

18.4 配置

yaml
spring:
  datasource:
    url: jdbc:h2:mem:stream_demo;DB_CLOSE_DELAY=-1
    username: sa
    password: ""
  cloud:
    function:
      definition: reserveInventory
    stream:
      bindings:
        orderCreatedOutput:
          destination: order.created.v1
          content-type: application/json
          producer:
            partition-key-expression: payload.orderId
            partition-count: 12
        reserveInventory-in-0:
          destination: order.created.v1
          group: inventory-service
          consumer:
            max-attempts: 3
            back-off-initial-interval: 1000
      kafka:
        binder:
          brokers: localhost:9092
        bindings:
          reserveInventory-in-0:
            consumer:
              enable-dlq: true
              dlq-name: order.created.inventory.dlq
management:
  endpoints:
    web:
      exposure:
        include: health,info,metrics,bindings

启动本地Kafka后运行应用,再调用:

bash
curl -i -X POST http://localhost:8080/orders/9001/events

Demo会向order.created.v1发送事件,再由同一应用的reserveInventory函数消费并扣减H2库存。若要直接验证重复消息,可临时让Controller接收固定eventId并发送两次;最终库存只能扣一次。生产测试还必须覆盖:消费函数抛异常后的实际尝试次数、DLQ消息内容、Broker重启、应用在业务提交后确认前崩溃、以及幂等表与业务表是否始终同事务回滚。

十九、JDK 8与Spring Boot 2.7存量线差异

存量线可以继续使用函数式模型,但代码和版本不能照搬现代线:

项目JDK 8 + Boot 2.7Java 17 + Boot 3.x
Java最低版本817
事件DTO普通POJO可用Record,也可继续POJO
EE包名javax.*jakarta.*
Cloud发布列车样本2021.0.82023.0.1
Stream样本3.2.94.1.1
旧注解模型遗留项目可能仍有@EnableBinding应迁移到函数式模型
可观测性Sleuth时代常见Micrometer Tracing方向

JDK 8事件必须写成POJO:

java
public class OrderCreatedEvent {
    private String eventId;
    private Long orderId;

    public OrderCreatedEvent() {
    }

    public OrderCreatedEvent(String eventId, Long orderId) {
        this.eventId = eventId;
        this.orderId = orderId;
    }

    public String getEventId() {
        return eventId;
    }

    public void setEventId(String eventId) {
        this.eventId = eventId;
    }

    public Long getOrderId() {
        return orderId;
    }

    public void setOrderId(Long orderId) {
        this.orderId = orderId;
    }
}

不要为了迁移而同时升级JDK、Boot、Cloud、Broker客户端、序列化协议和所有业务代码。更稳妥的顺序是:先建立契约与回归测试,迁移旧注解模型到函数式模型,再升级Java和依赖平台,最后逐项启用现代能力。跨版本滚动发布期间,新旧实例必须能互相读取同一事件格式。

二十、Bus的职责边界

Spring Cloud Bus把多个应用实例连接到同一系统事件通道,典型用途是广播刷新或环境变更。它不是:

  • 配置内容的权威存储;
  • 所有实例必达的部署系统;
  • 数据库事务协调器;
  • 业务消息总线的替代品;
  • 敏感配置明文分发工具。

权威配置仍应位于Config Server、Nacos、Git、数据库或密钥系统。Bus只发送“发生了变化”或受控变更事件,实例重启后必须能从权威源重新加载并收敛。

二十一、Bus启动时自动做了什么

BusEnvironmentPostProcessor在环境准备阶段追加并映射函数:

text
spring.cloud.function.definition += busConsumer
busConsumer-in-0 -> springCloudBusInput
springCloudBusInput.destination -> springCloudBus
springCloudBusOutput.destination -> springCloudBus

默认常量包括:

名称默认值
输入BindingspringCloudBusInput
输出BindingspringCloudBusOutput
Broker destinationspringCloudBus
消费函数busConsumer

它还会推导spring.cloud.bus.id。ID用于目标匹配和判断事件是否来自自己,生产环境必须足够唯一;若多个实例ID碰撞,可能错误忽略事件或把单实例事件投递给错误对象。

二十二、/actuator/busrefresh发送链

mermaid
flowchart TD
    A["调用busrefresh端点"] --> B["发布RefreshRemoteApplicationEvent"]
    B --> C["源实例按目标执行本地刷新"]
    C --> D["远程事件监听器继续传播"]
    D --> E["StreamBusBridge发送输出Binding"]
    E --> F["Binder发布到springCloudBus"]

关键细节:

  1. 端点先在源实例本地发布远程事件。
  2. 如果目标匹配源实例,本地RefreshListener可以立即调用ContextRefresher.refresh()
  3. RemoteApplicationEventListener忽略ACK事件,并避免把不该发送的事件再次传播。
  4. StreamBusBridgeStreamBridgespringCloudBusOutput发送。
  5. Binder再把事件编码并发布到Kafka或RabbitMQ。

因此HTTP端点返回成功通常只能说明事件已在源实例触发并尝试发送,不能证明所有目标实例都刷新成功。

二十三、目标实例接收与刷新链

mermaid
flowchart TD
    A["Bus Broker事件"] --> B["busConsumer接收"]
    B --> C["不匹配忽略,匹配则本地发布"]
    C --> D["刷新监听器处理"]
    D --> E["刷新环境与相关Bean"]
    E --> F["可选发送ACK与trace"]

BusConsumer.accept()还会判断事件是否来自自己。源实例已经本地执行过刷新,同一个事件从Broker绕回时不会再次本地发布,从而避免源实例重复刷新。

刷新不是“重建整个ApplicationContext”。ContextRefresher会刷新环境并处理可刷新的作用域;没有使用正确刷新机制、在初始化时复制到普通单例字段、或第三方组件自己缓存的值,不一定自动变化。每个动态配置项都应通过测试证明实际生效路径。

二十四、Bus目标匹配怎样工作

Bus使用路径式目标匹配。现代样本的PathDestinationFactory会把空目标视为**,简单服务名扩展为service:**,并以冒号作为Ant风格路径分隔符。

可以理解为三种范围:

范围用途风险
**所有Bus实例影响面最大,应严格授权
order-service:**某服务所有实例适合服务级配置刷新
具体Bus ID单个实例依赖ID唯一且运维可定位

实际端点路径和destination编码要按当前Spring Cloud Bus版本验证,尤其注意冒号、端口和实例ID在URL中的转义。不要把用户输入直接拼成广播目标。

二十五、为什么Bus能广播,又为什么离线实例会错过

Bus输入Binding默认没有让所有实例共享一个普通业务group。根据Binder契约,空group是匿名、非共享订阅,因此每个在线Bus实例都能接收事件。

mermaid
flowchart TD
    A["springCloudBus事件"] --> B["每个在线实例有独立匿名订阅"]
    B --> C["实例A、B、C各自收到一份"]

代价是匿名订阅的持久性取决于Binder/Broker实现。实例在发布时离线、网络隔离或订阅尚未建立,可能错过这次广播。正确恢复机制是:

  1. 配置中心保存权威版本。
  2. 实例启动时主动拉取最新配置。
  3. Bus用于加速在线实例感知变化。
  4. 监控每个实例实际配置版本,而不是只看端点HTTP状态。
  5. 发现版本不一致时重试定向刷新或安全重启实例。

这就是“事件通知”和“状态收敛”的区别。事件可能丢,权威状态必须可重读。

二十六、ACK与trace能证明什么

现代样本中Bus ACK默认启用,trace默认关闭。目标实例可发送AckRemoteApplicationEvent,trace可记录发送与接收事件。

ACK能证明某个Bus消费者看到了并处理了事件链的一部分,但不能证明:

  • 所有目标实例都在线;
  • 所有Bean都读取到了新值;
  • 连接池、线程池等第三方组件已重建;
  • 新配置对应的外部服务可用;
  • 业务请求已经使用新配置成功执行。

生产验证应采用“配置版本 + 实例清单 + ACK/trace + 业务探针”四类证据。ACK数量少于预期要先与服务发现中的健康实例数对齐,再排查离线、ID冲突、目标不匹配和Broker订阅。

二十七、Bus安全治理

busrefresh和环境变更端点属于高权限控制面。若暴露到公网或普通业务网络,攻击者可能全局刷新、注入配置或触发大规模抖动。

最低要求:

  • Actuator端点默认不公开,只通过管理网或网关控制面访问。
  • 使用强认证、最小权限和操作审计。
  • 限制可修改键,禁止通过Bus传播密码、token和私钥。
  • 对全局广播设置审批、限流和变更窗口。
  • Broker Topic/Exchange设置独立ACL,业务应用只拥有必要的读写权限。
  • 事件日志脱敏,同时保留eventId、source、destination、配置版本和操作者。
  • 大规模刷新分批进行,防止所有实例同时重建连接形成惊群。

27.1 Java 17与Spring Boot 3的Bus双实例Demo

在前面Boot 3.2.4、Cloud 2023.0.1工程的基础上增加:

xml
<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-bus-kafka</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-config</artifactId>
</dependency>

配置中心保存order-service.yml,其中包含:

yaml
feature:
  checkout-mode: safe

Bus客户端配置:

yaml
spring:
  application:
    name: order-service
  config:
    import: optional:configserver:http://localhost:8888
  cloud:
    bus:
      id: ${spring.application.name}:${server.port}:${random.value}
      destination: springCloudBus
    stream:
      kafka:
        binder:
          brokers: localhost:9092

management:
  endpoints:
    web:
      exposure:
        include: health,busrefresh

random.value使同机不同实例的Bus ID更难碰撞;Kubernetes环境更适合拼接Pod UID或稳定实例标识。生产环境还要让运维系统能把ID映射回真实实例,否则虽然唯一却无法定向排查。

可刷新的配置读取Bean:

java
@RefreshScope
@RestController
public class CheckoutModeController {
    private final String checkoutMode;

    public CheckoutModeController(
            @Value("${feature.checkout-mode:safe}") String checkoutMode) {
        this.checkoutMode = checkoutMode;
    }

    @GetMapping("/config-view")
    public Map<String, String> configView() {
        return Map.of("checkoutMode", checkoutMode);
    }
}

这个Demo使用Java 17的Map.of;JDK 8应改用普通HashMap@RefreshScope的代理会在刷新时清除旧目标对象,下一次请求重新创建目标并读取新环境值。若把配置值复制到另一个普通单例的不可变字段里,那个单例不会因为此Controller刷新而自动改变。

按8081、8082启动两个实例后,先访问两个实例的/config-view。然后在权威配置中心把值改为fast并发布,再从受保护的管理网络调用:

bash
curl -X POST http://127.0.0.1:8081/actuator/busrefresh

最后再次读取两个实例的/config-view,预期都变为fast。如果只有8081变化,应按Bus部分实例不刷新Runbook检查Broker订阅、Bus ID、目标匹配、ACK和实际配置版本。

这个实验故意不开放busenv。直接远程写环境值绕过权威配置审批、版本记录和启动重载,容易制造“在线时生效、重启后消失”的配置漂移。商业系统应在配置中心完成审核和发布,Bus只广播刷新通知。

JDK 8 + Boot 2.7存量线应使用匹配的Cloud 2021.0.x与Bus 3.1.x;端点暴露、安全配置和依赖名以实际依赖树为准。不要把Bus 4.1.x强行塞入Boot 2.7项目。

二十八、Stream与Bus失败窗口表

窗口表面现象真实后果恢复手段
数据库提交后、消息发送前崩溃订单存在但无事件下游永远不知道Outbox、本地消息表扫描补发
Broker已收、生产方超时发送方认为失败重试可能产生重复eventId、生产幂等与消费幂等
业务提交后、offset/ACK前崩溃消息再次投递重复执行业务同事务去重记录、唯一约束
offset先提交、业务后失败Broker认为已完成业务数据永久丢失禁止错误确认顺序,审计对账补偿
DLQ发布失败主消费也失败消息可能反复或丢隔离证据监控DLQ发布、保留原始消息与告警
Bus源实例本地刷新、发送失败源实例新配置,其他实例旧配置集群配置分裂版本监控、重发、从权威源收敛
Bus目标实例离线没有ACK重启后可能仍是旧配置启动主动拉取最新版本
Bus ID冲突部分实例错误忽略定向投递不可靠生成唯一ID并审计实例清单
刷新触发连接重建失败环境值已变,组件仍不可用半刷新状态Bean级探针、回滚配置或安全重启

二十九、Stream生产排查Runbook

29.1 消息完全不消费

按由外到内的顺序取证:

  1. Broker里destination是否存在、是否真的有消息。
  2. 当前应用连接的是哪个Broker地址和命名空间。
  3. Binding是否处于started/bound状态,destination、group、binder是否正确。
  4. Kafka group是否分到partition,Rabbit Queue是否绑定到正确Exchange。
  5. 监听容器是否反复启动失败、认证失败或反序列化失败。
  6. 函数名是否在spring.cloud.function.definition中,Binding名称是否映射正确。
  7. 消息是否在业务入口前就被转换器拒绝。

29.2 消息重复

先收集同一eventId的Broker元数据、消费实例、时间、重试次数、业务事务结果和确认日志,再区分生产重发、框架重试、Broker重投、人工回放。不要用“MQ重复了”概括所有来源。

29.3 消息卡住或持续积压

同时看Lag分布、消费成功率、处理耗时分位数、线程池、连接池、GC、CPU、下游限流和重平衡。扩消费者前先确认Broker并行度和下游容量。

29.4 恢复动作

  • 毒消息:隔离到DLQ,避免无限阻塞。
  • 下游故障:降低并发或暂停消费,保护数据库,再恢复限速回放。
  • 热分区:调整业务key设计,短期可单独迁移热点流量。
  • 契约不兼容:部署兼容消费者后再回放,不能直接丢消息。
  • offset错误:先备份当前位置和影响范围,再受控重置;错误重置会重复或跳过数据。

三十、Bus部分实例不刷新Runbook

  1. 记录配置版本、事件ID、source、destination、触发时间和预期实例清单。
  2. 验证源实例本地是否刷新、Bus输出Binding是否发送成功。
  3. 在Broker确认Bus destination存在,各目标实例订阅是否在线。
  4. 核对Bus ID是否唯一,目标表达式是否匹配。
  5. 查看每个实例的接收、ACK、刷新异常和Bean重建日志。
  6. 直接读取实例暴露的安全配置版本或业务探针,不能只看ACK。
  7. 对未收敛实例定向刷新;仍失败则从权威配置源重新加载或滚动重启。
  8. 复盘是否存在离线窗口、Broker ACL、序列化版本、刷新作用域或连接重建问题。

全局再次广播不是第一反应。它可能让已正确实例重复刷新并造成连接风暴,却没有解决目标表达式或某个Bean不可刷新的根因。

三十一、必须监控哪些指标

指标或证据
生产端发送速率、失败率、确认耗时、重试次数、事件ID
BrokerTopic/Queue深度、分区Lag、ready/unacked、磁盘、副本、流控
消费端成功率、失败分类、重试、DLQ、处理耗时、并发、重平衡
业务端幂等命中、数据库事务失败、外部依赖延迟、对账差异
Bus事件ID、source、destination、预期实例、ACK、配置版本

消息链路中的trace必须把traceId放入安全的消息头并在消费端恢复上下文。不要把高基数的eventId、orderId直接作为时序指标标签,否则会撑爆指标系统;它们适合日志和追踪字段。

三十二、常见反模式与后果

反模式为什么错后果
不设置业务group实例变成独立订阅同一业务被每个实例执行一次
认为Stream完全屏蔽Broker忽略offset、ACK、分区、路由差异重试和确认设计错误
无限原地重试消费线程一直被毒消息占用同分区后续消息饥饿
先标记已消费再做业务两步不在同一事务业务失败却永远跳过重投
动态destination直接使用租户输入目的地基数不受控Broker和Binding资源爆炸
只看总Lag看不到热分区扩容消费者无效
Bus代替配置中心广播不是权威状态离线实例永久不一致
Bus端点公网开放控制面无保护全局刷新、配置注入和拒绝服务
ACK当作配置生效证明ACK只覆盖事件处理链Bean可能仍使用旧缓存值

三十三、源码阅读路线

现代线可按以下顺序定位,不必一开始钻进Kafka客户端:

  1. BindableFunctionProxyFactory:Binding名和通道怎样创建。
  2. FunctionConfiguration:函数怎样连接输入输出通道。
  3. BindingService:Binder选择与绑定调用。
  4. BinderAbstractMessageChannelBinder:统一契约和模板流程。
  5. KafkaMessageChannelBinderRabbitMessageChannelBinder:具体端点、确认和DLQ。
  6. PartitionHandler:分区键与选择算法。
  7. StreamBridge:主动发送与动态Binding。
  8. BusEnvironmentPostProcessor:Bus环境和函数映射。
  9. RemoteApplicationEventListenerStreamBusBridge:Bus发送链。
  10. BusConsumerPathServiceMatcherRefreshListener:Bus接收、匹配和刷新链。

源码版本必须与项目依赖树一致。同名类在3.2.x和4.1.x中可能存在属性、包结构和行为差异。

三十四、关联知识点

本章小结

Stream的价值是把业务函数、Binding配置和Broker适配分层,而不是取消Broker原理。生产可靠性仍由事件契约、发送确认、Broker持久化、消费确认、幂等、重试、DLQ、积压治理和对账共同组成。Bus复用这条消息链广播系统事件,但它只能加速在线实例感知变化,不能替代权威配置状态。真正可维护的系统必须同时回答“事件怎样流动”和“某一步失败后怎样证明、恢复并最终收敛”。