Skip to content

Stream与Bus独立面试题

本页只保留面试中的标准回答、追问方向和原理入口。回答时先给结论,再按“抽象层次、执行链、失败边界、生产治理”展开;完整源码与Demo统一跳转到内部原理与生产治理

版本口径必须先说清楚:存量系统常见JDK 8、Spring Boot 2.7、Spring Cloud 2021.0.x和Stream 3.2.x;现代系统以Java 17+、Spring Boot 3.x、匹配的Spring Cloud发布列车和Stream 4.x为主。不要把旧注解模型或废弃属性当成现代推荐方案。

一、核心模型

1. Spring Cloud Stream是什么

标准回答:Stream是消息驱动微服务的编程模型抽象。业务代码面对SupplierFunctionConsumer和Spring Message,Binding描述函数端口到destination/group的映射,Binder把统一模型适配到Kafka、RabbitMQ等Broker。它减少样板代码,但不会消除分区、offset、ACK、死信和Broker选型差异。

常见追问:Stream与Spring Kafka有什么区别?能否完全切换Broker?

原理入口:四层心智模型Binder对比

2. Function、Binding、Binder和Broker分别是什么

标准回答:Function承载业务处理;Binding是逻辑输入输出端口及其配置;Binder负责把端口连接到具体中间件;Broker负责消息存储、路由、分区、副本和确认。业务成功与Broker确认不在同一层,必须额外设计事务和幂等。

常见追问:哪个对象负责创建消费者?哪个对象真正保存消息?

原理入口:核心源码对象

3. 为什么推荐函数式模型

标准回答:函数式模型使用标准SupplierFunctionConsumer,降低对旧Stream注解API的耦合,便于函数组合、测试和Spring Cloud Function集成。Stream 4.x应优先采用函数式模型;遗留@EnableBinding代码应分阶段迁移,而不是在升级Boot 3时继续扩写。

常见追问:响应式函数和普通Consumer执行方式一样吗?

原理入口:函数Binding名称命令式与响应式

4. orderConsumer-in-0怎样生成

标准回答:BindableFunctionProxyFactory<functionName>-in-<index><functionName>-out-<index>生成默认Binding名,多输入输出时索引递增。可用spring.cloud.stream.function.bindings映射为业务名,再在bindings下配置destination和group。

常见追问:函数名、Binding名和Topic名是否相同?

原理入口:Binding名称生成

5. 应用启动时怎样创建消费者

标准回答:框架先发现函数并创建输入通道,Binding生命周期调用BindingService.bindConsumer;它读取destination/group、选择Binder、合并通用和Binder专属属性,随后具体Binder供应Broker资源、创建并启动监听端点,最后把端点输出连接到输入通道和业务函数。

常见追问:Broker暂不可达时会怎样?函数为什么有Bean却不消费?

原理入口:启动Binding主链

6. 生产者Binding怎样创建

标准回答:BindingService.bindProducer选择Binder并调用bindProducer;Binder供应目的地、创建并启动生产消息处理器,再把处理器订阅到输出通道。业务函数返回对象或StreamBridge.send进入通道后,消息才沿处理器发到Broker。

常见追问:方法返回成功是否代表Broker和消费者都成功?

原理入口:生产端启动链

7. Binding重试和消息重试有什么区别

标准回答:Binding重试是在Broker不可用时重试建立生产/消费连接;消息重试是在某条消息进入业务函数后重试处理。前者重试“基础设施绑定”,后者重试“业务消息”,再加上Broker重投和补偿任务可能形成多层重试乘法。

常见追问:怎样避免3×5×3次重复执行?

原理入口:四类重试边界消费重试

8. 多Binder如何选择

标准回答:Binding可以显式指定binder;否则使用default-binder或按候选和目标类型推断。多Binder系统应显式标注,并把不同集群的环境、认证和专属属性隔离,避免依赖推断导致消息发往错误Broker。

常见追问:Binder子上下文有什么代价?

原理入口:Binder选择与多Binder

二、消费组、分区与消息契约

9. 同组和异组怎样消费

标准回答:同一个非空group共享逻辑订阅,组内实例竞争处理;不同group各得到一份逻辑消息。因此同一业务服务的实例应使用相同group,不同下游业务使用不同group。

常见追问:库存和积分如何都收到订单事件?

原理入口:Group真实语义

10. 不配置group会怎样

标准回答:Binder契约把空group视为匿名、非共享订阅,每个实例可能独立收到消息;具体持久性取决于Binder和Broker。普通业务消费者不应随意省略group,否则扩一个实例可能多执行一次业务。

常见追问:为什么Bus恰好需要类似广播语义?

原理入口:匿名订阅Bus广播

11. Stream分区是怎样计算的

标准回答:PartitionHandler先通过提取策略或表达式得到分区键,再通过选择器或默认哈希计算原始值,最后按partitionCount取模归一化。Binder把结果映射到Kafka分区或对应的分区资源。

常见追问:为什么修改分区数会让同一key映射变化?

原理入口:分区内部原理

12. 配置分区后能保证全局有序吗

标准回答:不能。它最多保证相同key进入同一处理通道,Kafka也只保证单分区顺序。多分区、多消费者、失败重试、重试Topic和人工补偿都可能改变全局顺序,业务还需状态机和版本号防止旧事件覆盖新状态。

常见追问:订单状态事件怎样防乱序?

原理入口:分区限制

13. 消费者扩容上限由什么决定

标准回答:Kafka同一group有效并行度不超过Topic分区数,且热分区可能成为单点瓶颈;Rabbit还受Queue、prefetch、消费者并发和下游吞吐约束。实例数只是容量变量之一,不能脱离Broker并行度和数据库容量盲目扩容。

常见追问:为什么新实例处于空闲?

原理入口:消息积压与扩容

14. contentType和原生编码有什么区别

标准回答:默认路径由Spring消息转换依据contentType把Java对象和消息载荷互转;启用原生编码后更多交给Kafka Serializer等Broker客户端。原生能力更直接,但会增加对具体Binder和序列化协议的耦合。

常见追问:跨版本DTO怎样兼容?

原理入口:消息转换与契约

15. 事件为什么要有eventId和eventVersion

标准回答:eventId用于端到端追踪和消费幂等,eventVersion用于契约演进和消费者兼容判断。只有业务主键不能区分同一订单的多次不同事件,也不能可靠识别生产重发。

常见追问:能否直接发送数据库实体?

原理入口:消息契约消费幂等

三、发送、重试、死信与一致性

16. StreamBridge解决什么问题

标准回答:它让HTTP请求、事务Service等非函数输出场景主动发送消息。它会查找Binding通道,必要时选择Binder、转换消息、计算分区并动态创建生产Binding。

常见追问:什么时候用Supplier,什么时候用StreamBridge?

原理入口:StreamBridge内部链

17. StreamBridge.send()返回true代表发送成功吗

标准回答:它首先表示当前消息通道接受了消息,实际Broker确认取决于通道、Binder和客户端配置;它不代表下游消费者处理成功,更不代表下游数据库事务已提交。业务接口不能把这个布尔值宣传成端到端成功。

常见追问:怎样获得真正的发送可靠性?

原理入口:send返回值边界失败窗口

18. 动态destination有什么风险

标准回答:StreamBridge可按运行时名称创建并缓存生产Binding,但无限租户ID或用户输入会造成Binding缓存、Topic/Exchange、权限和指标标签爆炸。生产上应固定目的地,或对白名单、基数、命名和生命周期做强约束。

常见追问:缓存淘汰时会怎样?

原理入口:动态目的地治理

19. maxAttempts=3是重试三次还是总共三次

标准回答:总共最多尝试三次,包含第一次业务调用;设为1才关闭框架重试。还要检查Kafka/Rabbit Broker重投和业务补偿层,避免总尝试次数相乘。

常见追问:退避参数怎样配置?

原理入口:消费重试

20. 错误通道等于DLQ吗

标准回答:不等于。错误通道是Spring消息错误处理基础设施,默认不保证像Broker队列一样持久。需要长期隔离、告警和回放时,应使用Kafka/Rabbit DLQ、重试Topic或数据库补偿表。

常见追问:错误处理器自己失败怎么办?

原理入口:错误通道与死信

21. 哪些异常不应原地重试

标准回答:永久参数错误、不支持的事件版本、明确业务规则失败和永久权限错误不应长时间原地重试;它们会占住消费线程并阻塞同分区消息。应分类后快速隔离,修复数据或消费者再受控回放。

常见追问:数据库超时和库存不足分别怎么处理?

原理入口:重试决策

22. 消费端为什么必须幂等

标准回答:至少一次投递下,业务提交后ACK/offset提交前宕机、生产重试和人工回放都会产生重复。消费者应以eventId和consumer group建立唯一记录,并让去重记录与业务更新处于同一本地事务。

常见追问:Redis SETNX去重是否足够?

原理入口:消费幂等落地

22.1 Spring Cloud Stream消费成功、ACK和业务事务是什么关系

标准回答:Stream消费成功不能只理解成Consumer方法没有抛异常。生产上应先让业务数据和消费日志在同一个本地事务中提交,再让Binder提交Kafka offset或Rabbit ACK。业务提交后ACK失败会导致重复投递,消费者必须用consumer_group + event_id唯一键幂等吸收;如果先ACK再写库,写库失败就会造成消息丢失。Stream统一编程模型,但不会取消Kafka offset和Rabbit ACK的底层差异。

常见追问:业务成功但ACK失败怎么办?为什么不能异常被捕获后直接返回成功?Kafka和Rabbit确认语义有什么差异?

原理入口:消费确认、Offset/ACK与业务事务顺序消费幂等落地

23. 为什么不能先写Redis已消费再更新数据库

标准回答:如果Redis成功而数据库失败,消息重投会被Redis挡住,形成业务数据丢失。Redis可用于快速过滤,但核心去重状态必须与业务结果具有同一事务边界,或有可验证的补偿机制。

常见追问:唯一索引冲突怎样处理?

原理入口:幂等事务边界

24. 数据库提交和消息发送如何保证一致

标准回答:普通本地事务无法原子覆盖数据库与Broker。常用方案是Outbox/本地消息表:业务数据和待发送事件同事务提交,再由可靠任务投递并按eventId幂等;也可按Broker和业务条件选择事务消息。不能仅靠“提交后立即send”。

常见追问:发送超时但Broker已收到怎么办?

原理入口:失败窗口分布式事务

四、Kafka Binder与Rabbit Binder

25. Kafka Binder怎样映射Stream概念

标准回答:destination通常映射Topic,group映射consumer group,分区信息映射Kafka partition;监听容器负责poll和offset提交。消费可靠性必须结合ack mode、offset提交时机、DLQ和数据库事务分析。

常见追问:新group从earliest还是latest开始?

原理入口:Kafka Binder映射

26. Kafka的ackMode为什么重要

标准回答:它决定监听容器何时确认和提交offset,直接影响失败后重复还是丢失。现代版本应查看实际ackMode配置,不能继续套用已废弃的autoCommitOffset旧文章。

常见追问:业务事务提交后offset失败会怎样?

原理入口:Kafka确认边界失败窗口

27. Rabbit的requeueRejected=true有什么风险

标准回答:失败消息可能立即回到原队列并再次投递,永久错误会形成无退避热循环,占满消费者并刷爆日志。生产中通常采用有限重试后DLQ或延迟重试队列,并确保消费幂等。

常见追问:republishToDlq与Broker死信有什么区别?

原理入口:Rabbit Binder映射

28. autoBindDlqrepublishToDlq分别做什么

标准回答:autoBindDlq让Binder供应DLX/DLQ拓扑;republishToDlq让Binder把失败消息重新发布并附带异常头。前者是资源创建,后者是失败消息处理,两者不能混为一谈,且都要考虑权限和DLQ发布失败。

常见追问:异常头太大有什么影响?

原理入口:Rabbit死信配置

29. Kafka和Rabbit如何选

标准回答:Kafka基于分区追加日志,适合高吞吐事件流、保留与重放;Rabbit通过Exchange到Queue灵活路由,适合工作队列和复杂路由。应比较重放、顺序、路由、积压模型、团队运维能力和业务语义,而不是只比较峰值TPS。

常见追问:使用Stream后是否可以无成本切换?

原理入口:Kafka与Rabbit对比

五、积压与生产排查

30. 什么是Lag分布

标准回答:Kafka Lag分布是按Topic、consumer group、partition观察log-end-offset - committed-offset及其变化率,而不只是总Lag。均匀增长代表总容量不足的可能性较高,集中在少数分区通常代表热key、慢消息或分区倾斜。

常见追问:总Lag相同为什么恢复时间不同?

原理入口:Lag分布与积压

31. 扩容消费者后短暂有效,后来又积压是什么原因

标准回答:扩容只暂时提高了消费速率,随后生产速率再次超过有效确认速率;也可能受分区上限、热分区、数据库/连接池瓶颈、重试风暴、GC、频繁rebalance或Rabbit unacked倾斜限制。必须比较每分区Lag、成功率、处理耗时和下游容量。

常见追问:先加分区还是先扩实例?

原理入口:扩容短暂有效的十类原因积压Runbook

32. 为什么消费者越多数据库越慢

标准回答:消费者增加并发请求后,数据库连接池、锁、索引页、日志刷盘和CPU可能成为共享瓶颈;等待和冲突使单条耗时上升,最终有效吞吐反而下降。应从端到端容量约束确定并发,而不是让消息线程无限抢资源。

常见追问:连接池大小和消费并发如何配合?

原理入口:积压诊断

33. 消息完全不消费从哪里开始查

标准回答:先确认Broker目的地和消息,再核对应用连接的集群、Binding状态、destination/group/binder、分区或Queue绑定、监听容器异常、函数定义和反序列化错误。应由Broker到Binder再到函数逐层收缩,不要只在业务方法打断点。

常见追问:有消息但group没有成员说明什么?

原理入口:Stream生产Runbook

34. 如何安全回放DLQ

标准回答:先修复根因并确认消费者向后兼容,再统计影响、备份原消息、启用幂等、限速小批回放,观察业务成功率和下游负载;不能直接把全部DLQ倒回主Topic。回放事件要保留原eventId并记录回放批次。

常见追问:回放期间新消息怎么办?

原理入口:恢复动作

六、Spring Cloud Bus

35. Spring Cloud Bus是什么,与Stream什么关系

标准回答:Bus是Spring Cloud系统事件总线,用于配置刷新等控制事件;它复用Stream和Kafka/Rabbit Binder传播RemoteApplicationEvent。Stream是通用业务消息抽象,Bus是构建在其上的系统控制面能力。

常见追问:能否用Bus发订单消息?

原理入口:Bus职责边界

36. /actuator/busrefresh之后发生什么

标准回答:端点在源实例发布RefreshRemoteApplicationEvent,源实例按目标可先本地刷新;远程事件监听器通过StreamBusBridge和输出Binding发到Broker。目标实例的busConsumer接收、匹配目标、在本地发布事件,RefreshListener再调用上下文刷新逻辑。

常见追问:源实例为什么不会刷新两次?

原理入口:Bus发送链接收链

37. Bus为什么不需要自己手写Consumer

标准回答:BusEnvironmentPostProcessor会把busConsumer追加到函数定义,把busConsumer-in-0映射为springCloudBusInput,并把输入输出绑定到默认springCloudBus目的地,同时推导Bus ID。

常见追问:默认输入输出Binding叫什么?

原理入口:Bus启动自动配置

38. Bus如何定向到一个服务或实例

标准回答:Bus使用基于服务ID的路径匹配,空目标相当于广播,简单服务目标可匹配该服务全部实例,具体Bus ID可定向实例。ID必须足够唯一,并注意冒号等路径字符的URL转义。

常见追问:Bus ID冲突会怎样?

原理入口:Bus目标匹配

39. 为什么Bus实例不能共享普通消费组

标准回答:控制事件要让每个在线实例都收到;如果所有实例共享同一group,一条事件只会由其中一个实例处理。Bus默认依赖匿名非共享订阅形成fanout,但持久性要按具体Binder判断。

常见追问:这样会带来什么可靠性问题?

原理入口:Bus广播与匿名订阅

40. 离线实例会不会错过Bus刷新

标准回答:可能。Bus是事件广播,不是持久配置数据库;匿名订阅的离线保留取决于Binder/Broker。实例启动时必须从Config Server、Nacos等权威源读取最新状态,Bus只负责加速在线实例感知变化。

常见追问:如何判断集群最终收敛?

原理入口:离线失败窗口

41. Bus ACK能证明配置已生效吗

标准回答:不能完全证明。ACK说明实例处理了Bus事件链的一部分,不代表所有Bean、连接池和第三方组件都使用了新值,也不代表业务探针成功。应结合配置版本、实例清单、ACK/trace和业务验证。

常见追问:ACK数量少于实例数怎么查?

原理入口:ACK与trace边界Bus Runbook

42. 为什么Bus不能代替配置中心

标准回答:Bus传播变化事件,不保存权威配置状态;广播可能因离线、网络和订阅窗口被错过。配置中心提供可重复读取的版本化状态,实例重启时靠它收敛,二者是“状态源”和“通知通道”的关系。

常见追问:源实例刷新成功、其他实例失败怎么办?

原理入口:Bus职责失败窗口

43. Bus端点有哪些安全风险

标准回答:它属于高权限控制面,未保护会导致全局刷新、环境篡改和大规模拒绝服务。应限制在管理网,启用强认证、最小权限、操作审计、Broker ACL、键白名单和变更限流,敏感值不能经Bus明文传播。

常见追问:为什么全量刷新也会造成事故?

原理入口:Bus安全治理

Demo入口:Java 17与Boot 3双实例刷新

44. 部分实例不刷新怎样排查

标准回答:保存事件ID、目标、配置版本和预期实例清单;依次验证源实例本地刷新、输出Binding、Broker订阅、Bus ID与目标匹配、目标实例接收和ACK、刷新异常及实际业务配置版本。最后定向刷新或从权威源重载,不能只反复全局广播。

常见追问:为什么ACK有了业务仍使用旧值?

原理入口:Bus部分刷新Runbook

七、版本与迁移

45. JDK 8 / Boot 2.7与Java 17 / Boot 3如何区分

标准回答:Boot 3最低Java 17并迁移到jakarta.*,应配套相应Spring Cloud发布列车和Stream 4.x;Boot 2.7仍可运行在JDK 8并常配Cloud 2021.0.x、Stream 3.2.x。现代代码可用Record,但跨版本事件契约仍应保持兼容,旧系统不能直接复制Java 17语法。

常见追问:只升级Boot版本为什么会失败?

原理入口:双版本Demo与差异JDK 8存量线

46. Stream旧注解模型怎样迁移

标准回答:先固定消息契约和集成测试,再把输入输出迁移为SupplierFunctionConsumer及显式Binding配置,验证group、分区、重试和DLQ行为;之后再升级Java、Boot和Cloud平台。不要一次同时更换框架、Broker和序列化协议。

常见追问:滚动发布期间如何兼容?

原理入口:存量线迁移边界

八、场景题回答模板

47. 订单事件偶发重复扣库存,怎么回答

标准回答:先按eventId确认重复来自生产重发、框架重试、Broker重投还是人工回放;检查业务提交与ACK/offset的先后窗口。修复方案是在库存数据库中以consumer group和eventId建立唯一记录,并与扣减同事务提交,同时保留唯一业务约束和对账补偿。

原理入口:消费幂等失败窗口

48. 消费扩容后数据库被打挂,怎么回答

标准回答:消费吞吐受最慢下游约束。先降低并发或暂停消费保护数据库,再检查SQL、连接池等待、锁、批量写和重试风暴;按数据库可承受QPS反推消费者并发,限速恢复并持续观察Lag斜率,不能继续无脑扩实例。

原理入口:积压Runbook

49. 配置刷新后只有一半实例生效,怎么回答

标准回答:把它视为配置版本不一致事件,按预期实例清单核对Bus ID、目标匹配、订阅在线状态、ACK和每实例实际配置版本;离线实例从权威配置源重载,刷新作用域不生效的Bean需重建或滚动重启。Bus广播本身不能保证状态收敛。

原理入口:Bus失败窗口Bus Runbook

50. 怎样证明一条消息端到端成功

标准回答:单一发送日志无法证明。至少关联生产业务事务、eventId、Broker确认、消费group与Broker位置、消费幂等记录、业务数据库结果和必要的下游事件;对关键链路再做对账。可观测性提供证据,幂等和补偿提供恢复能力。

原理入口:可观测性失败窗口

九、继续学习