Chapter 06
自测题库与面试准备
前五章建起了一整条认知链:Streams 是嵌进应用的库(01),状态住在本地 RocksDB 但真相源是 changelog 日志(02),DSL 把这套机制跑成可运行的代码(03),违反机制时它如何崩坏(04),综合实战里如何在多个机制间判别取舍(05)。本章不再讲新东西——它把那条链反过来用:合上教程,靠提取而非重读检验哪些环节真正进了长期记忆。题目按概念→原理→判别三层递进,并为最易在面试里露馅的几道题列出"普通答案 vs 资深必须点到"。能在没有提示的情况下重建这些答案,才算把 Streams 框成"消费组 + 本地状态 + changelog 日志"而非"又一个流处理引擎"。
本章你将检验的 schema
- 三层提取梯度:概念层(回忆事实)→ 原理层(解释机制)→ 判别层(迁移到新场景)
- 面试主战场是判别层——面试官在意对取舍的推理,而非背 API
- 五六道最易露馅题的"资深必须点到"清单:状态容错、EOS 真实边界、共分区、线程模型并行上限、stream-time 冻结
总题数 20 道,分三层:概念层 7 道(对应 01)、原理层 8 道(对应 02)、应用判别层 5 道(综合,面试主战场)。所有答案集中在文末一个 <details> 里,按三层分节。题目区只有题——先把能想到的答案写在纸上或编辑器里,写完再展开对照;直接点开等于把这套机制当再读一遍,提取练习的效果归零。每道概念层与原理层题都带一条 提示链,回指对应章节的锚点,卡住时去那里复习,而不是直接看答案。
6.A概念层(对应 01 概念)
每题只考一个原子事实,对应 01 章某个概念。答得卡壳,就回提示链指的锚点复习——这一层是后面两层的地基,地基空了上面的推理也立不住。
- 用一句话各自概括 KStream、KTable、GlobalKTable 的语义差异(分别是什么"日志/视图",同一个 key 来多条时各自怎么处理)。
提示:参考 01 章 §1.2 KStream / KTable / GlobalKTable - 什么是流表对偶?为什么说 table 是 stream 的物化、stream 是 table 的差异日志?这个双向性和"状态可恢复"有什么关系?
提示:参考 01 章 §1.3 流表对偶 - 一个 Streams 应用的 task 数量由谁决定?是配置项
num.stream.threads吗?运行时能改吗?
提示:参考 01 章 §1.4 应用即库 - Streams 应用的并行上限是多少?把实例数或线程数加到超过这个上限会发生什么?
提示:参考 01 章 §1.4 应用即库 - 已经有 KafkaConsumer + KafkaProducer 了,为什么还要用 Streams?它在裸 consumer 之上多给了哪几样东西?
提示:参考 01 章 §1.1 为什么是 Streams - 下列算子哪些是 stateless、哪些是 stateful:
map、filter、flatMap、branch、aggregate、count、reduce、join、窗口算子?判断依据是什么?
提示:参考 01 章 §1.5 无状态 vs 有状态 - 什么是 topology(拓扑)?什么情况下一个拓扑会被切成多个 sub-topology(子拓扑)?
提示:参考 01 章 §1.4 拓扑
6.B原理层(对应 02 原理)
这一层不问"是什么",问"怎么实现的、代价在哪、什么时候失效"。每题对应 02 章某条机制。答案要能说出底层路径,而不是复述定义。
- 有状态算子的状态怎么容错?一个有状态 task 故障转移到新实例后,新 owner 手上没有本地 RocksDB 副本,它如何恢复到正确状态?说清"真相源"在哪。
提示:参考 02 章 §2.1 状态存储 - 为什么
selectKey/map/groupBy这类改 key 的操作会触发 repartition?Streams 为此自动做了什么?这条往返的代价是什么?
提示:参考 02 章 §2.3 repartition - stream-time 何时前进?它由什么驱动、用来做什么?一个低流量或空闲的分区上,stream-time 会怎样,进而对窗口结果有什么后果?
提示:参考 02 章 §2.4 窗口与时间 - 辨析 grace period、retention、suppress 三者:各自管什么?过了
windowEnd + grace的迟到事件会怎样?为什么 retention 必须 ≥ 窗口 + grace?
提示:参考 02 章 §2.4 窗口与时间 - 说出 四种 join(KStream-KStream、KStream-KTable、KTable-KTable、外键 join)各自的触发语义;其中哪些要求共分区、哪些不要求?共分区到底要求两侧满足什么?
提示:参考 02 章 §2.5 join 与共分区 exactly_once_v2到底把哪三件事绑进一个 Kafka 事务?相比 v1 的"一 task 一 producer",v2 改了什么?为什么 v2 在高分区数下不会爆?
提示:参考 02 章 §2.6 EOS v2- RocksDB 的内存为什么不受
-Xmx管?它分配在哪里?这件事在容器里会以什么形式炸出来?为什么 Streams 选 RocksDB+changelog 而不是一个远程状态库?
提示:参考 02 章 §2.1 状态存储;失败形态见 04 章 §4.C 堆外内存 - standby 副本(
num.standby.replicas)做什么?它如何把故障转移从"分钟级回放"压到"秒级"?代价是什么?
提示:参考 02 章 §2.2 再平衡与缩放
6.C应用判别层(综合 · 面试主战场)
这一层给场景、要决策。没有"标准答案"对应某一章某一节——要把前两层的机制迁移到一个没见过的取舍里。面试官在这一层听的是推理过程:你凭什么排除另一个选项。
- Streams vs Flink / Spark Streaming:团队要做实时处理,有人主张上 Flink。从"部署形态"和"处理模型"两个维度说清 Streams 与 Flink/Spark 的根本差异;什么场景下 Streams 是更省心的选择,什么场景下它明显不够用?
- Streams vs 裸 consumer:一个服务只是"从 topic A 读、做个无状态转换、写 topic B"。该上 Streams 还是直接用 KafkaConsumer/Producer?把判据说成一条可操作的线——出现什么需求时才值得引入 Streams?
- 何时不用 Streams:举出三类即使数据在 Kafka 里、也不该用 Streams 的场景,并说明每一类的根本障碍在哪。
- groupByKey vs groupBy:两个聚合,一个用
groupByKey(),一个用groupBy((k,v)->...)。它们在拓扑和运行成本上有什么不同?默认该优先用哪个、为什么?什么情况下不得不用另一个? - KTable-KTable join vs GlobalKTable join:要用一张维表去丰富主流。维表是"小而慢变"。选普通 KTable join 还是 GlobalKTable join?从共分区要求、每实例内存、数据规模三个角度说清取舍;维表如果变大到几十 GB,结论会怎么变?
合上教程,凭记忆画出有状态聚合一条记录的完整路径:记录进入 → 更新本地 RocksDB → 写 changelog topic →(故障时)从 changelog 或 standby 副本恢复。只画这四个环节加它们之间的箭头就行。画完回 02 章 §2.1 对照——你的图里,那个箭头标对了吗:你标出真相源是 changelog 而不是本地 RocksDB 了吗?如果你的箭头是"恢复时从本地 RocksDB 读",说明对偶那一层还没真进去,回 01 章 §1.3 重看。
6.D面试加餐 · 强答案必须包含什么
下面六道是最容易"答得像对、其实露馅"的题。面试官靠它们区分"用过 Streams"和"懂 Streams"。每道列出普通答案(不算错,但谁都能背)与资深必须点到(点到才证明你摸过它的边界)。这一节不放在文末 <details> 里——它本身就是"答案该长什么样"的标尺,对照上面三层的题用。
| 题 | 普通答案(谁都能背) | 资深必须点到(摸过边界才说得出) |
|---|---|---|
| 状态容错 | 状态存在 RocksDB,会持久化,故障了能恢复。 | 本地 RocksDB 只是读缓存式的物化视图,真相源是 log-compacted 的 changelog topic;故障转移时新 owner 回放 changelog 重建 RocksDB(或用 standby 副本已追平的副本秒级接管)。措辞上要说"changelog 复制",不能说"checkpoint / 快照"——那是 Flink/Spark 的模型,说错就暴露没分清两套体系。 |
| EOS 真实边界 | exactly_once_v2 保证精确一次,不重不丢。 |
事务只绑定三件事:offset 提交 + 状态 changelog 写 + 输出 produce,下游必须 read_committed 才看不到 aborted 数据;只在单个 Kafka 集群内原子,不延伸到外部副作用(DB / HTTP / 短信仍可能重复,要 outbox 或按 event-id 幂等)。还要点出"变慢"的真因是 commit.interval.ms 在 EOS 下默认从 30000ms 降到 100ms,而非事务本身重。 |
| 共分区 | join 两边要按同一个 key。 | 等值 join 要两侧相同分区数 + 相同分区方式(同一 producer 分区器);分区数不符启动时拓扑校验失败,分区方式不符则静默零输出、不报错(最坑的失败形态)。逃生口:GlobalKTable join 和外键 join 免共分区。能把"静默无输出"这个症状讲出来,才证明真踩过。 |
| 线程模型 / 并行上限 | 调 num.stream.threads 加线程、加实例就能扩。 |
并行天花板 = 子拓扑最大输入分区数,运行时钉死不可改;task = 分区,一个 task 是一个分区组、不跨线程拆;线程/实例多于 task 则多出的闲置空转。要再扩只能重分区输入 topic(破坏性)。层级关系要说全:instance ⊃ StreamThread ⊃ task ⊃ partition。 |
| stream-time 冻结 | 窗口按事件时间关闭。 | stream-time = 每 task 见过的最大时间戳,只前进、只在有记录到达时推进;没有新记录 → 时间冻结 → 窗口永不关、结果永不发。所以低流量/空闲分区会让窗口结果"卡住不出来"——这是个真实的生产现象,不是理论。点到这一条,面试官就知道你跑过它而不是只读过文档。 |
| RocksDB 内存模型 | RocksDB 是嵌入式状态存储,落盘。 | block cache 与 memtable 分配在 JVM 堆外,-Xmx 管不住;每分区每 store 各占一份,分区一多堆外内存线性涨。容器里的表现是 JVM 堆指标全程正常、Pod 却被内核 OOMKilled。修复是用一个共享 LRUCache+WriteBufferManager 设上限,并把 Pod 内存算成"堆 + 堆外缓存 + 余量"。 |
对比上表两列,资深列的共同点不是"知道更多 API",而是每条都点到了一个边界或失败形态:changelog 才是真相源、EOS 不覆盖外部、共分区不符会静默无输出、并行被分区数钉死、空闲分区冻结 stream-time、堆外内存绕过 -Xmx。面试官在意对取舍的推理胜过背 API,想听的是失败模式的故事(恢复卡顿、空闲分区、EOS 边界)。把每个机制都连到"它在什么时候、以什么方式坏掉",答案就从"用过"升到"懂"。
把五道判别题压成一棵决策树
不看教程,把 6.C 的五道判别题合并成一张选型决策树:根节点是"我要做流处理",叶子是 {裸 consumer、Kafka Streams、Flink/Spark、Streams + GlobalKTable、Streams + KTable join}。每个分叉用一个二元判据(如"需要状态/窗口/join 吗?""源汇主要是 Kafka 吗?""维表能塞进每个实例内存吗?")。画完检查:有没有哪条判据其实在重复另一条?有没有哪个叶子永远走不到?
提示(卡住再展开)
先按"源汇是否主要在 Kafka"分一刀——否则连 Streams 的门槛都够不到,落向 Flink/Spark 或别的方案。再按"是否需要状态/窗口/join/EOS"分第二刀——只读简单消费的落向裸 consumer。进了 Streams 之后才轮到"维表多大"决定 GlobalKTable(小而慢变、免共分区、每实例全量)还是 KTable join(大、要共分区)。判据顺序很重要:先排除门槛,再在门槛内部细分。
答案(三层全部 · 先做完再展开)
概念层(6.A)
- KStream = append-only 的 INSERT 日志,同一个 key 来多条记录全部保留(每条都是一个独立事件)。KTable = 按 key UPSERT 的 changelog,同 key 后到的值覆盖先到的,
null值是 tombstone(删除该 key)。GlobalKTable = 每个实例都消费该 topic 的所有分区、持有全量副本,启动时 eager 加载——因此 join 时免共分区,但只适合小而慢变的维表。 - 流表对偶 = 流和表是同一条日志的两种视图:table 是 changelog "每 key 最新值" 的物化(把日志按 key 折叠成当前状态),stream 是 table 的差异日志(把状态变化展开成事件序列)。正因为这个双向性,状态才可恢复——把 changelog 从头回放一遍就能重建出 table,所以本地状态丢了也不怕。
- task 数 = 子拓扑输入 topic 的分区数(取最大输入分区数),运行时钉死不可改。
num.stream.threads只控制每个实例起几个 StreamThread,不决定 task 总数;task 在所有实例的所有线程间分配。要改 task 数只能重分区输入 topic。 - 并行上限 = 最大输入分区数(= task 数)。实例数或线程数加到超过它,多出来的实例/线程会被分到零个 task、CPU 空转,吞吐不再增长。越过这个上限的唯一办法是重分区输入 topic(破坏性操作)。
- 裸 consumer 只给你"读到消息"。Streams 在其上加了:有状态算子(聚合/count/reduce)及其容错(changelog + standby)、窗口与事件时间处理、join(含自动 repartition 与共分区校验)、EOS v2 端到端精确一次、以及把这些架在消费组之上的缩放/再平衡。只要不需要这些,裸 consumer 就够。
- stateless:
map、filter、flatMap、branch(还有 merge/peek)——逐条处理、不依赖其它记录。stateful:aggregate、count、reduce、join、所有窗口算子——需要跨记录维护状态,因此要状态存储。判断依据:处理这一条时是否需要"记住"别的记录。 - topology = source → processor → sink 的算子有向图,Streams 把 DSL 编译成它。当流中出现需要重分区的边界(改 key 后接 stateful 算子,要经 repartition topic 重新洗牌)时,拓扑会在那里被切成多个 sub-topology——每个子拓扑独立调度成 task。
原理层(6.B)
- 每次对状态存储的 put 同时写一条内部 changelog topic(log-compacted)。本地 RocksDB 只是物化视图,真相源是 changelog。故障转移:新 owner 没有本地副本,就从 changelog 回放逐条重建 RocksDB 后再处理;配了 standby 副本则由一直尾随 changelog 的热备秒级接管,跳过大部分回放。
- 聚合和等值 join 要求相同 key 的记录落到同一个 task / 分区(共分区)。改 key 的操作打乱了原有分区,所以 Streams 在下游 stateful 算子前自动建一个内部
<app>-<name>-repartitiontopic,把数据按新 key produce 进去再重新消费回来。代价:每个改 key 的操作前置一次 produce→re-consume 往返(延迟 + broker 磁盘 + 额外 topic)。注意selectKey本身是 lazy 的,只有下游真依赖 key 的算子才物化出 topic。 - stream-time = 该 task 见过的所有记录的最大时间戳,只前进不回退,用来驱动窗口关闭与迟到判定。它只在有记录到达时推进。低流量或空闲分区上没有新记录,stream-time 就冻结,窗口永远等不到
windowEnd + grace、结果永不发——表现为"结果卡住不出来"。 - grace period = 窗口结束后迟到事件还能更新该窗口结果的时长;过
windowEnd + grace的迟到事件被静默丢弃(只计入 dropped 指标)。retention = 窗口状态/changelog 在存储里保留多久,必须 ≥ 窗口 + grace,否则窗口还在 grace 内状态就被清了。suppress(untilWindowCloses) = 抑制中间结果,每个窗口只发一个最终结果(代价:要等窗口关,增加延迟;内存压力下若用emitEarlyWhenFull会发非最终结果,要"最终"必须shutDownWhenFull)。 - KStream-KStream:窗口化 join,两侧入库缓冲,窗口+grace 内每对匹配各发一次。KStream-KTable:查表 join,流记录探当前表值;表更新不反向触发。KTable-KTable:维护式,任一侧更新都重新发结果。外键 join:按 value 提取的 key join,Streams 建两个隐藏 topic(subscription + response)+ 复合 RocksDB 存储。前三种等值 join 要求共分区(相同分区数 + 相同分区方式);GlobalKTable join 和外键 join 免共分区。
- 把 offset 提交 + 状态 changelog 写 + 输出 produce 三件事绑进一个 Kafka 事务,下游
read_committed跳过 aborted。v2(KIP-447)相比 v1 改成一个线程一个 producer 覆盖该线程的多个 task,而 v1 是一个 task 一个 producer。高分区数下 v1 会产生海量 producer(爆炸),v2 因此不爆。 - RocksDB 的 block cache(≈50MiB)和 memtable 分配在 JVM 堆外(off-heap),
-Xmx只管堆内、管不到它;而且每分区每 store 各占一份,分区多则堆外内存线性增长。容器里表现为 JVM 堆指标正常、Pod 却被内核 OOMKilled。选 RocksDB + changelog 而非远程状态库:本地读亚毫秒、无逐条网络跳,且 Kafka 的 compacted 日志本身已是复制/HA 层,不必再引第二个集群化数据库。 - standby 副本 = 在另一个实例上预先维护一份尾随 changelog 的热备状态。主 task 故障时,已追平的 standby 直接接管,跳过从零回放 changelog,把故障转移从分钟级压到秒级。代价:额外的磁盘 + 网络(每个 standby 都在持续拉 changelog 保温)。
应用判别层(6.C)
- 部署形态:Streams 是嵌进你应用的库,没有独立集群,跑在消费组上、用普通 JVM 部署;Flink/Spark 是独立集群,要单独的 JobManager/TaskManager 或 Spark 集群运维。处理模型:Streams 逐条处理、容错靠 changelog 复制;Flink 逐条但容错靠分布式快照 checkpoint;Spark Streaming 是微批。Streams 更省心的场景:源和汇都在 Kafka、想避免再养一个集群、Java 后端就近嵌入。Streams 不够用的场景:需要跨多种数据源/汇、复杂 CEP、ML pipeline、或要 Flink 那种独立 checkpoint/状态后端能力。
- 无状态单纯转发不该上 Streams——裸 consumer/producer 就够,还少一层抽象和内部 topic 开销。值得引入 Streams 的那条线:出现有状态需求(聚合/count/reduce)、窗口/事件时间、join、或端到端 EOS 中任意一项时。判据是"是否需要跨记录的状态或时间语义",不是数据量大小。
- 三类即使数据在 Kafka 也不该用 Streams:① 源汇主要不在 Kafka(比如主要从 DB/HTTP 拉、写第三方系统)——Streams 的所有机制都围绕 Kafka topic,离开 Kafka 它的优势全失;② 需要复杂 CEP / ML / 跨多集群的能力——Streams 的算子模型撑不起,该上 Flink;③ 只是简单消费没有状态/窗口/join 需求——上 Streams 是过度设计,徒增内部 topic 与运维面。
groupByKey()不改 key,下游聚合不需要 repartition,没有额外内部 topic。groupBy((k,v)->...)改 key,下游聚合会自动建 repartition topic 把数据重洗一遍——多一次 produce→re-consume 往返(延迟 + broker 磁盘)。默认优先groupByKey。只有当确实需要按一个不同于当前 key 的字段聚合时,才不得不用groupBy(这时重分区是必要代价,不是浪费)。- "小而慢变"的维表选 GlobalKTable join:免共分区(主流不必按维表 key 重分区)、查表式语义简单;代价是每个实例都持有全量副本(磁盘 + 内存 + 启动 bootstrap 延迟)。普通 KTable-KTable join 要求共分区、状态按分区切分(每实例只存自己那份),适合维表较大、放不进单实例内存的情况。维表涨到几十 GB 时结论翻转:GlobalKTable 的"每实例全量"变得不可承受(内存爆 + 启动慢),应改用按 key 共分区的 KTable join,或用外键 join(按 value 提取 key、免共分区、状态分区化)。