Chapter 01
核心概念:配置驱动的数据搬运
起点页给出了本质:Connect 把"系统间搬数据"从写码变成写配置——connector 搬运、converter 定格式、SMT 变换,控制面全在 Kafka 自己的 topic 里。这一章把这句话拆成八个能在面试里准确定义、在配置里准确落地的概念。
本章你将建立的 schema
- connector 只管"搬",序列化格式由 converter 决定——这两者解耦,一个 connector 可配任意格式。
- worker 是运行进程(standalone 单机无容错 / distributed 集群有 REST 与自动再平衡),connector 在 worker 里派生出 ≤
tasks.max个 task 作并行单位。 - 偏移分两套:source 存进
connect-offsetstopic、按 connector 名索引;sink 用普通消费组——改名即丢位置。 - 错误默认
none(一条坏记录停整个 task);DLQ 仅 sink、只抓 converter/SMT 错误。EOS source 是 worker 级、全集群一致的开关。
八个概念按依赖顺序排列:connector(搬什么)→ worker(在哪跑)→ task(多快)→ converter(什么格式)→ SMT(路上改什么)→ offset(读到哪了)→ 错误处理 / DLQ(坏记录去哪)→ EOS source(不重不漏)。前四个回答"一条记录怎么从外部系统流进 Kafka 又流出去",后四个回答"出错、重启、扩容时它还正确吗"。读者若已学过主 Kafka 教程的 partition / 消费组 / offset,这一章会反复借用那套词汇。
1.1connector:搬运插件,POST 一段 JSON 就跑
connector 是一个可复用的搬运插件:source connector 把外部系统的数据写进 Kafka,sink connector 把 Kafka 的数据写进外部系统。
没有 Connect,"把 Postgres 表同步进 Kafka"这种活得自己写一个 producer 程序:轮询数据库、把行转成记录、处理位点跟踪、处理重启续传、处理并行、再写一套部署和监控。每接一个新系统就重写一遍这套胶水。Connect 把这套胶水抽成插件契约——JDBC、S3、Debezium、Elasticsearch 等厂商各实现一个 connector,工程师把 JAR 放进 plugin path,再 POST 一段 JSON 配置,搬运就跑起来,几乎不写码。
底层机制(比文档深一层)
文档说"connector 负责搬数据"。再深一层:connector 类本身不搬数据。它只做两件事——验证配置、把工作切成若干份。真正读写数据的是它派生出的 task(§1.3)。所以一个 connector 实例 = 一份配置 + 一个"把活切成 N 份"的计划,框架拿着这个计划把 task 分发到各 worker 上跑。这解释了一个反直觉的事实:connector 配置里写 "tasks.max": "10" 不保证有 10 个 task 在干活——connector 自己决定能切出几份(JDBC source 受表数限、sink 受分区数限),tasks.max 只是上限。
还有一层:connector 只管"搬",不管序列化格式。一条记录在 connector 内部是结构化对象(Connect 的内部 SchemaAndValue 表示),变成 Kafka 里的字节是 converter(§1.4)的事。这个解耦是 Connect 设计的关键,后面专门讲。
connector 像货运公司的标准化集装箱接口:你不关心箱子里装的是冰箱还是钢材,只要按接口装卸即可。边界:集装箱本身是被动的,而 connector 会主动验证你的运单(配置)、决定派几辆车(task 数)。它继承了承运方的能力上限——connector 有 bug 或限制,你就有;调优时仍要懂它下一层的配置。
场景走查
把一个 Postgres 的 public.orders 表增量同步进 Kafka。装好 JDBC source connector 的 JAR 后,POST 这段配置到 worker 的 REST 端点:
{
"name": "orders-jdbc-source",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"connection.url": "jdbc:postgresql://db:5432/shop",
"table.whitelist": "orders",
"mode": "incrementing",
"incrementing.column.name": "id",
"topic.prefix": "pg-",
"tasks.max": "1"
}
}
connector.class 选定哪个插件;mode: incrementing 让它靠自增主键追新行;topic.prefix 把表写进 pg-orders topic。注意这里没有任何序列化配置——格式留给 worker 或 connector 级的 converter,下面 §1.4 才出现。这正是"connector 与格式解耦"在配置上的体现。
Kafka 自带的跨集群复制工具 MirrorMaker 2 不是独立程序,而是三个 source connector:MirrorSourceConnector(复制数据)、MirrorCheckpointConnector(翻译消费位点供故障转移)、MirrorHeartbeatConnector(探活)。它因此白拿了 Connect 的扩缩容与高可用——这是"connector 即插件"抽象复用性的最好证据。
与下一个概念的关系:配置 POST 到哪?谁来跑这个 connector、谁来响应这个 REST 请求?那是运行进程——worker。
1.2worker:运行进程,standalone 还是 distributed
worker 是真正运行 connector 与 task 的 JVM 进程;standalone 是单进程、文件配置、无容错,distributed 是多进程共享一个集群、靠 REST 管理、自动再平衡。
connector 是一份配置和一段插件代码,它得有个进程来加载、运行、监控。worker 就是这个进程。两种模式回答的是同一个运维问题的两端:开发期单机搬一个日志文件,要的是简单——一个进程、一个 properties 文件就够;生产期要的是容错与可管理——某台机器宕了任务自动转移、不重启进程就能加 connector。standalone 服务前者,distributed 服务后者。
底层机制(比文档深一层)
文档说 distributed 模式"有容错和负载均衡"。再深一层:它靠 Kafka 自己实现,没有外部协调器。多个 worker 配同一个 group.id,通过 Kafka 的消费组协议加入同一个 Connect 集群;connector 配置、source 偏移、connector/task 状态全写进三个 compacted internal topic。这意味着 Connect 集群不需要 ZooKeeper、不需要数据库——它的全部状态就是几条 Kafka topic。代价直接:这三个 internal topic 配置不一致或损坏,整个集群就坏;group.id 配错会把一个集群裂成两个。
standalone 则把这套全省了:偏移存本地文件(offset.storage.file.filename),无 REST 集群管理,进程一死状态随它去。所以 standalone 不是"小一号的 distributed",而是另一套存储与协调模型——这是面试常考的分水岭。
REST API:connector 的生命周期入口
distributed worker 暴露一个 REST 端点(默认 8083),它是管理 connector 的唯一正道。整个生命周期都走它:
# 创建(POST 上一节那段 JSON)
curl -X POST -H "Content-Type: application/json" \
--data @jdbc-source.json http://worker:8083/connectors
# 看状态:connector 与每个 task 是 RUNNING 还是 FAILED
curl http://worker:8083/connectors/orders-jdbc-source/status
# 暂停 / 恢复 / 重启失败的 task
curl -X PUT http://worker:8083/connectors/orders-jdbc-source/pause
curl -X POST http://worker:8083/connectors/orders-jdbc-source/restart?includeTasks=true
# 删除
curl -X DELETE http://worker:8083/connectors/orders-jdbc-source
请求发给任意一个 worker 都行——它会把变更写进 connect-configs topic,集群里所有 worker 读到后协同执行。这就是"配置即数据":管理动作本质是往 Kafka 写一条记录。
group.id 与 REST,把全部状态写进三个 compacted internal topic。
注意:distributed 没有外部协调器——它的"集群大脑"就是右下角那三条 Kafka topic,group.id 配错会把一个集群裂成两个。与下一个概念的关系:图 1.1 里每个 worker 上跑着 task a/b/c——这些 task 从哪来、怎么被分到不同 worker 上?这是并行的单位:task。
1.3task:并行单位,tasks.max 是上限不是保证
task 是 connector 派生出的、真正读写数据的执行单位;一个 connector 最多派生 tasks.max 个 task,它们被均摊到集群各 worker 上并行跑。
若一个 connector 只能单线程搬,那一张大表、一个高吞吐 topic 就被单线程卡死,加机器也没用。task 是 Connect 的并行原语:把一个 connector 的工作切成多份,分到多个 worker 上同时跑。"并行单位是 task,不是 connector"——这一句决定了怎么估算吞吐、怎么扩容。
底层机制(比文档深一层)
关键、且面试高频:tasks.max 是上限,不是保证。实际 task 数由 connector 根据"能切成几份"决定,再对 tasks.max 取下限。两类典型约束:
- sink connector:并行度受订阅 topic 的分区数限——每个 task 是一个消费组成员,分区数即并行上限。
tasks.max=10但 topic 只有 3 个分区,只有 3 个 task 拿到分区,其余 7 个空转。 - JDBC source connector:并行度受表数限——它一张表给一个 task,同步 2 张表最多 2 个 task。
task 分到哪个 worker,由集群的 assignor 决定,目标是均摊。worker 加入或离开会触发再平衡重新分配——Connect 自 2.3 起用 incremental cooperative 再平衡(只暂停被收回/移动的 task,其余继续),机制细节是 02 章 §2.2 的事,这里只需记住:task 是被调度的最小单位。
场景走查
一个 sink connector 要把 pg-orders(4 个分区)写进 Elasticsearch。配 tasks.max 等于分区数最划算:
{
"name": "orders-es-sink",
"config": {
"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
"topics": "pg-orders",
"connection.url": "http://es:9200",
"tasks.max": "4"
}
}
这里 4 个 task 各认领 1 个分区、并行写 ES。若把 tasks.max 写成 8,多出的 4 个 task 不会报错,只是空转——没有分区可认领。`GET /status` 会显示 8 个 task 都 RUNNING,但其中 4 个不干活,是排查吞吐问题时的常见误判点。
一个 sink connector 配了 "tasks.max": "10",但它订阅的 topic 只有 3 个分区。实际有几个 task 在搬数据?
展开答案(先停 10 秒再点)
3 个。sink 的每个 task 是消费组成员,并行度被分区数封顶。框架会启动 10 个 task,但只有 3 个能分到分区,剩下 7 个拿不到分区、空转。
这道题指向的设计要点:tasks.max 是上限不是保证,sink 真实并行度 = min(tasks.max, 分区数)。想提高 sink 吞吐,先加分区,再加 task——单加 tasks.max 只是制造空转的进程。
与下一个概念的关系:task 把结构化记录搬到了 Kafka 边界,但 Kafka 里存的是字节。结构化对象怎么变字节、又怎么变回来?这是和 connector 解耦的那一层——converter。
1.4converter:决定序列化格式,与 connector 解耦
converter 负责记录在 Connect 内部对象与 Kafka 字节之间的相互转换;key 和 value 各配一个,它决定线上的序列化格式——且与 connector 完全解耦。
如果序列化格式烧死在每个 connector 里,那"换 JSON 为 Avro"就得改 connector 代码、每个 connector 各实现一遍各种格式。Connect 把序列化抽成独立的 converter 层:connector 只产出/消费结构化对象,converter 负责对象↔字节。于是任意 connector × 任意格式自由组合,换格式只动配置不动代码。
底层机制(比文档深一层)
三个常被忽略、却最容易出事的点:
- key 和 value 各一个 converter,独立配置。
key.converter和value.converter互不相干——线上常见 key 是 String、value 是 Avro。两边配错任一个,就只有半边反序列化失败。 - converter 可配在 worker 级或 connector 级。worker 级是集群默认,connector 级覆盖它。一个 connector 完全可以用与集群默认不同的格式。
JsonConverter≠JsonSchemaConverter。前者是org.apache.kafka.connect.json.JsonConverter(纯 JSON,可选每条内嵌 schema);后者是io.confluent.connect.json.JsonSchemaConverter(走 Schema Registry、字节里带 schema id 的 wire format)。名字像,wire format 完全不同,不可互换。
再深一层是 JsonConverter 的 schemas.enable:设 true 时每条消息都内嵌一份 {"schema":...,"payload":...} 结构,体积成倍膨胀但自描述;设 false 只发裸 JSON。这个开关与"消费端期望什么结构"必须对齐,否则 sink 报 JsonConverter requires "schema" and "payload"。
AvroConverter 的字节带 magic byte + schema id,正是它与 JsonConverter 互不兼容的根源。场景走查
给 §1.1 那个 JDBC source 加上序列化:key 用 String、value 用 Avro 并接 Schema Registry。converter 配在 connector 级(覆盖 worker 默认):
{
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"value.converter": "io.confluent.connect.avro.AvroConverter",
"value.converter.schema.registry.url": "http://schema-registry:8081"
}
worker 级默认则写在 properties 里,对全集群生效:
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
sink 端的 converter 必须匹配 topic 里实际的字节格式,而不是你以为的格式。producer 写的是 Avro(字节以 magic byte + schema id 开头),sink 却配 JsonConverter,反序列化会失败并报 Unknown magic byte!——它把 Avro 的第一个字节当成 JSON 起始字符去解。修复:sink 的 value.converter 改 AvroConverter 并对齐 schema.registry.url。这类"反序列化风暴"是 sink connector 最常见的失败模式。
上游 producer 用 Avro 写入一个 topic(每条字节带 magic byte + schema id)。一个 sink connector 配 value.converter=JsonConverter 去消费它。会发生什么?
展开答案(先停 10 秒再点)
sink 的每个 task 在反序列化阶段就失败,报 Unknown magic byte!。JsonConverter 拿到 Avro 字节的第一个字节(值为 0x00 的 magic byte),按 JSON 解析立刻失败。由于 errors.tolerance 默认是 none,第一条记录就让整个 task 停掉(见 §1.7)。
这道题指向的设计要点:converter 与 connector 解耦带来灵活,也带来一个新约束——sink converter 必须与线上真实格式一致。converter 不会"猜"格式,配错就是反序列化失败,而非默默降级。
与下一个概念的关系:converter 决定记录的"外壳格式",但有时要在搬运途中改记录的内容——加个字段、脱敏、改目标 topic。这是逐条的轻量变换:SMT。
1.5SMT:逐条变换,无状态、无 join
SMT(Single Message Transform,单消息变换)是配置在 connector 里的链式逐条变换,每条记录流过时被原地修改,可加 predicate 谓词做条件门控。
很多集成需求只是对每条记录做点轻微编辑:加一个摄入时间戳、把某字段脱敏、按内容改写目标 topic 名。为这点事单起一个 Kafka Streams 应用是杀鸡用牛刀。SMT 让这类逐条编辑变成 connector 配置里的几行——零代码、内联在搬运路径上。
底层机制(比文档深一层)
SMT 以链的形式声明:transforms 列出名字,每个名字配 type 和参数,记录按顺序流过整条链。常用的有 InsertField(插字段)、MaskField(脱敏)、RegexRouter(按正则改 topic 名)、ReplaceField(删改字段)。自 2.6 起每个 SMT 可挂一个 predicate,只对满足条件的记录生效。
关键边界、且是面试判别点:SMT 是单条、无状态的。它一次只看一条记录,没有跨记录的内存。所以 SMT 做不了聚合、join、按 key 关联另一个 topic、一条变多条——这些是有状态的流处理,属于 Kafka Streams / ksqlDB。把"逐条无状态编辑用 SMT,有状态关联用 Streams"刻进判断里,选型就不会错。代价:很重的 SMT 链会增加每条记录的 CPU 开销,因为它在搬运热路径上同步执行。
场景走查
给搬进来的订单记录做两件事:插入一个固定来源标记字段,并把所有记录从 pg-orders 改写到 ingested-orders topic。两个 SMT 串成一条链:
{
"transforms": "addSource,route",
"transforms.addSource.type":
"org.apache.kafka.connect.transforms.InsertField$Value",
"transforms.addSource.static.field": "source_system",
"transforms.addSource.static.value": "postgres-shop",
"transforms.route.type":
"org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "pg-(.*)",
"transforms.route.replacement": "ingested-$1"
}
记录先过 addSource(value 里多出 source_system: "postgres-shop"),再过 route(topic 名 pg-orders 被正则改写成 ingested-orders)。整条链对每条记录同步执行,无任何跨记录状态。
需求若是"按订单 key 关联另一个 topic 的用户信息再聚合",SMT 做不到——它无状态、看不到别的记录。强行用 SMT 会卡在"拿不到第二条记录"上。这类有状态关联是 Kafka Streams / ksqlDB 的职责。SMT 与 Streams 的边界(单条无状态 vs 有状态 join/聚合)是面试判别题的常客。
与下一个概念的关系:connector 搬、converter 转、SMT 改——这条流水线跑起来后,重启时怎么知道"上次搬到哪了"?这就是 offset,而且 source 和 sink 的存法完全不同。
1.6offset:source 存 topic、sink 用消费组
offset 记录"搬到哪了"用于重启续传;source connector 把自定义偏移存进 connect-offsets topic、按 connector 名索引,sink connector 用普通 Kafka 消费组的位点。
搬运进程会重启、会崩溃。没有偏移记录,重启后要么从头重灌全部数据、要么漏掉中间的。offset 让 connector 重启后从上次提交的位置续上,做到不重不漏(或至少不漏)。手写 producer/consumer 时这套位点跟踪得自己实现——Connect 把它白送了,这本身就是用 Connect 而非手写的一大理由。
底层机制(比文档深一层)
面试高频:source 和 sink 的偏移是两套完全不同的机制。
- source connector:外部系统(数据库、文件)没有 Kafka 式 offset 的概念,所以 connector 自己定义"位点"长什么样——比如"读到
id=10042"或"文件读到 byte 4096"。框架把这组{sourcePartition → sourceOffset}键值对存进connect-offsetstopic,按 connector 名索引。 - sink connector:它就是个 Kafka 消费者,用普通消费组的位点,消费组名是
connect-<connector 名>。位点存在 Kafka 的__consumer_offsets里,和任何消费者一样。
"按 connector 名索引"埋了个直接后果:给 source connector 改名 = 框架认不出旧偏移 = 从头重灌。新名字对应一组空偏移,connector 会把外部系统从头再搬一遍。KIP-875(3.6)加了 REST 的 GET/PATCH/DELETE /connectors/{name}/offsets,可以查看和手动改偏移,但前提仍是名字稳定。
想"重建"一个 source connector 而顺手改了它的名字(如 orders-jdbc-source → orders-jdbc-source-v2),结果它把整张表从头重灌一遍——因为偏移按旧名索引,新名查不到任何已提交位点。修复:保持 connector 名稳定;确实要重置时,用 REST /offsets 显式管理,而不是靠改名。
场景走查
用 REST 查一个 source connector 当前的偏移,确认它搬到了哪条主键:
curl http://worker:8083/connectors/orders-jdbc-source/offsets
# 返回(节选):source 自定义的位点,按 connector 名归档
# {
# "offsets": [
# { "partition": { "table": "orders" },
# "offset": { "incrementing": 10042 } }
# ]
# }
返回里 partition 和 offset 都是 connector 自定义的结构——这正是"source 偏移是自定义格式、框架只负责存取"的体现。同一个端点对 sink connector 返回的则是普通消费组位点({topic, partition, offset})。
与下一个概念的关系:偏移让 connector 知道搬到哪了。但搬运途中遇到一条处理不了的记录怎么办?默认行为出乎意料地严厉——这是错误处理与 DLQ。
1.7错误处理 / DLQ:默认一条坏记录停整个 task
errors.tolerance 控制遇错行为,默认 none——一条处理失败的记录就让整个 task 停;DLQ(死信队列)把坏记录转去另一个 topic,但仅 sink 可用、且只抓部分错误。
真实数据流里总有"毒丸记录":格式不符、字段缺失、反序列化失败。如果一条坏记录直接让搬运停摆,运维要半夜起来手动跳过。错误处理与 DLQ 让 connector 能容忍坏记录、把它们隔离到一边继续跑,而不是整条管道停摆。
底层机制(比文档深一层)
三个必须记准的点:
errors.tolerance默认是none。意思是零容忍——任何一条记录在 converter / SMT / 反序列化阶段出错,整个 task 立刻进入FAILED。改成all才会跳过坏记录继续。这个默认值常让人措手不及:一条脏数据就能让搬运在凌晨停掉。- DLQ 仅 sink connector 可用。
errors.deadletterqueue.topic.name指定一个 topic,坏记录被转发过去;配errors.deadletterqueue.context.headers.enable=true还会把"为什么失败"写进消息 header。source connector 没有 DLQ。 - DLQ 只抓 converter / SMT / key-value 错误,不抓 sink 写外部系统的失败。这是最大的认知误区——以为开了 DLQ 就万无一失。实际上"ES 拒绝写入""JDBC 主键冲突"这类投递失败绕过 DLQ,由
errors.retry.*重试逻辑处理,重试耗尽仍会让 task 失败。
场景走查
给 ES sink 加上容错:跳过坏记录、转去 DLQ、带上失败原因,并设重试窗口:
{
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "dlq-orders-es",
"errors.deadletterqueue.context.headers.enable": "true",
"errors.retry.timeout": "60000",
"errors.log.enable": "true"
}
此时一条反序列化失败的记录会被写进 dlq-orders-es(header 里带异常信息),task 继续搬后面的记录。但若 ES 本身拒绝一次合法写入,那不是 converter/SMT 错误——它走 errors.retry.timeout 的 60 秒重试窗口,不会进 DLQ;重试窗口内没成功,task 仍会失败。
把 DLQ 理解成"数据形状的隔离区",不是"所有故障的兜底"。它拦的是"这条记录读不懂/变换不了"(converter、SMT),放过的是"目标系统暂时不收"(写外部失败)。两类问题用两套机制——DLQ 管前者,errors.retry.* 管后者。混淆这条边界是错误处理配置里最常见的误判。
与下一个概念的关系:容错保证了"坏记录不停摆"。但还有更强的正确性诉求——能不能保证每条记录恰好一次地搬进 Kafka,重启、再平衡都不重复?这是 EOS source。
1.8EOS source:worker 级的精确一次(KIP-618)
EOS source(exactly-once source,KIP-618)让 source connector 把记录与偏移原子地写进 Kafka,配合僵尸隔离做到精确一次;它是 worker 级开关,仅 distributed 可用、且必须全集群一致开启。
默认情况下 source connector 是"至少一次":崩溃重启时,已写进 Kafka 但偏移没来得及提交的记录会被重发。对账务、计费这类不容重复的场景,重复记录是正确性事故。EOS source 把"写记录"和"提交偏移"绑成一个原子事务,让重启/再平衡时不产生重复。
底层机制(比文档深一层)
它靠两个 Kafka 原语:
- 事务:task 把"一批数据记录"和"对应的偏移更新(写进
connect-offsets)"放进同一个 Kafka 事务,原子提交。要么记录和偏移一起生效,要么都不生效——杜绝了"记录写了、偏移没写、于是重启重发"这个窗口。 - fencing(僵尸隔离):再平衡后旧的 task 实例可能还"活着"想继续写(僵尸)。框架给每代 task 分配递增的事务标识,旧代的写入被 broker 拒绝,确保同一份工作只有最新一代能提交。
配置上两个约束必须记住:它是 worker 级开关 exactly.once.source.support,不是 per-connector;并且仅 distributed 模式可用。开启要两阶段滚动:先把所有 worker 设成 preparing、滚动重启,再设成 enabled、再滚动一次——中途全集群必须一致,不能一半开一半不开。connector 侧再声明 exactly-once.support=required 表示它要求这个保证。
场景走查
worker 级开启(properties,全集群一致):
# 两阶段升级:先全员 preparing 滚动重启,再全员 enabled
exactly.once.source.support=enabled
connector 侧声明要求这个保证:
{
"exactly-once.support": "required"
}
exactly.once.source.support 是 worker 级、全集群一致的开关——不能只给某一个 connector 单独开。在 standalone 模式下它不可用。试图"只给账务 connector 开 EOS"会失败:要么整个 distributed 集群一起开,要么都不开。EOS 机制(事务双写 + fencing)的完整推导留给 02 章 §2.5。
承接下一章:这八个概念给了词汇表。它们背后的机制与取舍——分布式控制面为何能只靠 Kafka topic、incremental cooperative 再平衡解决了什么、序列化在字节层如何工作、偏移如何恢复、EOS 的事务与 fencing 如何咬合——是 02 章的主题。
§本章 self-check
先合上教程,把你能想到的答案写在纸上或编辑器里。 写完再点开答案对照——直接点开等于把这一节当再读一遍。
- 一句话说清 source connector 与 sink connector 的方向差异;再说清 connector 与 converter 各管什么、为什么要解耦。
- standalone 与 distributed 各把 connector 配置、source 偏移、connector 状态存在哪里?distributed 为什么不需要 ZooKeeper?
- source connector 和 sink connector 的偏移分别存在哪、按什么索引?为什么给 source connector 改名会导致从头重灌?
- (设计题)要把一张 Postgres 表增量同步进 Elasticsearch,途中给每条记录脱敏一个手机号字段、并容忍偶发的脏数据。你会用哪些概念(connector / converter / SMT / 错误处理),各自配什么关键项?哪类失败 DLQ 抓不到?
答案(先做完再展开)
- source connector 把外部系统数据写进 Kafka,sink connector 把 Kafka 数据写进外部系统。connector 只管"搬"(读写哪个系统、切几个 task),converter 管"序列化格式"(结构化对象 ↔ Kafka 字节)。解耦是为了让任意 connector 配任意格式——换序列化只动配置不动 connector 代码。
- standalone:配置在 properties 文件、偏移在本地文件(
offset.storage.file.filename)、无独立状态存储,进程死即没。distributed:配置存connect-configs、source 偏移存connect-offsets、状态存connect-status,三个都是 compacted internal topic。不需要 ZooKeeper 是因为 worker 靠 Kafka 消费组协议协调、靠这三个 topic 持久化全部状态——集群大脑就是几条 Kafka topic。 - source 偏移存进
connect-offsetstopic,是 connector 自定义格式,按 connector 名索引;sink 用普通消费组位点(组名connect-<name>),存在__consumer_offsets。改 source connector 名 = 新名字查不到任何已提交偏移 = 框架认为它从未搬过 = 从头重灌整个外部源。 - 用 JDBC source connector(
mode: incrementing追新行)搬,用MaskField这个 SMT 脱敏手机号字段,converter 按线上格式配(如 value 用 Avro + Schema Registry,key 用 String)。容忍脏数据:errors.tolerance=all+errors.deadletterqueue.topic.name(仅 sink 端有效)+context.headers.enable=true。DLQ 抓不到的:ES 端拒绝写入这类投递失败——它走errors.retry.*重试,重试耗尽仍会让 task 失败。
为什么 EOS 是 source 专属,sink 端的"精确一次"是另一回事?
本章讲的 EOS(KIP-618)只覆盖 source connector——它能把"写记录 + 提交偏移"放进一个 Kafka 事务原子完成。但 sink connector 要把数据写进的是外部系统(ES、JDBC),那不是 Kafka,没法纳入 Kafka 事务。想一想:为什么 source 端能用一个 Kafka 事务搞定原子性,而 sink 端的精确一次必须依赖目标系统本身的能力?sink 端要做到精确一次,外部系统需要提供什么?
提示(卡住再展开)
source 端的"双写"两个对象(数据记录、偏移)都在 Kafka 里,所以一个 Kafka 事务能同时覆盖。sink 端的两个对象是"外部系统的写入"和"消费位点"——前者在 Kafka 之外,Kafka 事务管不到。线索:要让 sink 精确一次,外部系统得支持幂等写入(同一条记录写多次效果等同一次,比如按主键 upsert)或自己的事务,让重复投递不产生重复效果。这就是为什么 sink EOS 不是一个统一框架开关,而是"看目标系统支不支持"。