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> 里,按三层分节。题目区只有题——先把能想到的答案写在纸上或编辑器里,写完再展开对照;直接点开等于把这套机制当再读一遍,提取练习的效果归零。每道概念层与原理层题都带一条 提示链,回指对应章节的锚点,卡住时去那里复习,而不是直接看答案。

认知要求 ↑ 概念层 · 回忆 7 题 · 对应 01 概念 原理层 · 理解 8 题 · 对应 02 原理 判别层 迁移 · 5 题 面试主战场 能复述 能解释为什么 能迁移判别
图 6.0三层提取梯度:底层最宽(回忆的事实最多、最易答),越往上认知要求越高、题量越少。注意:朱红的顶层"判别层"题量最少却是面试主战场——它要的不是背 API,而是把下两层的机制迁移到一个没见过的取舍场景里。底层答得顺不代表顶层答得出,三层一起做才看得出 schema 哪里是空的。

6.A概念层(对应 01 概念)

每题只考一个原子事实,对应 01 章某个概念。答得卡壳,就回提示链指的锚点复习——这一层是后面两层的地基,地基空了上面的推理也立不住。

  1. 用一句话各自概括 KStream、KTable、GlobalKTable 的语义差异(分别是什么"日志/视图",同一个 key 来多条时各自怎么处理)。
    提示:参考 01 章 §1.2 KStream / KTable / GlobalKTable
  2. 什么是流表对偶?为什么说 table 是 stream 的物化、stream 是 table 的差异日志?这个双向性和"状态可恢复"有什么关系?
    提示:参考 01 章 §1.3 流表对偶
  3. 一个 Streams 应用的 task 数量由谁决定?是配置项 num.stream.threads 吗?运行时能改吗?
    提示:参考 01 章 §1.4 应用即库
  4. Streams 应用的并行上限是多少?把实例数或线程数加到超过这个上限会发生什么?
    提示:参考 01 章 §1.4 应用即库
  5. 已经有 KafkaConsumer + KafkaProducer 了,为什么还要用 Streams?它在裸 consumer 之上多给了哪几样东西?
    提示:参考 01 章 §1.1 为什么是 Streams
  6. 下列算子哪些是 stateless、哪些是 stateful:map、filter、flatMap、branch、aggregate、count、reduce、join、窗口算子?判断依据是什么?
    提示:参考 01 章 §1.5 无状态 vs 有状态
  7. 什么是 topology(拓扑)?什么情况下一个拓扑会被切成多个 sub-topology(子拓扑)?
    提示:参考 01 章 §1.4 拓扑

6.B原理层(对应 02 原理)

这一层不问"是什么",问"怎么实现的、代价在哪、什么时候失效"。每题对应 02 章某条机制。答案要能说出底层路径,而不是复述定义。

  1. 有状态算子的状态怎么容错?一个有状态 task 故障转移到新实例后,新 owner 手上没有本地 RocksDB 副本,它如何恢复到正确状态?说清"真相源"在哪。
    提示:参考 02 章 §2.1 状态存储
  2. 为什么 selectKey / map / groupBy 这类改 key 的操作会触发 repartition?Streams 为此自动做了什么?这条往返的代价是什么?
    提示:参考 02 章 §2.3 repartition
  3. stream-time 何时前进?它由什么驱动、用来做什么?一个低流量或空闲的分区上,stream-time 会怎样,进而对窗口结果有什么后果?
    提示:参考 02 章 §2.4 窗口与时间
  4. 辨析 grace period、retention、suppress 三者:各自管什么?过了 windowEnd + grace 的迟到事件会怎样?为什么 retention 必须 ≥ 窗口 + grace?
    提示:参考 02 章 §2.4 窗口与时间
  5. 说出 四种 join(KStream-KStream、KStream-KTable、KTable-KTable、外键 join)各自的触发语义;其中哪些要求共分区、哪些不要求?共分区到底要求两侧满足什么?
    提示:参考 02 章 §2.5 join 与共分区
  6. exactly_once_v2 到底把哪三件事绑进一个 Kafka 事务?相比 v1 的"一 task 一 producer",v2 改了什么?为什么 v2 在高分区数下不会爆?
    提示:参考 02 章 §2.6 EOS v2
  7. RocksDB 的内存为什么不受 -Xmx 管?它分配在哪里?这件事在容器里会以什么形式炸出来?为什么 Streams 选 RocksDB+changelog 而不是一个远程状态库?
    提示:参考 02 章 §2.1 状态存储;失败形态见 04 章 §4.C 堆外内存
  8. standby 副本(num.standby.replicas)做什么?它如何把故障转移从"分钟级回放"压到"秒级"?代价是什么?
    提示:参考 02 章 §2.2 再平衡与缩放

6.C应用判别层(综合 · 面试主战场)

这一层给场景、要决策。没有"标准答案"对应某一章某一节——要把前两层的机制迁移到一个没见过的取舍里。面试官在这一层听的是推理过程:你凭什么排除另一个选项。

  1. Streams vs Flink / Spark Streaming:团队要做实时处理,有人主张上 Flink。从"部署形态"和"处理模型"两个维度说清 Streams 与 Flink/Spark 的根本差异;什么场景下 Streams 是更省心的选择,什么场景下它明显不够用?
  2. Streams vs 裸 consumer:一个服务只是"从 topic A 读、做个无状态转换、写 topic B"。该上 Streams 还是直接用 KafkaConsumer/Producer?把判据说成一条可操作的线——出现什么需求时才值得引入 Streams?
  3. 何时不用 Streams:举出三类即使数据在 Kafka 里、也不该用 Streams 的场景,并说明每一类的根本障碍在哪。
  4. groupByKey vs groupBy:两个聚合,一个用 groupByKey(),一个用 groupBy((k,v)->...)。它们在拓扑和运行成本上有什么不同?默认该优先用哪个、为什么?什么情况下不得不用另一个?
  5. 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> 里——它本身就是"答案该长什么样"的标尺,对照上面三层的题用。

表 6.1 · 六道易露馅题:普通答案 vs 资深必须点到
题普通答案(谁都能背)资深必须点到(摸过边界才说得出)
状态容错 状态存在 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)

  1. KStream = append-only 的 INSERT 日志,同一个 key 来多条记录全部保留(每条都是一个独立事件)。KTable = 按 key UPSERT 的 changelog,同 key 后到的值覆盖先到的,null 值是 tombstone(删除该 key)。GlobalKTable = 每个实例都消费该 topic 的所有分区、持有全量副本,启动时 eager 加载——因此 join 时免共分区,但只适合小而慢变的维表。
  2. 流表对偶 = 流和表是同一条日志的两种视图:table 是 changelog "每 key 最新值" 的物化(把日志按 key 折叠成当前状态),stream 是 table 的差异日志(把状态变化展开成事件序列)。正因为这个双向性,状态才可恢复——把 changelog 从头回放一遍就能重建出 table,所以本地状态丢了也不怕。
  3. task 数 = 子拓扑输入 topic 的分区数(取最大输入分区数),运行时钉死不可改。num.stream.threads 只控制每个实例起几个 StreamThread,不决定 task 总数;task 在所有实例的所有线程间分配。要改 task 数只能重分区输入 topic。
  4. 并行上限 = 最大输入分区数(= task 数)。实例数或线程数加到超过它,多出来的实例/线程会被分到零个 task、CPU 空转,吞吐不再增长。越过这个上限的唯一办法是重分区输入 topic(破坏性操作)。
  5. 裸 consumer 只给你"读到消息"。Streams 在其上加了:有状态算子(聚合/count/reduce)及其容错(changelog + standby)、窗口与事件时间处理、join(含自动 repartition 与共分区校验)、EOS v2 端到端精确一次、以及把这些架在消费组之上的缩放/再平衡。只要不需要这些,裸 consumer 就够。
  6. stateless:map、filter、flatMap、branch(还有 merge/peek)——逐条处理、不依赖其它记录。stateful:aggregate、count、reduce、join、所有窗口算子——需要跨记录维护状态,因此要状态存储。判断依据:处理这一条时是否需要"记住"别的记录。
  7. topology = source → processor → sink 的算子有向图,Streams 把 DSL 编译成它。当流中出现需要重分区的边界(改 key 后接 stateful 算子,要经 repartition topic 重新洗牌)时,拓扑会在那里被切成多个 sub-topology——每个子拓扑独立调度成 task。

原理层(6.B)

  1. 每次对状态存储的 put 同时写一条内部 changelog topic(log-compacted)。本地 RocksDB 只是物化视图,真相源是 changelog。故障转移:新 owner 没有本地副本,就从 changelog 回放逐条重建 RocksDB 后再处理;配了 standby 副本则由一直尾随 changelog 的热备秒级接管,跳过大部分回放。
  2. 聚合和等值 join 要求相同 key 的记录落到同一个 task / 分区(共分区)。改 key 的操作打乱了原有分区,所以 Streams 在下游 stateful 算子前自动建一个内部 <app>-<name>-repartition topic,把数据按新 key produce 进去再重新消费回来。代价:每个改 key 的操作前置一次 produce→re-consume 往返(延迟 + broker 磁盘 + 额外 topic)。注意 selectKey 本身是 lazy 的,只有下游真依赖 key 的算子才物化出 topic。
  3. stream-time = 该 task 见过的所有记录的最大时间戳,只前进不回退,用来驱动窗口关闭与迟到判定。它只在有记录到达时推进。低流量或空闲分区上没有新记录,stream-time 就冻结,窗口永远等不到 windowEnd + grace、结果永不发——表现为"结果卡住不出来"。
  4. grace period = 窗口结束后迟到事件还能更新该窗口结果的时长;过 windowEnd + grace 的迟到事件被静默丢弃(只计入 dropped 指标)。retention = 窗口状态/changelog 在存储里保留多久,必须 ≥ 窗口 + grace,否则窗口还在 grace 内状态就被清了。suppress(untilWindowCloses) = 抑制中间结果,每个窗口只发一个最终结果(代价:要等窗口关,增加延迟;内存压力下若用 emitEarlyWhenFull 会发非最终结果,要"最终"必须 shutDownWhenFull)。
  5. KStream-KStream:窗口化 join,两侧入库缓冲,窗口+grace 内每对匹配各发一次。KStream-KTable:查表 join,流记录探当前表值;表更新不反向触发。KTable-KTable:维护式,任一侧更新都重新发结果。外键 join:按 value 提取的 key join,Streams 建两个隐藏 topic(subscription + response)+ 复合 RocksDB 存储。前三种等值 join 要求共分区(相同分区数 + 相同分区方式);GlobalKTable join 和外键 join 免共分区。
  6. 把 offset 提交 + 状态 changelog 写 + 输出 produce 三件事绑进一个 Kafka 事务,下游 read_committed 跳过 aborted。v2(KIP-447)相比 v1 改成一个线程一个 producer 覆盖该线程的多个 task,而 v1 是一个 task 一个 producer。高分区数下 v1 会产生海量 producer(爆炸),v2 因此不爆。
  7. RocksDB 的 block cache(≈50MiB)和 memtable 分配在 JVM 堆外(off-heap),-Xmx 只管堆内、管不到它;而且每分区每 store 各占一份,分区多则堆外内存线性增长。容器里表现为 JVM 堆指标正常、Pod 却被内核 OOMKilled。选 RocksDB + changelog 而非远程状态库:本地读亚毫秒、无逐条网络跳,且 Kafka 的 compacted 日志本身已是复制/HA 层,不必再引第二个集群化数据库。
  8. standby 副本 = 在另一个实例上预先维护一份尾随 changelog 的热备状态。主 task 故障时,已追平的 standby 直接接管,跳过从零回放 changelog,把故障转移从分钟级压到秒级。代价:额外的磁盘 + 网络(每个 standby 都在持续拉 changelog 保温)。

应用判别层(6.C)

  1. 部署形态:Streams 是嵌进你应用的库,没有独立集群,跑在消费组上、用普通 JVM 部署;Flink/Spark 是独立集群,要单独的 JobManager/TaskManager 或 Spark 集群运维。处理模型:Streams 逐条处理、容错靠 changelog 复制;Flink 逐条但容错靠分布式快照 checkpoint;Spark Streaming 是微批。Streams 更省心的场景:源和汇都在 Kafka、想避免再养一个集群、Java 后端就近嵌入。Streams 不够用的场景:需要跨多种数据源/汇、复杂 CEP、ML pipeline、或要 Flink 那种独立 checkpoint/状态后端能力。
  2. 无状态单纯转发不该上 Streams——裸 consumer/producer 就够,还少一层抽象和内部 topic 开销。值得引入 Streams 的那条线:出现有状态需求(聚合/count/reduce)、窗口/事件时间、join、或端到端 EOS 中任意一项时。判据是"是否需要跨记录的状态或时间语义",不是数据量大小。
  3. 三类即使数据在 Kafka 也不该用 Streams:① 源汇主要不在 Kafka(比如主要从 DB/HTTP 拉、写第三方系统)——Streams 的所有机制都围绕 Kafka topic,离开 Kafka 它的优势全失;② 需要复杂 CEP / ML / 跨多集群的能力——Streams 的算子模型撑不起,该上 Flink;③ 只是简单消费没有状态/窗口/join 需求——上 Streams 是过度设计,徒增内部 topic 与运维面。
  4. groupByKey() 不改 key,下游聚合不需要 repartition,没有额外内部 topic。groupBy((k,v)->...) 改 key,下游聚合会自动建 repartition topic 把数据重洗一遍——多一次 produce→re-consume 往返(延迟 + broker 磁盘)。默认优先 groupByKey。只有当确实需要按一个不同于当前 key 的字段聚合时,才不得不用 groupBy(这时重分区是必要代价,不是浪费)。
  5. "小而慢变"的维表选 GlobalKTable join:免共分区(主流不必按维表 key 重分区)、查表式语义简单;代价是每个实例都持有全量副本(磁盘 + 内存 + 启动 bootstrap 延迟)。普通 KTable-KTable join 要求共分区、状态按分区切分(每实例只存自己那份),适合维表较大、放不进单实例内存的情况。维表涨到几十 GB 时结论翻转:GlobalKTable 的"每实例全量"变得不可承受(内存爆 + 启动慢),应改用按 key 共分区的 KTable join,或用外键 join(按 value 提取 key、免共分区、状态分区化)。