Chapter 02

机制与取舍:分布式集群怎么协调自己

01 章认识了 connector / converter / task 这些部件,知道了它们各管什么、配置长什么样。这一章往下钻一层:一个 distributed 集群没有外部协调器,凭什么能协调多个 worker?再平衡为什么不再 stop-the-world?序列化在字节层如何工作、偏移如何在重启后恢复、EOS 的事务与 fencing 如何咬合?每个机制都配一张备选方案对比表——看清 Connect 为什么这么设计,而不是只会照着配。

本章你将建立的 schema

  • distributed 集群的控制面就是三个 compacted topic(connect-configs/connect-offsets/connect-status)+ 消费组协议,REST 只是读写它们的门面——无 ZooKeeper、无数据库。
  • 再平衡是 incremental cooperative(KIP-415, 2.3):拓扑变更只暂停被移动的 task,其余继续搬——不再全员停摆。
  • 序列化由 converter 在字节层决定(Schema Registry 用 magic byte + schema id),与 connector 解耦;JsonConverter 与 JsonSchemaConverter 是不同 wire format。
  • 偏移分两套并各自恢复:source 存 connect-offsets 按 connector 名索引、sink 用消费组;EOS source 用事务双写 + 代际 fencing,是 worker 级、仅 distributed 的开关。

五个机制按从控制面到数据正确性的顺序排列:分布式协调(集群怎么不靠外部协调器活着)→ 再平衡(拓扑变了怎么重分 task)→ 序列化机制(字节层如何编解码)→ 偏移管理与恢复(重启从哪续)→ EOS source(不重不漏怎么做到)。每个机制先讲运行方式,再用备选方案对比表说清"为什么没选别的设计",然后是代价与失效模式,最后一道预测题。01 章给的是词汇表,这一章给的是这些词背后的工作机制与取舍——这是从"会配"到"理解"的那一层。

REST API :8083 读写控制面 WORKER 面 共享 group.id worker 1 task a worker 2 task b worker 3 task c 任一 worker 控制面 · 全在 Kafka(compacted topic) connect-configs 连接器/任务配置 connect-offsets source 偏移 connect-status 连接器/任务状态 读写状态 · 消费组协议协调 无 ZooKeeper 无外部数据库
图 2.0distributed 控制面:REST 是门面,worker 靠消费组协议加入同一集群,全部配置 / 偏移 / 状态都写进下方三个 compacted topic。 注意:这三个 topic 就是集群的"大脑"——没有任何外部协调器。它们的配置或压缩一旦不一致,整张图就裂成两个互不相认的集群。

2.1分布式协调:控制面全在 Kafka 自己的 topic 里

distributed 集群不靠任何外部协调器协调,而是用 Kafka 的消费组协议加上三个 compacted internal topic——控制面就是 Kafka 本身。

运行方式

多个 worker 配同一个 group.id,它们通过 Kafka 的消费组成员协议(与普通消费者加入消费组用的是同一套 group membership 机制)加入同一个 Connect 集群。集群里有一个被选出的 leader worker 负责计算 task 分配方案,其余 worker 执行分到自己头上的 task。三类需要持久化、需要跨 worker 一致的状态,各写进一个 compacted topic:

  • connect-configs:connector 与 task 的配置。POST 一个 connector 的本质,就是往这个 topic 追加一条配置记录。
  • connect-offsets:source connector 的偏移(见 §2.4)。
  • connect-status:connector 与 task 的运行状态(RUNNING/FAILED/PAUSED),GET /status 读的就是它。

这三个 topic 全是 compacted(日志压缩)——压缩保留每个 key 的最新值,正好匹配"配置 / 偏移 / 状态都只关心当前值"的语义。REST API(§1.2 已见)不是一个独立服务,它只是这套机制的门面:一次管理请求落到任意 worker,被翻译成对这三个 topic 的读或写,集群里所有 worker 读到后协同执行。这就是 01 章那句"配置即数据"在机制层的展开。

备选方案对比

"集群协调"是分布式系统的经典难题,常规答案是引入一个外部协调器(ZooKeeper、etcd)。Connect 选了另一条路。

表 2.1 · distributed 集群协调:三种设计
方案怎么协调 / 存状态为什么没选
外部协调器 ZooKeeper / etcd worker 把成员、配置、偏移托管给独立的协调集群 多引入一套要单独部署、监控、备份、扩容的有状态基建;Kafka 自己都在去 ZooKeeper(KRaft),Connect 再依赖它是逆流。被放弃。
独立数据库存配置/偏移 配置、偏移、状态写进关系库或 KV 存储 引入一个不在 Kafka 可用性域内的故障点——Kafka 活着但库挂了,集群照样瘫;还要解决库与 Kafka 之间的一致性。被放弃。
Kafka topic + 消费组协议 消费组成员协议管协调,三个 compacted topic 存全部状态 选中:零额外基建,控制面与数据面同生共死,复用 Kafka 已有的复制 / 持久化 / 压缩能力。
洞察 · 为什么是 compacted

三个 topic 用日志压缩而非按时间删除,是因为它们存的是"状态"不是"事件流"。状态只关心每个 key 的最新值——connector orders-source 当前的配置是什么、当前搬到哪个偏移、当前是不是 RUNNING。压缩保证 worker 重启后能从 topic 重建出完整的当前世界,又不会让 topic 无限增长。把这三个 topic 设成普通的按时间过期,旧配置 / 偏移会被删掉,集群重启即失忆——这正是它们必须 compacted 的原因。

代价与失效模式

"控制面就是几条 topic"换来零外部依赖,代价是这几条 topic 的健康直接等于集群的健康。两种典型失效:

  • internal topic 配置 / 压缩不一致 = 裂集群。如果新加的 worker 配了不同的 internal topic 名,或这些 topic 没设成 compacted、副本数不足,集群会读不到一致的控制面,表现为 connector 莫名消失、配置回滚、状态错乱。
  • group.id 配错 = 一个集群裂成两个。两批本应同属一个集群的 worker 配了不同 group.id,会各自组成独立集群、各跑一份 connector——同一个 source 被搬两遍。
陷阱

三个 internal topic 必须在全部 worker 上配成完全相同的名字 + 都 compacted + 副本数 ≥ 3。生产事故的常见形态是:手动建 topic 时漏设 cleanup.policy=compact,或不同 worker 的 config.storage.topic / offset.storage.topic / status.storage.topic 写得不一致——集群于是间歇性"分裂",connector 时有时无。修复方向:核对所有 worker 的这三项配置一字不差,并确认 topic 的 cleanup.policy。

想一想

运维把 connect-offsets 这个 topic 误建成了按时间删除(cleanup.policy=delete、保留 7 天)而非 compacted。集群短期看起来正常。两周后所有 source connector 重启了一次,会发生什么?

展开答案(先停 10 秒再点)

source connector 会大面积从头重灌。connect-offsets 存的是每个 source connector 的当前偏移(按 connector 名索引的 key)。按时间删除会把超过 7 天没更新的偏移记录直接删掉;compaction 本该保留每个 key 的最新值,这里却被删除策略丢了。重启后 worker 从这个 topic 读不到旧偏移,等于这些 connector 从未搬过,于是把外部源整个重灌一遍。

这道题指向的设计要点:三个 internal topic 必须 compacted,不是可选优化而是正确性前提。控制面在 Kafka 里的代价,就是这几个 topic 的配置错误会直接变成数据正确性事故。

与下一个机制的关系:worker 靠消费组协议加入集群——那么一个 worker 加入或离开时,它名下的 task 怎么重新分配给其他 worker?这就是再平衡,而 Connect 的再平衡方式经历过一次关键演进。

2.2再平衡:incremental cooperative 只暂停被移动的 task

再平衡是 worker 加入 / 离开 / 配置变更时重新分配 task 的过程;自 KIP-415(2.3)起用 incremental cooperative 策略,只暂停被收回或移动的那部分 task,其余继续搬。

运行方式

集群拓扑变化(一个 worker 宕机、新 worker 加入、connector 配置改了 task 数)会触发再平衡——leader worker 重新计算"哪个 task 该在哪个 worker 上"。关键在于怎么过渡到新方案:

  • 旧的 eager 策略(2.3 之前):再平衡一开始,所有 worker 立即放下手上全部 task(revoke everything),等 leader 算出新分配后再统一领回。整个集群在这期间停止搬运——一个 worker 的加入会让所有 connector 短暂停摆,这就是 stop-the-world。
  • incremental cooperative(2.3 起,KIP-415):leader 先算出新旧分配的差集——只有真正需要换主的 task 才被收回,其余 task 完全不受影响、继续搬。收回的 task 在下一轮再分配给目标 worker。整个过程"增量"完成,集群绝大部分吞吐不中断。

对一个跑着几十个 connector 的集群,差别是质变:扩容时加一台 worker,eager 会让全部 connector 抖一下,cooperative 只动那几个被迁移的 task。

触发:worker 3 加入集群 eager (2.3 前) worker 1 全部暂停 worker 2 全部暂停 worker 3 新加入 集群整体停摆 stop-the-world cooperative (2.3 起) worker 1 继续搬 worker 2 仅迁出 task d worker 3 接收 task d 其余 task 不中断 只动差集 迁移 1 个
图 2.1同一个触发(worker 3 加入)下两种再平衡的差别:eager 让全集群放下所有 task,cooperative 只迁移真正需要换主的 task d。 注意:cooperative 的关键不是"更快",而是只动新旧分配的差集——绝大多数 task 在再平衡期间完全不知情、持续搬运。

备选方案对比

表 2.2 · 再平衡策略:两代设计
方案过渡方式为什么(没)选
eager(stop-the-world) 再平衡时全员先放弃所有 task,再统一重新领取 实现简单、分配逻辑无需算差集;但任何拓扑变更都让整个集群停摆,connector 越多抖动越大。2.3 起被取代。
incremental cooperative(KIP-415) 只收回需要换主的 task,分多轮增量收敛,其余 task 不停 选中:拓扑 / 配置变更的影响面缩到最小,扩缩容与单点故障不再波及无关 connector。代价是协议更复杂、收敛要多轮。
洞察 · 与消费者再平衡同源

Connect 的 incremental cooperative 与 Kafka 普通消费者的 cooperative sticky assignor 是同一个思路的两次落地:都把"全员重来"改成"只动差集"。读者若已理解主 Kafka 教程里消费组的再平衡,这里只是把"被重分的是分区"换成"被重分的是 task"。这种一致性不是巧合——Connect 复用的就是 Kafka 那套 group membership 机制。

代价与失效模式

cooperative 不是没有锋利的边。worker 数在再平衡过程中继续变动会引入分配倾斜(KAFKA-12495):滚动重启时若 worker 一个接一个地进出,多轮增量再平衡可能收敛到一个 task 分布不均的中间态,部分 worker 过载、部分空闲。它不影响正确性,但会让吞吐不均,需要等集群稳定后再触发一次再平衡才会重新摊平。

陷阱

滚动升级或扩容时如果让 worker 进出过快(一个还没稳定下一个又动),cooperative 再平衡会卡在倾斜的中间分配上——表现为某些 worker CPU 打满、另一些几乎空闲,而 task 全是 RUNNING。这不是配置错误,是 KAFKA-12495 这类已知的再平衡时序问题。处理方向:操作 worker 时留足稳定窗口、避免并发进出;集群稳定后必要时手动触发一次再平衡让分配重新摊匀。

想一想

一个 distributed 集群跑着 20 个 connector、共 60 个 task。运维加入第 4 台 worker。用 incremental cooperative 策略,再平衡期间这 60 个 task 里大约有多少会被暂停?换成旧的 eager 策略呢?

展开答案(先停 10 秒再点)

cooperative:只有需要迁到新 worker 的那一小部分被暂停——大约是为了让新 worker 分到公平份额而迁移的 task 数(数量级在十几个上下,取决于均摊目标),其余四十多个继续搬。eager:全部 60 个先被收回、集群整体停摆,等新分配算完再统一领回。

这道题指向的设计要点:再平衡的"代价"不是发生频率,而是每次波及的 task 范围。cooperative 把这个范围从"全部"压到"差集",这正是它在大集群里取代 eager 的根本原因。

与下一个机制的关系:再平衡决定 task 落在哪个 worker 上跑——而 task 真正搬运时,要把结构化记录变成 Kafka 里的字节。这一步的机制独立于 connector,发生在字节层:序列化机制。

2.3序列化机制:converter 在字节层定格式,与 connector 解耦

序列化由 converter 在字节层完成、与 connector 完全解耦;走 Schema Registry 的格式在字节前缀里写 magic byte + schema id,而 JsonConverter 与 JsonSchemaConverter 是两种不可互换的 wire format。

运行方式

01 章(§1.4)讲了 converter 与 connector 解耦"是什么"。这里看字节层"怎么工作"。一条记录在 Connect 内部是带 schema 的结构化对象(SchemaAndValue);converter 负责把它变成 Kafka 里的字节,反向亦然。关键在于不同 converter 产出的字节布局完全不同:

  • Schema Registry 系(Avro / Protobuf / JSON Schema converter):字节不是裸数据,而是带前缀的 wire format——第 1 个字节是 magic byte(值 0x00),紧跟 4 字节的 schema id,之后才是序列化后的 payload。消费端读到 schema id,去 Schema Registry 拉对应 schema 再解码。schema 本身不随每条消息走,只走一个 id——这是它比内嵌 schema 省体积的原因。
  • JsonConverter(org.apache.kafka.connect.json):产出纯 JSON 字节,没有 magic byte。schemas.enable=true 时每条消息内嵌一份 {"schema":...,"payload":...},false 时只发裸 JSON。
  • JsonSchemaConverter(io.confluent.connect.json):名字像 JSON,但它走 Schema Registry、字节带 magic byte + schema id。与 JsonConverter 的 wire format 完全不同,不可互换——这是最容易混淆、也最容易出事的一对。

解耦体现在:同一个 connector,换 converter 就换了线上字节格式,connector 代码与配置一行不动(图 1.2 已示)。这把"序列化"从每个 connector 各实现一遍,变成一个可插拔的横切层。

Avro / Protobuf / JsonSchema converter magic 0x00 schema id 4 字节 序列化 payload Schema Registry 按 id 查 schema JsonConverter(纯 JSON,无前缀) { "id": 42, "amount": 9.9 } ← 直接是 JSON 字节
图 2.2两种 wire format 的字节布局:Schema Registry 系在最前面塞 magic byte + 4 字节 schema id,再接 payload;JsonConverter 从第一个字节起就是 JSON。 注意:sink 端若用 JsonConverter 去读上面那种带 0x00 开头的字节,第一个字节就解析失败(Unknown magic byte!)——这正是两类 wire format 不可互换的物理根源。

备选方案对比

表 2.3 · 序列化放在哪:两种设计
方案序列化逻辑的位置为什么(没)选
序列化烧进 connector 每个 connector 内部自己实现各种格式 换格式要改 connector 代码;每个 connector 各实现一遍 JSON/Avro/Protobuf,重复且不一致;同一 connector 无法跨格式复用。被放弃。
固定单一线上格式 整个 Connect 强制一种序列化(如只许 Avro) 无法对接已有的异构 topic(有的 JSON、有的 Avro、key 与 value 不同格式),现实里行不通。被放弃。
独立 converter 层(key/value 各一) connector 只产出/消费结构化对象,converter 负责对象↔字节 选中:任意 connector × 任意格式自由组合,换格式只动配置;key 与 value 可用不同 converter。代价是多一层、且 sink 端必须配对线上真实格式。

代价与失效模式

解耦的代价是sink 端的 converter 必须与 topic 里实际的字节格式严格一致,否则就是反序列化阶段的硬失败。最典型的失效是 sink converter 与线上格式不符引发的反序列化风暴:上游 producer 写 Avro(字节以 0x00 开头),sink 配了 JsonConverter,每一条记录在反序列化时都报 Unknown magic byte!。由于 errors.tolerance 默认 none(§1.7),第一条就让整个 task FAILED。converter 不会"猜"格式、不会默默降级——配错就是失败。

修复:sink converter 对齐线上 Avro 格式 JSON
{
  "value.converter": "io.confluent.connect.avro.AvroConverter",
  "value.converter.schema.registry.url": "http://schema-registry:8081",
  "key.converter": "org.apache.kafka.connect.storage.StringConverter"
}
陷阱

把 JsonConverter(org.apache.kafka.connect.json)当成 JsonSchemaConverter(io.confluent.connect.json)配,两端就此错位——一个发裸 JSON、一个期待 magic byte + schema id 的 JSON。它们名字只差一个词,wire format 却互不兼容。配 converter 时核对完整类名,尤其分清 org.apache.kafka 与 io.confluent 两个命名空间。

想一想

一个 source connector 用 JsonSchemaConverter(Confluent,带 magic byte)写一个 topic。下游有人新建一个 sink,把 value.converter 配成了 JsonConverter(Apache,纯 JSON),心想"反正都是 JSON"。sink 起得来吗?

展开答案(先停 10 秒再点)

起不来——sink 的 task 在反序列化阶段失败。JsonSchemaConverter 写出的字节以 magic byte 0x00 + 4 字节 schema id 开头,JsonConverter 从第一个字节就按 JSON 解析,撞上 0x00 立刻报错(典型是 Unknown magic byte! 或 JSON 解析异常)。"都是 JSON"是错觉:两者的 wire format 一个有 Schema Registry 前缀、一个没有,物理布局不同。

这道题指向的设计要点:converter 解耦带来灵活,也把"格式契约"完全压在配置上。org.apache.kafka...JsonConverter 与 io.confluent...JsonSchemaConverter 不可互换——选 converter 时认的是字节布局,不是名字里的"Json"。

与下一个机制的关系:converter 把记录序列化进 Kafka 后,task 还得记住"搬到哪了"以便重启续传。而 source 与 sink 的偏移存法和恢复路径完全不同——这是偏移管理与恢复。

2.4偏移管理与恢复:source 存 topic、sink 用消费组,两套各自续传

source connector 把自定义偏移存进 connect-offsets、按 connector 名索引,sink connector 用消费组(组名 connect-<name>)的位点;重启时各自从上次提交处续,REST /offsets(KIP-875)可读写两者。

运行方式

01 章(§1.6)建立了"两套偏移"这件事。这里看恢复路径怎么走。两套机制的差别源于一个事实:外部系统没有 Kafka 式的 offset 概念。

  • source 恢复:外部系统(DB、文件)的"位点"由 connector 自己定义——"读到 id=10042""文件读到 byte 4096"。框架把这组 {sourcePartition → sourceOffset} 键值对周期性提交进 connect-offsets topic,key 含 connector 名。worker 重启或 task 被再平衡到别处,新实例先按自己的 connector 名去 connect-offsets 读回最近一次提交的偏移,再从那里继续向外部系统要数据。
  • sink 恢复:sink 就是一个 Kafka 消费者,用普通消费组位点,组名固定为 connect-<connector 名>,位点存在 __consumer_offsets。重启即普通消费者重启——从消费组上次提交的位点接着消费,与任何 Kafka 消费者完全一样。

"按 connector 名索引"在 source 端埋了一个直接后果:给 source connector 改名 = 新名字在 connect-offsets 里查不到任何偏移 = 框架认为它从未搬过 = 从头重灌整个外部源。sink 端同理——改名意味着新消费组,从头消费。KIP-875(3.6)加了 REST GET/PATCH/DELETE /connectors/{name}/offsets,可以查看和手动调整两类偏移,但前提仍是名字稳定。

SOURCE source connector 自定义位点 框架提交 connect-offsets topic key 含 connector 名 重启按名读回 改名=丢偏移 SINK sink connector 就是个消费者 消费组提交 __consumer_offsets 组名 connect-<name> 重启续消费 普通消费者语义
图 2.3两套偏移机制并排:source 把自定义位点按 connector 名存进 connect-offsets,sink 用名为 connect-<name> 的消费组位点存进 __consumer_offsets。 注意:两套都以 connector 名为锚——改名在 source 端查不到旧偏移、在 sink 端是个新消费组,两边都从头开始。名字是偏移的主键。

备选方案对比

表 2.4 · source 与 sink 偏移:统一还是分两套
方案偏移怎么存为什么(没)选
统一用消费组偏移 source 也强行套 Kafka 消费组位点语义 source 读的是外部系统,位点是"DB 主键 / 文件字节",根本不是"topic+分区+offset",套消费组语义表达不了。被放弃。
统一存进外部存储 source 与 sink 偏移都写进一个独立的偏移库 引入 Kafka 之外的故障点(与 §2.1 同理),还要为 sink 放弃 Kafka 原生消费组这套成熟机制。被放弃。
source 存 topic / sink 用消费组 source 自定义偏移进 connect-offsets;sink 复用消费组位点 选中:各自用最贴合语义的机制——source 的位点框架可任意建模,sink 直接白拿消费者重启续传。代价是两套机制、按名索引,改名即丢位置。

代价与失效模式

分两套的代价集中在一个失效模式:改 source connector 名 = 丢偏移 = 重灌。想"重建"一个 source connector 而顺手把名字从 orders-jdbc-source 改成 orders-jdbc-source-v2,新名字在 connect-offsets 里对应一组空偏移,connector 会把整张表从头再搬一遍——对大表是一次代价高昂的全量重灌,还可能给下游灌入重复数据。

用 REST 显式管理偏移(KIP-875, 3.6+)· 别靠改名 bash
# 查看当前偏移(source 返回自定义位点;sink 返回消费组位点)
curl http://worker:8083/connectors/orders-jdbc-source/offsets

# 要重置时:先 stop,再 DELETE 偏移,而不是改名
curl -X PUT    http://worker:8083/connectors/orders-jdbc-source/stop
curl -X DELETE http://worker:8083/connectors/orders-jdbc-source/offsets
陷阱

把"重置一个 source connector 的进度"和"给它改名"当成一回事,是偏移管理里代价最高的失误。改名不会重置进度——它是创建了一个没有任何偏移记录的新 connector,于是从头重灌。要重置进度,保持名字不变,用 REST 先 stop 再 DELETE /offsets;要保留进度地重建,更要保持名字一字不动。

想一想

一个 JDBC source connector 已把一张 5 亿行的表增量同步了三个月。团队想"换个更规范的名字",于是删掉旧 connector、用新名字 POST 了一份配置完全相同(除了 name)的 connector。下一刻发生什么?

展开答案(先停 10 秒再点)

新 connector 会从头重灌整张 5 亿行的表。source 偏移在 connect-offsets 里按 connector 名索引,新名字查不到任何已提交偏移,框架认为这是个全新的、从未搬过的 connector,于是从 incrementing 列的起点重新拉全量——海量重复数据涌向下游,下游若无幂等还会被污染。

这道题指向的设计要点:connector 名是偏移的主键,不是一个可随意美化的标签。"两套偏移、按名索引"的代价就是名字必须当成契约对待——要改进度用 REST /offsets,绝不靠改名。

与下一个机制的关系:默认的偏移提交是"先写记录、再提交偏移",崩溃时这两步之间的窗口会导致重发——即"至少一次"。要做到"精确一次",必须把这两步绑成原子操作。这是 EOS source。

2.5EOS source:事务双写 + 代际 fencing,worker 级开关

EOS source(KIP-618)把"写记录"与"提交偏移"放进同一个 Kafka 事务原子完成,再用代际 fencing 拒绝僵尸 task 的写入;它是 worker 级开关 exactly.once.source.support,仅 distributed 可用、须全集群一致。

运行方式

默认 source 是"至少一次":task 先把一批记录写进目标 topic,再把对应偏移提交进 connect-offsets。如果崩溃发生在"记录已写、偏移未提交"这个窗口,重启后 task 从旧偏移续传,已写过的那批记录被重发。EOS source 用两个 Kafka 原语消除这个窗口:

  • 事务双写:task 把"一批数据记录"和"对应的偏移更新(写进 connect-offsets)"放进同一个 Kafka 事务,原子提交。要么记录与偏移一起生效,要么都不生效——"记录写了、偏移没写"的窗口不复存在。
  • 代际 fencing(僵尸隔离):再平衡后旧 task 实例可能还"活着"想继续写(僵尸)。框架给每代 task 分配递增的事务标识,broker 只接受最新代的事务、拒绝旧代的写入。同一份工作因此只有最新一代能提交,僵尸写不进去。

开启是两阶段滚动:先把所有 worker 的 exactly.once.source.support 设为 preparing 滚动重启一轮,再设为 enabled 滚动重启一轮。两阶段是为了让集群在升级过程中始终保持一致——不能一半 worker 开、一半没开。connector 侧再声明 exactly-once.support=required 表示它要求这个保证。

task(当前代) 事务 id = N 同一个 Kafka 事务 写数据记录 → topic 写偏移 → connect-offsets broker 原子提交 记录与偏移同生共死 僵尸 task(旧代) 事务 id = N-1 broker 拒绝旧代 fencing:写不进去
图 2.4EOS source 两件事:当前代 task 把"写记录"和"写偏移"装进同一个事务原子提交;旧代僵尸 task 因事务 id 过期被 broker fencing 拒绝。 注意:消除重复靠两件事咬合——事务保证"记录与偏移同生共死"(无重发窗口),fencing 保证"只有最新一代能写"(僵尸不重写)。少任一件都不成立。

备选方案对比

表 2.5 · source 端怎么做到精确一次:两种设计
方案怎么消除重复为什么(没)选
应用层 / 下游去重 接受 source 重发,靠下游按业务主键幂等或去重 把正确性责任推给每一个下游,各写各的去重逻辑、容易漏;source 仍在源头制造重复。作为框架级保证被放弃。
只提交偏移、不用事务 更频繁地提交偏移,缩小重发窗口 只是把窗口变小,没有消除——崩溃时机不巧仍会重发,做不到"精确一次"。被放弃。
Kafka 事务 + 代际 fencing 记录与偏移原子双写,僵尸 task 被 broker fencing 拒绝 选中:框架在源头保证精确一次,下游无需去重。代价:仅 distributed、worker 级全集群一致、两阶段升级、有事务开销。

代价与失效模式

EOS source 的代价主要是约束而非性能:仅 distributed 模式可用(standalone 没有这套机制)、是 worker 级开关、必须全集群一致(不能 per-connector 单独开)、开启要两阶段滚动。最常见的失效是认知性的:以为能"只给账务这一个 connector 开 EOS"——做不到,要么整个 distributed 集群一起开,要么都不开。此外事务提交本身有开销,会略增延迟。

connect-distributed.properties(worker 级、全集群一致) properties
# 两阶段:先全员 preparing 滚动重启,再全员 enabled 滚动重启
exactly.once.source.support=enabled
陷阱

exactly.once.source.support 不是 per-connector 配置,写进 connector 的 JSON 里不会生效——它是 worker 级、要写进每个 worker 的 properties、且全集群取值必须一致。在 standalone 模式下它根本不可用。"给某个 connector 单独开精确一次"这个需求,在 Connect 的 source EOS 模型里不存在;能控制的只有 connector 侧 exactly-once.support=required(声明该 connector 要求集群提供这个保证)。

想一想

某团队想让一个计费 source connector 精确一次,于是只在这个 connector 的 JSON 配置里加了 "exactly.once.source.support": "enabled",其余 worker 配置没动。重启后这个 connector 是精确一次吗?

展开答案(先停 10 秒再点)

不是。exactly.once.source.support 是 worker 级配置,只在 worker 的 properties 里读;把它塞进 connector 的 JSON 不会开启 EOS——框架不会从 connector 配置里读这个 worker 级开关。这个计费 connector 仍是默认的"至少一次",崩溃时照样可能重发。要真正开启,得把所有 worker 的 properties 设成 enabled(经 preparing 两阶段滚动),并在 connector 侧声明 exactly-once.support=required。

这道题指向的设计要点:EOS source 是集群级能力,不是单个 connector 的属性。它要 fence 整个集群里同名 task 的所有代、要全员用事务语义,所以只能整集群开——这是它与 per-connector 配置(如 tasks.max、errors.tolerance)的根本区别。

2.6跨机制综合:一条记录如何穿过整个集群

五个机制不是孤立的——一次真实的搬运同时用到协调、再平衡、序列化、偏移恢复。把它们串成一个场景,看它们如何协同。

场景:一个 distributed 集群(worker 1/2/3,exactly.once.source.support 未开)跑着一个 JDBC source connector,从 orders 表增量读、写进 pg-orders topic。某时刻 worker 2 宕机。跟着一条记录走一遍:

  1. task 分配(协调 §2.1 + 再平衡 §2.2):connector 启动时,leader worker 按 tasks.max 与可切分单元算出 task,均摊到三个 worker——假设 orders 表的 task d 落在 worker 2 上。这份分配方案写进 connect-configs。
  2. 读取 + 序列化(序列化 §2.3):worker 2 上的 task d 从 orders 读到一行(id=10042),在 Connect 内部是结构化对象;value.converter(设为 Avro)把它序列化成带 magic byte + schema id 的字节,写进 pg-orders。
  3. 提交偏移(偏移 §2.4):框架周期性把 task d 的位点 {table:orders → incrementing:10042} 提交进 connect-offsets,key 含 connector 名。
  4. worker 2 宕机 → 再平衡(§2.2):消费组协议检测到 worker 2 离开,触发 incremental cooperative 再平衡。worker 1/3 上无关的 task 继续搬,只有 task d 被重新分配——假设迁到 worker 3。
  5. 从偏移续传(§2.4 + §2.1):worker 3 上新起的 task d 实例,先按 connector 名去 connect-offsets 读回最近提交的偏移(incrementing:10042),从 id>10042 继续向 orders 要数据——不重读已搬过的行。整个恢复没碰任何外部协调器,全程只读写 Kafka 的 topic。

这一圈把控制面(§2.1 三个 topic)、再平衡(§2.2 只动 task d)、序列化(§2.3 converter 定字节)、偏移恢复(§2.4 按名读回)咬合在一起。注意一个边界:因为这个集群没开 EOS(§2.5),第 3 步与第 2 步之间存在窗口——若 worker 2 恰在"记录已写、偏移未提交"时宕机,task d 在 worker 3 上会从上一个已提交偏移续传,重发 10042 这批里已写的记录(至少一次)。要消除这点,才需要 §2.5 的事务双写。

洞察 · 五个机制的分工

一句话各归其位:协调决定状态存在哪(三个 topic)、再平衡决定 task 跑在哪(只动差集)、序列化决定记录长什么样(converter 定字节)、偏移恢复决定从哪续(按名读回)、EOS决定续得重不重(事务 + fencing)。前四个让搬运能在故障后正确地继续,第五个把"正确"从"至少一次"提到"精确一次"。

§本章 self-check

先合上教程,把你能想到的答案写在纸上或编辑器里。 写完再点开答案对照——直接点开等于把这一节当再读一遍。

  1. distributed 集群的控制面由哪三个 topic 构成、各存什么?为什么它们必须是 compacted、且全集群配置一致?
  2. incremental cooperative 再平衡相比 eager 改了什么?说清"只动差集"的含义,以及它在大集群里为什么是质变。
  3. 走 Schema Registry 的 wire format 字节开头是什么?为什么 sink 端把 JsonConverter 用在 Avro topic 上会失败?
  4. EOS source 靠哪两个 Kafka 原语做到精确一次?为什么它只能 worker 级、全集群一致地开,不能 per-connector?
  5. (跨机制综合题)一个 distributed 集群跑着一个 JDBC source(未开 EOS),某 worker 宕机后一条记录被重新分配、从偏移续传。把这个过程涉及的协调、再平衡、序列化、偏移恢复四个机制串起来讲一遍,并指出在哪一步可能产生重复、要消除它需要开什么。
答案(先做完再展开)
  1. connect-configs(connector/task 配置)、connect-offsets(source 偏移)、connect-status(connector/task 状态)。必须 compacted 是因为它们存的是"状态"——只关心每个 key 的最新值,压缩保证 worker 重启能从 topic 重建当前世界又不无限增长;按时间删除会丢掉旧 key 的最新值导致失忆(如偏移被删→重灌)。全集群配置一致(topic 名相同、都 compacted、group.id 一致)否则集群会读到不一致控制面或裂成两个。
  2. eager 在再平衡时让全员先放弃所有 task 再统一重领(stop-the-world);cooperative 只收回新旧分配的差集(真正需要换主的 task),其余 task 完全不停。质变在于:波及范围从"全部 task"压到"被迁移的少数"——大集群里加一台 worker,eager 让所有 connector 抖动,cooperative 只动那几个 task。
  3. 走 Schema Registry 的字节以 magic byte(0x00)+ 4 字节 schema id 开头,之后才是 payload。JsonConverter 是纯 JSON、没有这个前缀,用它读 Avro topic 时,它把第一个字节(0x00)按 JSON 解析立刻失败(Unknown magic byte!);因 errors.tolerance 默认 none,第一条就让 task FAILED。两者 wire format 物理布局不同,不可互换。
  4. ① Kafka 事务:把"写记录"与"写偏移(进 connect-offsets)"放进同一事务原子提交,消除"记录写了偏移没写"的重发窗口。② 代际 fencing:给每代 task 递增事务 id,broker 拒绝旧代僵尸 task 的写入。只能全集群开是因为它要 fence 整个集群里同名 task 的所有代、全员用事务语义——这是集群级能力,不是单 connector 属性;且仅 distributed 可用、两阶段(preparing→enabled)滚动升级。
  5. ① leader 按 tasks.max 与可切分单元算出 task 均摊到各 worker,方案写进 connect-configs(协调 §2.1);② task 从表读行、converter 序列化成字节写进 topic(序列化 §2.3);③ 框架周期性把位点按 connector 名提交进 connect-offsets(偏移 §2.4);④ worker 宕机触发 incremental cooperative 再平衡,无关 task 不停、只有该 task 被迁到别的 worker(再平衡 §2.2);⑤ 新 task 实例按 connector 名读回最近偏移、从该位点续传,全程只读写 Kafka topic、不碰外部协调器(§2.4+§2.1)。产生重复的点:第②与第③步之间有窗口——若崩溃在"记录已写、偏移未提交"时,续传会重发那批已写记录(至少一次)。要消除它需开 EOS source(§2.5):把记录与偏移原子双写进同一事务,并对僵尸 task 做 fencing。
进阶挑战 · 刚好够不着

EOS source 用 Kafka 事务搞定原子双写——那 EOS sink 为什么不能照搬同一招?

§2.5 的精确一次靠"把记录和偏移放进同一个 Kafka 事务"。这一招成立的前提,是这两个被写的对象都在 Kafka 里。现在反过来想 sink 端:sink 要把数据写进的是外部系统(ES、JDBC),消费位点在 Kafka 的消费组里。想一想:为什么"用一个 Kafka 事务把外部写入和消费位点绑在一起"这条路走不通?sink 端要做到精确一次,必须依赖什么、由谁提供?

提示(卡住再展开)

Kafka 事务只能覆盖"写进 Kafka 的操作"。source 的两个对象(数据记录、偏移)都是写进 Kafka 的,所以一个事务能同时管。sink 的两个对象是"外部系统的写入"(在 Kafka 之外)和"消费位点"(在 Kafka 内)——前者根本不在 Kafka 事务的管辖范围,事务提交了也无法回滚一次已发生的 ES 写入。线索:sink 端的精确一次必须靠外部系统自己的能力——幂等写入(同一记录写多次效果等同一次,如按主键 upsert)或外部系统自己的事务,让重复投递不产生重复效果。这就是为什么 EOS source 是一个统一的框架开关,而"sink 精确一次"取决于目标系统支不支持、没有对应的全局开关。可对照主教程的投递语义与事务一节。