Chapter 03

自测题库:从回忆到判别

01 章给了词汇表(connector / worker / task / converter / SMT / offset / DLQ / EOS source),02 章给了机制与取舍(控制面在 Kafka topic、incremental cooperative 再平衡、序列化 wire format、偏移恢复、EOS 事务 + fencing)。读得顺不等于答得出。这一章把那两章的内容翻过来——不再喂结论,而是逼着读者从空白里把结论提取出来,并在真实选型场景里做判别。

这一章怎么用

  • 三层梯度共 16 道题:概念层 6 道(对应 01)、原理层 6 道(对应 02)、应用判别层 4 道(跨 01 + 02 的选型场景)。
  • 每道题先合上教程、把答案写在纸上或编辑器里,写完再翻文末答案对照——直接看答案等于把这两章当再读一遍,提取的训练效果归零。
  • 所有答案集中在文末一个折叠块里,按三层分组。题目区只有题,没有夹带答案。
  • 应用判别层是这套 concept-focused 教程的迁移训练主战场——它检验的不是"记没记住",而是"换个真实场景还选得对吗"。
概念层 · 回忆 原子定义 · 对应 01 章 · 6 道 原理层 · 理解 机制为何如此 · 对应 02 章 · 6 道 判别 · 迁移 真实选型 · 跨 01+02 · 4 道 最难 认得出 讲得清 选得对 认知要求 自下而上
图 3.0三层题库不是难度随机堆叠,而是一座金字塔:底层"认得出"概念、中层"讲得清"机制、顶层"选得对"场景,认知要求自下而上递增。 注意:塔尖最窄、题量最少(4 道判别),却是这套 concept-focused 教程唯一训练迁移的地方——底层答得再溜,塔尖答不出就等于没真正理解取舍。

A概念层 · 对应 01 章(6 道)

每题只考一个原子概念——一句话就该答得出。答不上来的,顺着提示链回 01 章对应小节补。

  1. source connector 与 sink connector 的数据流方向各是什么?一句话说清。 提示:参考 01 章 §1.1 connector
  2. standalone 与 distributed 两种 worker 模式,关键差异是什么?各自适合什么场景? 提示:参考 01 章 §1.2 worker
  3. task 和 connector 是什么关系?为什么说"并行单位是 task 不是 connector"? 提示:参考 01 章 §1.3 task
  4. converter 到底干什么?为什么序列化格式由它决定、而不是由 connector 决定? 提示:参考 01 章 §1.4 converter
  5. SMT 能不能按 key 关联另一个 topic 做 join 或聚合?为什么?这类需求该交给谁? 提示:参考 01 章 §1.5 SMT
  6. source connector 的偏移存在哪里、按什么索引?这与 sink connector 的偏移有什么不同? 提示:参考 01 章 §1.6 offset

B原理层 · 对应 02 章(6 道)

这一层不问"是什么",问"为什么这么设计、代价是什么"。答案要触及机制,而不是复述定义。

  1. distributed 模式的三个 internal topic(connect-configs / connect-offsets / connect-status)为什么都必须是 compacted?如果不压缩会怎样? 提示:参考 02 章 §2.1 控制面协调
  2. incremental cooperative 再平衡(KIP-415)相比早期的 eager stop-the-world 再平衡,到底解决了什么问题?代价是什么? 提示:参考 02 章 §2.2 再平衡
  3. org.apache.kafka.connect.json.JsonConverter 与 io.confluent.connect.json.JsonSchemaConverter 名字相近,为什么说它们的 wire format 完全不同、不可互换?字节层差在哪? 提示:参考 02 章 §2.3 序列化机制
  4. 给一个正在运行的 source connector 改名,为什么会导致它从头重灌整个外部源?从偏移恢复机制讲清因果。 提示:参考 02 章 §2.4 偏移恢复
  5. EOS source(KIP-618)靠哪两个 Kafka 原语做到精确一次?事务保证了什么、fencing 又拦住了什么? 提示:参考 02 章 §2.5 EOS 机制
  6. distributed worker 的 group.id 配错(一半 worker 用 A、一半用 B)会发生什么?为什么这是配置一致性的红线? 提示:参考 02 章 §2.1 控制面协调

C应用判别层 · 跨 01 + 02 场景(4 道,capstone 替代)

这一层把概念放进真实选型里。每题给一个场景,回答选哪个、为什么——理由要落到前两章的具体机制上,而不是"感觉这个更好"。这四道题承担整套 concept-focused 教程的判别训练职责,是塔尖。

  1. Connect 还是自写 producer?要把一张 Postgres 表持续同步进 Kafka。一名工程师提议直接写个 producer 程序轮询表、把行发进 topic。该用 Kafka Connect 还是自写 producer/consumer?给出理由——Connect 在这件事上白送了哪些自写时必须自己实现的东西?
  2. SMT 还是 Kafka Streams?需求是:把订单 topic 按订单 key 关联用户 topic 的用户信息,再按地区聚合下单金额,结果写进一个新 topic。这该用 connector 的 SMT 链实现,还是另起一个 Kafka Streams 应用?为什么?
  3. SMT 还是 Streams(轻量版)?需求变了:只是在数据摄入 Kafka 时,把每条记录里的手机号字段脱敏、并给字段改个名,不做任何跨记录关联。这次该用 SMT 还是 Streams?和上一题的判断为什么相反?
  4. standalone 还是 distributed?两个场景:(a) 开发机上临时把一个本地日志文件搬进 Kafka 做联调;(b) 生产环境跑 20 个 connector,要求某台机器宕机后任务自动转移、且能用 REST 在线增删 connector。各该选哪种 worker 模式?分界点是什么?

D面试加餐:这几道题的"强答案"必须包含什么

下面五道题最容易"答得像懂了、实则露馅"——普通答案停在表层、资深答案点到机制与边界。对照这张表检查自己上面写的答案够不够硬。这是面试官区分"用过"和"理解"的分水岭。

表 3.1 · 五道最易露馅题:普通答案 vs 资深答案
题目 普通答案(停在表层) 资深答案(必须点到的机制与边界)
EOS source 怎么实现的? "它能保证精确一次、不重复。" Kafka 事务把数据记录 + 偏移更新原子双写;fencing 给每代 task 递增标识、broker 拒绝僵尸旧代的写入。是 worker 级开关 exactly.once.source.support、仅 distributed、全集群一致开(两阶段 preparing→enabled),不能 per-connector(KIP-618, 3.3 GA)。
source 与 sink 的偏移有什么不同? "都是记录搬到哪了,用来重启续传。" source 是 connector 自定义格式的偏移,存进 connect-offsets topic、按 connector 名索引——所以改名 = 查不到旧偏移 = 从头重灌。sink 是普通消费组位点(组名 connect-<name>),存在 __consumer_offsets,和任何消费者一样。两套机制,不是一回事。
tasks.max 怎么映射并行度? "设成几就有几个 task 并行。" tasks.max 是上限不是保证。实际 task 数由 connector 决定能切几份,再对上限取下限:sink 受订阅分区数封顶(每 task 一个消费组成员)、JDBC source 受表数限。超出的 task 不报错只空转,GET /status 仍显示 RUNNING——排查吞吐的常见误判点。
DLQ 能抓到哪些错误? "开了 DLQ,坏记录都进死信队列,就不丢了。" DLQ 仅 sink 可用,且只抓 converter / SMT / key-value 反序列化错误——即"这条记录读不懂/变换不了"。不抓 sink 写外部系统的投递失败(ES 拒写、JDBC 主键冲突),那类走 errors.retry.* 重试、耗尽仍让 task 失败。再加 errors.tolerance 默认 none(一条坏记录停整 task)这层。
三个 internal topic 为什么 compacted? "为了持久化,不丢配置。" 它们是 key-value 状态存储,要的是每个 key 的最新值(最新配置 / 最新偏移 / 最新状态),不是完整历史——compaction 正好保留每 key 最新、回收旧值。若不压缩,topic 无限膨胀、worker 重启重放全部历史变慢。配套红线:group.id 不一致会把一个集群裂成两个。
洞察 · 面试官在听什么

面试官真正在意的不是术语背诵,而是对取舍的推理和失败模式的故事。能讲出"converter 配错导致反序列化风暴""source 改名丢偏移从头重灌""毒丸记录默认停整个 task"这类具体失败链,比复述十个定义更能证明真用过、真理解。强答案 = 机制 + 边界 + 它在什么时候会咬人。

亲手画一张图

合上教程,在纸上或 Excalidraw 里画出一条记录的完整路径:从外部源 → source connector → converter 序列化 → Kafka topic → converter 反序列化 → sink connector → 外部目标。然后在图上标出两处偏移分别存在哪——source 端的偏移、sink 端的偏移。画完回到 起点页的概念图(§概念图)对照:这两处偏移你都标对了吗?(提示:一处在 connect-offsets topic 按 connector 名索引,另一处是普通消费组位点——画错任一处,说明 §1.6 那条边界还没真正内化。)

§答案(三层都做完再展开)

下面是全部 16 道题的答案,按三层分组。先把你自己的答案写完——直接展开对照,提取练习的效果就没了。

展开全部答案(16 道,按概念层 / 原理层 / 应用判别层分组)

概念层(对应 01 章)

  1. source → Kafka,sink → 外部。source connector 把外部系统(数据库、文件、消息队列)的数据写进 Kafka;sink connector 把 Kafka 的数据写进外部系统(ES、S3、JDBC)。方向相反,是同一套插件契约的两端。
  2. standalone 是单进程、文件配置、无容错,偏移存本地文件,进程死即状态没——只适合开发期或单机日志搬运。distributed 是多 worker 共享一个 group.id 组成集群,靠 REST 管理、worker 加入/离开/故障自动再平衡,状态存 Kafka topic——生产环境用它。两者不是大小号关系,而是两套不同的存储与协调模型。
  3. connector 是"一份配置 + 把工作切成 N 份的计划",它本身不搬数据;真正读写数据的是它派生出的 task。一个 connector 最多派生 tasks.max 个 task,被均摊到各 worker 并行跑。所以扩容、估吞吐都按 task 算——并行单位是 task。
  4. converter 负责记录在 Connect 内部结构化对象与 Kafka 字节之间的相互转换(序列化/反序列化),key 和 value 各配一个。格式由它决定而非 connector,是因为 Connect 把序列化抽成了独立的一层、与 connector 解耦——于是任意 connector × 任意格式自由组合,换格式只动配置不动 connector 代码。
  5. 不能。SMT 是单条、无状态的逐条变换,一次只看一条记录、没有跨记录内存,所以做不了 join、聚合、按 key 关联另一个 topic、一条变多条。这类有状态关联是 Kafka Streams / ksqlDB 的职责。
  6. source 偏移是 connector 自定义格式,存进 connect-offsets topic,按 connector 名索引。sink 偏移是普通消费组位点(组名 connect-<connector 名>),存在 Kafka 的 __consumer_offsets 里,和任何消费者一样。两套完全不同的机制。

原理层(对应 02 章)

  1. 这三个 topic 是key-value 状态存储,需要的是每个 key 的最新值——最新的 connector 配置、最新的 source 偏移、最新的 task 状态——而不是完整的变更历史。compaction 正好保留每个 key 的最新记录、回收旧值。若不压缩:topic 会无限膨胀,且 worker 启动时要重放全部历史才能重建当前状态,恢复越来越慢。
  2. 它解决的是"一点变动就全停"的问题。eager 再平衡在任何拓扑/配置变更时停掉所有 task重新分配(stop-the-world),变更越频繁停摆越多。incremental cooperative 只暂停被收回或被移动的那部分 task,其余继续搬。代价:分配过程分多轮收敛、逻辑更复杂;且 worker 数中途变动有已知的 skew bug(KAFKA-12495)。
  3. JsonConverter 产出的是纯 JSON 字节(可选每条内嵌 {"schema":...,"payload":...}),自描述、不依赖外部注册中心。JsonSchemaConverter 走 Schema Registry,字节里带 magic byte + schema id(Confluent wire format),靠 id 去注册中心取 schema。字节布局根本不同——用一个写、用另一个读,会在反序列化阶段失败(典型 Unknown magic byte!)。名字像,不可互换。
  4. source 偏移按 connector 名索引存在 connect-offsets 里。改名后,框架拿新名字去查偏移,查到的是一组空偏移——它认为这个 connector 从未搬过任何数据,于是从外部源的起点重新全量搬一遍。机制上"改名"等价于"新建一个没有历史的 connector"。要重置应该用 REST /offsets 显式管理,而不是靠改名。
  5. 靠事务和 fencing。事务:task 把一批数据记录 + 对应的偏移更新(写进 connect-offsets)放进同一个 Kafka 事务原子提交——要么都生效要么都不生效,杜绝"记录写了、偏移没写、重启重发"的窗口。fencing(僵尸隔离):再平衡后旧 task 实例可能还活着想写(僵尸),框架给每代 task 递增的事务标识,broker 拒绝旧代的写入,保证同一份工作只有最新一代能提交。
  6. group.id 是 distributed worker 加入同一个 Connect 集群的身份。一半用 A、一半用 B,会形成两个互不相识的集群:各自读不同的 internal topic、各自以为自己拥有全部 connector,导致 connector 被重复调度或"凭空消失"。这是红线,因为 distributed 的整个协调与状态共享都建立在"所有 worker 同 group.id + 同三个 internal topic 配置"之上。

应用判别层(综合 01 + 02)

  1. 用 Connect,别自写 producer。把表同步进 Kafka 是标准化搬运,Connect 白送了自写时必须自己实现的一整套:偏移跟踪与重启续传(source 偏移自动存 connect-offsets)、并行(task 切分 + 跨 worker 分摊)、容错与再平衡(worker 宕机任务自动转移)、REST 在线管理(增删改查、暂停恢复)、以及现成的 JDBC/Debezium connector。自写就是重造一条 ETL 管道,还得自己维护位点、容错、监控。只有当需求超出"搬运"(要复杂有状态加工)时才另说。
  2. 用 Kafka Streams,不是 SMT。这需求是有状态的:按 key join 两个 topic + 跨记录聚合。SMT 单条无状态,看不到第二条记录、做不了 join / 聚合 / 1→N,强行用会卡在"拿不到要关联的数据"上。join + 聚合 + 写出新 topic 正是 Streams / ksqlDB 的核心能力。
  3. 用 SMT,不是 Streams。这次是纯逐条、无状态的编辑——脱敏一个字段(MaskField)、改字段名(ReplaceField),不跨记录。SMT 几行配置内联在搬运路径上就完成,零代码、无需独立应用和 topic。为这点事起一个 Streams 应用是杀鸡用牛刀。和上一题相反的判断点就是这一个词:有没有状态——要不要看到"别的记录"。要,则 Streams;不要,则 SMT。
  4. (a) 开发机临时搬本地文件 → standalone:一个进程、一个 properties 文件就够,不需要集群和 REST。(b) 生产 20 个 connector + 故障自动转移 + REST 在线增删 → 必须 distributed。分界点是三个词:容错、可在线管理、生产规模——任一为真就上 distributed。standalone 没有自动再平衡、没有 REST 集群管理、进程死即状态没,扛不住生产。
收尾挑战 · 把整套串起来

一句话能不能说清 Connect 的"控制面"特殊在哪?

不看教程,用一句话回答:为什么说 Connect 的 distributed 集群"没有外部协调器、没有数据库",它的全部状态在哪?再追一层——这个设计省掉了什么基建,又因此引入了什么新的失败模式?答得出这一句,说明 01 + 02 两章的主线(控制面全在 Kafka compacted topic 里)已经长进认知里了。

提示(卡住再展开)

线索:distributed worker 靠 Kafka 自己的消费组协议协调,全部状态(配置 / 偏移 / 状态)写进三个 compacted internal topic——所以不需要 ZooKeeper、不需要外部数据库。省掉的基建:一套独立的协调/存储组件。引入的新失败模式:这三个 topic 配置不一致或损坏 = 集群坏;group.id 配错 = 集群裂成两个。"配置即数据、控制面在 Kafka 里"既是它的优雅之处,也是它的脆弱之处。