Stream与Bus内部原理与生产治理
Spring Cloud Stream不是“换一套注解发送MQ”,Spring Cloud Bus也不是“配置中心本身”。Stream在业务函数与Kafka、RabbitMQ之间增加了Binding和Binder抽象;Bus则复用Stream传播Spring Cloud系统事件。只有把应用启动、绑定创建、消息收发、重试确认、配置广播和失败恢复连成一条链,才能判断一条消息到底丢在哪里、为什么重复、为什么扩容后仍积压,以及配置为什么只刷新了部分实例。
本章明确区分两条版本线:
| 版本线 | Java与Spring Boot | Spring Cloud | Stream / Bus样本版本 | 主要用途 |
|---|---|---|---|---|
| 存量线 | JDK 8、Spring Boot 2.7.18 | 2021.0.8 | Stream 3.2.9、Bus 3.1.2 | 维护现有Java 8系统 |
| 现代线 | Java 17、Spring Boot 3.2.4 | 2023.0.1 | Stream 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确认实际版本。
一、学习目标
学完后应能独立回答:
- Function、Binding、Binder、Broker分别解决哪一层问题。
orderConsumer-in-0从哪里生成,应用启动时怎样变成Kafka消费者或Rabbit监听容器。- 同组、异组、匿名组分别怎样投递,为什么Bus不能使用所有实例共享的普通消费组。
- 分区键、分区选择器、Broker分区和消费者并发是什么关系。
- Binding重试、框架消费重试、Broker重投和业务补偿为什么不能混为一谈。
- Kafka Binder与Rabbit Binder在确认、重试、死信和积压上的根本差异。
StreamBridge怎样创建动态输出绑定,为什么无限动态目的地会耗尽资源。/actuator/busrefresh之后事件怎样发送、匹配、刷新和回传ACK。- 为什么Bus广播不能代替配置中心,离线实例怎样最终收敛。
- 如何用指标、日志、Broker工具和Runbook定位消息不消费、重复、积压与部分刷新。
二、先建立四层心智模型
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名称及通道 | 名称不对、函数没有绑定时先看这里 |
FunctionConfiguration | 从FunctionCatalog查找函数并把通道连接到函数 | 判断函数是否真正接入消息链 |
BindingService | 选择Binder,合并属性并调用绑定 | 中间件不可用、Binder选择错误时的核心入口 |
BinderFactory / DefaultBinderFactory | 按名称和目标类型取得Binder,可维护Binder子上下文 | 多Binder系统的选择与隔离 |
Binder | 定义bindConsumer、bindProducer契约 | 统一抽象与具体实现的分界线 |
AbstractMessageChannelBinder | 供应目的地、创建生产处理器或消费端点、错误基础设施 | 生产与消费启动主链 |
ConsumerProperties | group、并发、重试、分区等通用消费属性 | 判断通用重试是否生效 |
ProducerProperties | 分区、requiredGroups、编码、错误通道等通用生产属性 | 判断消息如何选分区 |
PartitionHandler | 提取分区键并计算目标分区 | 同一业务键是否有序 |
StreamBridge | 在普通Service中发送,必要时动态创建输出Binding | 动态目的地与缓存治理 |
BusEnvironmentPostProcessor | 注入Bus函数定义、Binding映射和默认ID | Bus为什么无需手写Consumer |
StreamBusBridge | 把远程应用事件发入Stream输出Binding | Bus发送链入口 |
BusConsumer | 接收、目标匹配、本地发布、ACK和trace | Bus部分实例不刷新时的关键对象 |
四、函数Binding名称怎样生成
函数式模型的默认规则是:
<functionName>-in-<index>
<functionName>-out-<index>例如:
@Bean
public Function<OrderCreated, InventoryCommand> reserveInventory() {
return event -> new InventoryCommand(event.orderId(), event.skuId(), event.quantity());
}会生成:
reserveInventory-in-0
reserveInventory-out-0源码中的BindableFunctionProxyFactory根据输入、输出数量建立代理Binding;FunctionConfiguration随后从FunctionCatalog找出reserveInventory并连接通道。多输入或多输出时索引递增,不要猜名称。
可以显式映射为稳定的业务名:
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
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)大致完成:
- 读取Binding的destination、group和通用
ConsumerProperties。 - 按显式
binder、默认Binder和目标通道类型选择具体Binder。 - 如果Binder支持扩展属性,再合并Kafka或Rabbit专属配置。
- 调用
binder.bindConsumer(destination, group, input, properties)。 - 具体Binder创建Topic/Queue等资源、监听容器和消息适配端点。
- 消费端点输出被连接到Stream输入通道,收到的Broker消息才能进入业务函数。
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、合并生产属性并调用bindProducer。AbstractMessageChannelBinder的典型过程是:
provisionProducerDestination供应或校验目的地。createProducerMessageHandler创建具体发送处理器。- 初始化并启动处理器。
- 将
SendingHandler订阅到输出通道。 - 发布
BindingCreatedEvent。
之后业务函数把消息发入输出通道,真正发送由处理器完成。业务方法“返回了对象”不等于Broker已经持久化,更不等于下游业务已经成功。
5.3 Binding创建重试不等于消息重试
如果应用启动时Broker暂不可达,BindingService可以按bindingRetryInterval延迟重试创建Binding,并用LateBinding占位。这一重试解决的是“连接和监听容器还没有建起来”,不是“某条消息处理失败”。
| 重试层 | 发生时间 | 重试对象 | 典型问题 |
|---|---|---|---|
| Binding创建重试 | 应用启动或重绑 | 建立生产者/消费者Binding | Broker暂不可达 |
| Stream消费重试 | 收到某条消息后 | 再调用业务函数 | 临时数据库异常 |
| Broker重投 | 消费失败或未确认后 | 同一Broker消息 | Kafka offset未提交、Rabbit未ACK |
| 业务补偿 | 进入重试Topic、任务表后 | 一次业务操作 | 长时间外部依赖故障 |
把四层同时打开可能产生乘法效应。例如框架3次、Broker再投5次、补偿任务又3次,最坏可能执行45次,因此幂等不是可选项。
六、Binder怎样选择,多Binder为什么容易出错
单Binder时依赖通常足够决定实现;多Binder时必须明确配置。选择过程可理解为:
flowchart TD
A["读取Binding配置"] --> B["优先使用显式binder"]
B --> C["否则使用defaultBinder"]
C --> D["仍无则按目标类型推断"]
D --> E["校验唯一候选并取得Binder"]现代线多Binder示意:
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要按匿名、互不共享的订阅处理。
flowchart TD
A["同一业务事件"] --> B["库存组与积分组各得一份"]
B --> C["每个组内由一个实例处理"]因此:
- 同一服务的多个实例通常使用相同group,实例间竞争,整组处理一份逻辑消息。
- 不同业务服务使用不同group,每个group各得到一份逻辑消息。
- 不写group会创建匿名非共享订阅,每个实例都可能收到一份;具体是否持久、离线是否保留取决于Binder和Broker配置,不能笼统承诺。
错误示例:库存服务每次发布都生成随机group。结果是旧group不断遗留、每个实例都消费全量消息,既重复扣库存又制造Broker资源。
Bus为了让每个在线实例都收到广播,正是利用非共享消费语义;这也带来“离线实例可能错过事件”的边界,后文会详细解释。
八、分区原理:同一订单为什么进入同一条处理通道
Stream通用分区链是:
flowchart TD
A["消息"] --> B["提取partition key"]
B --> C["选择器计算原始值"]
C --> D["对partitionCount取模并归一化"]
D --> E["写入目标分区信息"]
E --> F["Binder映射到Broker分区或队列"]PartitionHandler先使用自定义PartitionKeyExtractorStrategy或partitionKeyExpression取得键,再使用自定义选择器、表达式或默认哈希,最后归一化为0..partitionCount-1。默认实现还处理了Integer.MIN_VALUE绝对值异常边界。
spring:
cloud:
stream:
bindings:
publishOrder-out-0:
destination: order.created.v1
producer:
partition-key-expression: headers['orderId']
partition-count: 12必须理解四个限制:
- 只保证同一键稳定映射,不保证全局顺序。
- Kafka实际Topic分区数不能小于设计值;修改分区数后哈希映射会变化。
- 同一分区通常同时只能被组内一个消费者实例占有,消费者数超过分区数会有空闲实例。
- 热门键会形成热分区,整体消费者数量再多也无法并行处理该键。
若业务要求同一订单严格有序,应以orderId为键,并在消费端把数据库状态机作为最终约束。仅依靠队列顺序不能防止重试、人工补偿和跨Topic造成的乱序。
九、消息转换、contentType与原生编码
函数收到Java对象之前通常经历:Broker字节、Spring Message、内容类型转换、函数参数适配。contentType描述载荷如何序列化/反序列化,默认转换和Binder原生编码是两条不同路径。
flowchart TD
A["Java事件对象"] --> B["Message转换"]
B --> C["字节或原生客户端对象"]
C --> D["Broker"]
D --> E["消费端反序列化"]
E --> F["Consumer参数"]建议事件契约至少包含:eventId、eventType、eventVersion、occurredAt、traceId和业务主键。不要直接把数据库实体序列化为跨服务契约,因为字段删除、懒加载代理和内部状态会泄漏到消费者。
事件演进原则:
- 新增可选字段通常向后兼容。
- 删除、改名、改变类型要升版本或做双读。
- 消费者要容忍未知字段。
- 时间、金额、枚举必须明确格式与单位。
- PII、token和密钥不得写入消息体或日志。
启用useNativeEncoding后序列化更多交给具体Broker客户端。它能使用Kafka Serializer等原生能力,但也使迁移Binder更困难。切换前必须让生产者和消费者对序列化协议达成一致。
十、StreamBridge为什么适合业务Service发送消息
业务消息通常由HTTP请求或数据库事务触发,不适合轮询式Supplier。StreamBridge允许普通Service主动发送:
@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);
}
}内部主链是:
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=1、maxAttempts=3、初始退避1000ms、最大10000ms、倍数2.0。maxAttempts包含第一次调用,因此设为1才是关闭框架重试。
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一次失败的正确决策:
flowchart TD
A["业务处理抛异常"] --> B["判断是否短暂可恢复"]
B --> C["可恢复则有限指数退避"]
C --> D["仍失败则进入DLQ或重试Topic"]
D --> E["告警、修复、受控回放"]错误通道承载的是消息处理异常信息,不等于持久化死信队列。若错误处理器本身失败、实例重启或错误通道没有持久消费者,异常记录仍可能丢失。需要长期保留和回放时,应使用Broker DLQ、重试Topic或数据库补偿表。
以下异常通常不应盲目原地重试:参数校验失败、事件版本不支持、业务唯一约束明确冲突、永久权限错误。原地重试会占用消费线程并阻塞同分区后续消息。
十二、至少一次投递下消费幂等怎样落地
可靠消息系统通常追求at-least-once,这意味着“可能重复,但不能静默丢失”。消费者必须让同一eventId重复到达时得到相同业务结果。
数据库去重与业务操作应在同一个本地事务内:
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)
);@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 ACK | Broker认为这条消息或这个位置已经处理 |
这三件事没有天然原子性。数据库事务提交不了Kafka offset,Rabbit ACK也回滚不了MySQL事务,所以必须人为设计顺序和幂等。
flowchart TD
A["Broker投递消息"] --> B["Stream调用业务函数"]
B --> C["开启本地事务"]
C --> D["写消费日志和业务数据"]
D --> E["提交数据库事务"]
E --> F["提交Kafka offset或发送Rabbit ACK"]
F --> G["Broker不再正常重投"]最危险的顺序是“先确认Broker,再执行业务”:
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的差异要说清楚:
| 对比项 | Kafka | RabbitMQ |
|---|---|---|
| 确认对象 | consumer group在partition上的offset位置 | Queue中某条投递消息的delivery tag |
| 确认后果 | 下次从更后面的offset继续poll | Broker可从队列移除这条已ACK消息 |
| 未确认后果 | offset未前进,重启或重平衡后可能重读 | 连接断开或NACK后可重新入队或进死信 |
| 排查重点 | committed offset、log end offset、Lag、rebalance | ready、unacked、prefetch、redelivered、DLQ |
所以Stream消费端的生产原则是:
- 业务数据和消费日志放在同一个本地数据库事务里。
- 消费日志使用
consumer_group + event_id唯一键。 - 事务提交成功后再允许Binder提交Kafka offset或Rabbit ACK。
- 如果ACK失败导致重投,第二次根据消费日志直接返回成功,让Broker完成确认。
- 如果业务失败,不要吞异常伪装成功,否则框架可能提交确认。
- 对不可恢复错误要进入DLQ或补偿表,不能无限占住主分区或主队列。
一个更贴近商业项目的消费日志表应至少包含状态,而不是只有一条“已消费”:
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,直接确认,不会重复扣库存。
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、业务错误分类 |
面试里可以这样回答:
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和并发。
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属性已逐步废弃,不要用旧文章直接套现代线。enableDlq和dlqName可启用Kafka Binder死信发布,但要验证失败发布、offset提交和事务配置的组合。startOffset决定新group从哪里开始;resetOffsets会主动重置,误用可能大量重放或跳过数据。- 消费并发的有效上限受分区数限制。
- Kafka顺序只在单分区内成立。
Kafka消费成功的业务边界通常是“数据库事务提交后再允许offset确认”。如果先提交offset再写库,宕机会丢业务;如果先写库后offset提交失败,会重复,所以必须幂等。
完整底层原理见Kafka专栏。
十四、Rabbit Binder到底映射了什么
Rabbit Binder需要把destination、group映射为Exchange、Queue和Binding关系,再通过监听容器消费并ACK/NACK。
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怎么选
| 维度 | Kafka | RabbitMQ |
|---|---|---|
| 核心模型 | 分区追加日志,消费者维护offset | Exchange路由到Queue,Broker跟踪投递与确认 |
| 广播方式 | 不同consumer group各读一份 | 不同Queue绑定同一Exchange |
| 顺序 | 分区内有序 | 单Queue单消费者较直观,多消费者会并发 |
| 重放 | 保留期内调整offset重放 | 消息ACK后通常从队列删除,需另建归档/重试机制 |
| 高吞吐日志流 | 更擅长 | 可用但模型不同 |
| 灵活路由 | 依赖Topic、key和应用设计 | Exchange、routing key能力强 |
| 积压观察 | consumer group lag | ready、unacked、publish/deliver/ack速率 |
| 扩容上限 | 受Topic分区数约束 | 受Queue、prefetch、消费者和下游约束 |
选择不能只看“哪个吞吐高”。订单事件需要长时间重放、数据管道和大吞吐时常偏Kafka;复杂路由、工作队列、低延迟命令分发常偏Rabbit。若项目强依赖某一方高级语义,直接使用Spring Kafka或Spring AMQP有时比强行抽象更清晰。
十六、消息堆积为什么扩容后只短暂有效
积压变化可用一个简单关系理解:
积压增长速度 = 生产速率 - 实际确认速率扩容后短暂下降、随后继续增长,常见原因不是“扩容没生效”,而是实际确认速率再次低于生产速率:
- Kafka消费者数已超过分区数,多出的实例空闲。
- Lag集中在少数热分区,新增消费者无法拆分单分区。
- 数据库、Redis、HTTP下游到达瓶颈,消费者越多争用越严重。
- 连接池小于消费并发,大量线程等待连接。
- 失败消息反复重试,表面消费很快但有效成功率低。
- 单条消息体积或批次变大,平均处理时间上升。
- GC停顿、CPU限流、容器内存压力使实例实际吞吐下降。
- 扩容期间重平衡暂停消费,频繁自动伸缩造成抖动。
- Rabbit prefetch过大,消息堆在少数实例的unacked区。
- 生产侧流量也同步上涨,扩容增加量小于新增流量。
“Lag分布”不是只看总Lag,而是查看每个Topic、group、partition的log-end-offset - committed-offset及其变化率。总Lag 100万可能均匀分在100个分区,也可能99万集中在一个热分区,两个场景的扩容策略完全不同。
16.1 积压Runbook
第一步,冻结证据:记录时间窗、发布速率、确认速率、总Lag、分区Lag分布、消费者实例数、分区数、重平衡次数、失败率和下游延迟。
第二步,先算理论吞吐:
单实例理论吞吐约等于 并发数 / 平均处理秒数
有效吞吐还要乘以成功率,并受分区数和下游容量限制第三步,区分瓶颈:
| 证据 | 更可能的根因 | 合理动作 |
|---|---|---|
| 少数分区Lag很高 | 热key或分区倾斜 | 调整key、拆Topic、业务键分片 |
| 所有分区均匀增长 | 总容量不足 | 在下游允许范围内扩消费者 |
| DB连接池等待高 | 数据库/连接池瓶颈 | 优化SQL、批量写、限并发,不能盲目扩容 |
| 重试次数激增 | 毒消息或依赖故障 | 暂停热循环、隔离DLQ、修复依赖 |
| CPU低但消费慢 | I/O等待或锁竞争 | 查线程栈、连接池、锁与远程调用 |
| 扩容伴随频繁rebalance | 伸缩和消费参数不稳 | 稳定实例、调整poll与处理边界 |
第四步,恢复后受控回放。不要一次把消费并发拉满冲垮数据库;先限速消化,观察业务成功率、DB负载和Lag斜率,再逐步提高。
更完整专题见消息堆积、扩容与背压。
十七、命令式函数与响应式函数不能用同一套直觉
命令式Consumer<T>通常由可订阅通道触发一次消息调用;响应式Function<Flux<T>, Flux<R>>连接的是Publisher链。响应式函数不是“每条消息反射调用一次普通方法”,订阅、背压、线程切换和错误传播边界不同。
@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 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 事件与发送接口
启动类:
@SpringBootApplication
public class StreamOrderApplication {
public static void main(String[] args) {
SpringApplication.run(StreamOrderApplication.class, args);
}
}public record OrderCreated(
String eventId,
long orderId,
long skuId,
int quantity,
Instant occurredAt,
int eventVersion) {
}@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 幂等消费者
@Configuration
public class OrderConsumerConfiguration {
@Bean
Consumer<OrderCreated> reserveInventory(InventoryApplicationService service) {
return service::reserveOnce;
}
}@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实现:
public interface ConsumedEventRepository {
boolean tryInsert(String consumerGroup, String eventId);
}
public interface InventoryRepository {
int reserveIfEnough(long skuId, int quantity);
}@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;
}
}
}@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:
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 配置
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后运行应用,再调用:
curl -i -X POST http://localhost:8080/orders/9001/eventsDemo会向order.created.v1发送事件,再由同一应用的reserveInventory函数消费并扣减H2库存。若要直接验证重复消息,可临时让Controller接收固定eventId并发送两次;最终库存只能扣一次。生产测试还必须覆盖:消费函数抛异常后的实际尝试次数、DLQ消息内容、Broker重启、应用在业务提交后确认前崩溃、以及幂等表与业务表是否始终同事务回滚。
十九、JDK 8与Spring Boot 2.7存量线差异
存量线可以继续使用函数式模型,但代码和版本不能照搬现代线:
| 项目 | JDK 8 + Boot 2.7 | Java 17 + Boot 3.x |
|---|---|---|
| Java最低版本 | 8 | 17 |
| 事件DTO | 普通POJO | 可用Record,也可继续POJO |
| EE包名 | javax.* | jakarta.* |
| Cloud发布列车样本 | 2021.0.8 | 2023.0.1 |
| Stream样本 | 3.2.9 | 4.1.1 |
| 旧注解模型 | 遗留项目可能仍有@EnableBinding | 应迁移到函数式模型 |
| 可观测性 | Sleuth时代常见 | Micrometer Tracing方向 |
JDK 8事件必须写成POJO:
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在环境准备阶段追加并映射函数:
spring.cloud.function.definition += busConsumer
busConsumer-in-0 -> springCloudBusInput
springCloudBusInput.destination -> springCloudBus
springCloudBusOutput.destination -> springCloudBus默认常量包括:
| 名称 | 默认值 |
|---|---|
| 输入Binding | springCloudBusInput |
| 输出Binding | springCloudBusOutput |
| Broker destination | springCloudBus |
| 消费函数 | busConsumer |
它还会推导spring.cloud.bus.id。ID用于目标匹配和判断事件是否来自自己,生产环境必须足够唯一;若多个实例ID碰撞,可能错误忽略事件或把单实例事件投递给错误对象。
二十二、/actuator/busrefresh发送链
flowchart TD
A["调用busrefresh端点"] --> B["发布RefreshRemoteApplicationEvent"]
B --> C["源实例按目标执行本地刷新"]
C --> D["远程事件监听器继续传播"]
D --> E["StreamBusBridge发送输出Binding"]
E --> F["Binder发布到springCloudBus"]关键细节:
- 端点先在源实例本地发布远程事件。
- 如果目标匹配源实例,本地
RefreshListener可以立即调用ContextRefresher.refresh()。 RemoteApplicationEventListener忽略ACK事件,并避免把不该发送的事件再次传播。StreamBusBridge经StreamBridge向springCloudBusOutput发送。- Binder再把事件编码并发布到Kafka或RabbitMQ。
因此HTTP端点返回成功通常只能说明事件已在源实例触发并尝试发送,不能证明所有目标实例都刷新成功。
二十三、目标实例接收与刷新链
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实例都能接收事件。
flowchart TD
A["springCloudBus事件"] --> B["每个在线实例有独立匿名订阅"]
B --> C["实例A、B、C各自收到一份"]代价是匿名订阅的持久性取决于Binder/Broker实现。实例在发布时离线、网络隔离或订阅尚未建立,可能错过这次广播。正确恢复机制是:
- 配置中心保存权威版本。
- 实例启动时主动拉取最新配置。
- Bus用于加速在线实例感知变化。
- 监控每个实例实际配置版本,而不是只看端点HTTP状态。
- 发现版本不一致时重试定向刷新或安全重启实例。
这就是“事件通知”和“状态收敛”的区别。事件可能丢,权威状态必须可重读。
二十六、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工程的基础上增加:
<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,其中包含:
feature:
checkout-mode: safeBus客户端配置:
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,busrefreshrandom.value使同机不同实例的Bus ID更难碰撞;Kubernetes环境更适合拼接Pod UID或稳定实例标识。生产环境还要让运维系统能把ID映射回真实实例,否则虽然唯一却无法定向排查。
可刷新的配置读取Bean:
@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并发布,再从受保护的管理网络调用:
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 消息完全不消费
按由外到内的顺序取证:
- Broker里destination是否存在、是否真的有消息。
- 当前应用连接的是哪个Broker地址和命名空间。
- Binding是否处于started/bound状态,destination、group、binder是否正确。
- Kafka group是否分到partition,Rabbit Queue是否绑定到正确Exchange。
- 监听容器是否反复启动失败、认证失败或反序列化失败。
- 函数名是否在
spring.cloud.function.definition中,Binding名称是否映射正确。 - 消息是否在业务入口前就被转换器拒绝。
29.2 消息重复
先收集同一eventId的Broker元数据、消费实例、时间、重试次数、业务事务结果和确认日志,再区分生产重发、框架重试、Broker重投、人工回放。不要用“MQ重复了”概括所有来源。
29.3 消息卡住或持续积压
同时看Lag分布、消费成功率、处理耗时分位数、线程池、连接池、GC、CPU、下游限流和重平衡。扩消费者前先确认Broker并行度和下游容量。
29.4 恢复动作
- 毒消息:隔离到DLQ,避免无限阻塞。
- 下游故障:降低并发或暂停消费,保护数据库,再恢复限速回放。
- 热分区:调整业务key设计,短期可单独迁移热点流量。
- 契约不兼容:部署兼容消费者后再回放,不能直接丢消息。
- offset错误:先备份当前位置和影响范围,再受控重置;错误重置会重复或跳过数据。
三十、Bus部分实例不刷新Runbook
- 记录配置版本、事件ID、source、destination、触发时间和预期实例清单。
- 验证源实例本地是否刷新、Bus输出Binding是否发送成功。
- 在Broker确认Bus destination存在,各目标实例订阅是否在线。
- 核对Bus ID是否唯一,目标表达式是否匹配。
- 查看每个实例的接收、ACK、刷新异常和Bean重建日志。
- 直接读取实例暴露的安全配置版本或业务探针,不能只看ACK。
- 对未收敛实例定向刷新;仍失败则从权威配置源重新加载或滚动重启。
- 复盘是否存在离线窗口、Broker ACL、序列化版本、刷新作用域或连接重建问题。
全局再次广播不是第一反应。它可能让已正确实例重复刷新并造成连接风暴,却没有解决目标表达式或某个Bean不可刷新的根因。
三十一、必须监控哪些指标
| 层 | 指标或证据 |
|---|---|
| 生产端 | 发送速率、失败率、确认耗时、重试次数、事件ID |
| Broker | Topic/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客户端:
BindableFunctionProxyFactory:Binding名和通道怎样创建。FunctionConfiguration:函数怎样连接输入输出通道。BindingService:Binder选择与绑定调用。Binder、AbstractMessageChannelBinder:统一契约和模板流程。KafkaMessageChannelBinder或RabbitMessageChannelBinder:具体端点、确认和DLQ。PartitionHandler:分区键与选择算法。StreamBridge:主动发送与动态Binding。BusEnvironmentPostProcessor:Bus环境和函数映射。RemoteApplicationEventListener、StreamBusBridge:Bus发送链。BusConsumer、PathServiceMatcher、RefreshListener:Bus接收、匹配和刷新链。
源码版本必须与项目依赖树一致。同名类在3.2.x和4.1.x中可能存在属性、包结构和行为差异。
三十四、关联知识点
- Stream与Bus入门
- Stream与Bus独立面试题
- Kafka架构、可靠性与生产治理
- RabbitMQ ACK、重试与死信
- 消息堆积、扩容与背压
- 请求、消息与任务幂等
- 分布式事务与本地消息表
- 配置中心内部原理
- Nacos配置与服务治理
本章小结
Stream的价值是把业务函数、Binding配置和Broker适配分层,而不是取消Broker原理。生产可靠性仍由事件契约、发送确认、Broker持久化、消费确认、幂等、重试、DLQ、积压治理和对账共同组成。Bus复用这条消息链广播系统事件,但它只能加速在线实例感知变化,不能替代权威配置状态。真正可维护的系统必须同时回答“事件怎样流动”和“某一步失败后怎样证明、恢复并最终收敛”。
