起点 · 入口页
Kafka Connect:把"系统间搬数据"变成配置而非写码
基于 Apache Kafka 4.3(2026-05)。阅读时间约 1.5 小时。配置示例基于 Kafka 4.x,未在本机逐一运行。本教程是主 Kafka 教程的子教程,默认读者已掌握分区、消费组、offset。
·适合谁
这套教程为下面三类读者写——三条都对上,才会读得顺:
- 学过主 Kafka 教程或等同,理解分区 / 消费组 / offset / 再平衡,但没用过 Connect 的 Java 后端工程师。
- 要把数据库、对象存储、搜索引擎与 Kafka 之间的搬运标准化,正在纠结自写 producer/consumer 还是上 Connect 的人。
- 准备中高级面试,需要把 Connect 的架构、offset 语义、exactly-once、再平衡讲到能扛追问的人。
·不适合谁
三种情况下,别的资源更划算:
- Kafka 零基础:不清楚分区与消费组 offset 是什么,先回到主 Kafka 教程,那里建立的 schema 是本教程的地基。
- 要写流处理逻辑(按 key join、聚合、开窗、一进多出):那是 Kafka Streams 的领域,去 Kafka Streams 子教程,Connect 只搬运不做有状态计算。
- 找某个具体 connector 的配置手册(JDBC / S3 / Debezium 的全部参数表):去对应厂商文档;本教程讲框架机制与取舍,不逐一罗列 connector 参数。
·读完之后你能做到什么
你能把任意 Connect 故障定位到三层之一——是 connector 搬运、converter 序列化、还是控制面 topic 的问题;面试被问"Connect 怎么保证 exactly-once/不丢数据"时,能说清 source 偏移存哪、EOS 靠事务+fencing,而不是"它有容错"。具体到可验证的能力:
- 区分 source 与 sink connector、standalone 与 distributed worker、task 与 connector,并说出各自把配置 / 偏移 / 状态存在哪。
- 判断一个搬运需求该用 Connect、自写 producer/consumer、还是升级到 Streams,并讲出依据。
- 读懂一份 connector 的 JSON 配置:认出
key.converter/value.converter、tasks.max、errors.tolerance各自控制什么。 - 诊断
Unknown magic byte!、毒丸记录停 task、改名后重灌等典型失败,并定位到三层中的哪一层。 - 复述 EOS source 的成立条件:仅 distributed、全 worker 一致开启、事务原子双写记录与偏移、僵尸 task 被 fencing。
一句话本质
Kafka Connect 把"系统间搬数据"变成配置而非写码:connector 只管搬运、序列化格式由 converter 决定(不是 connector)、单条轻量变换由 SMT 做;而整个集群的配置/偏移/状态都存在 Kafka 自己的 compacted topic 里——没有外部协调器、没有数据库。
抓住"connector 搬运 + converter 定格式 + SMT 变换 + 控制面在 Kafka topic 里"这四块,整个 Connect 就清晰了。选型与面试的分水岭:知道何时用 Connect(标准化搬运)而非自写 producer/consumer,何时该升级到 Streams(有状态 / join / 聚合)。
·现状速览(截至 2026-06-03,Kafka 4.3)
| 状态 | 内容 |
|---|---|
| 稳定 | EOS source(KIP-618,3.3 GA,2022-09)、incremental cooperative 再平衡(KIP-415,2.3)、DLQ(KIP-298,2.0)、REST 偏移管理(KIP-875,3.6)。 |
| 近期变化 | 4.0(2025-03)起 Connect 需 Java 17(KIP-1032 转 Jakarta EE);ZooKeeper 已删(KRaft only);4.3(2026-05)加 KIP-1273 ConnectPlugin 接口(插件可发现性)。 |
| 生态 | JDBC / S3 / Debezium / Elasticsearch 等主流 connector 在核心 Kafka 之外(Confluent Hub / 厂商仓库),独立于 broker 版本演进。 |
这一页读着顺,不等于掌握了。三个假象,开读前先认清:
"我读得很顺"——Connect 的概念名词(converter、task、offset)单看都好懂,连起来的责任边界才是难点;顺着读完,合上页面能不能说清"序列化格式到底由谁决定"?
"我做题很快"——配置项眼熟不等于知道改错一个会触发哪种线上失败;自测的应用判别层才检验这个。
"我没卡壳"——没卡壳常常是因为还没遇到 source 与 sink offset 存两套、EOS 只在 distributed 成立这类反直觉点;卡壳是学到了的信号,不是没学会。
·概念地图
connect-offsets topic、sink 用消费组——两套不同机制;③ 控制面全在 Kafka compacted topic 里,无外部数据库。·学习路径建议
按目的选一条线,不必每页等量精读:
- 理解架构应付面试:01 概念 → 02 原理 → 03 自测,全程走完,自测的应用判别层重点做。
- 做数据集成选型:01 概念 → 02 原理 的备选方案对比表,抓 Connect / 自写 / Streams 的边界即可。
- 排查 Connect 故障:01 概念的 converter / offset / DLQ → 02 原理的序列化与 offset 恢复,对照三层定位法读。
·目录
·学完之后
把 Connect 这块 schema 立起来后,下面几个方向各自往上加一块:
- Debezium 与 CDC——基于 source connector 的变更数据捕获,在你的 schema 上加"从数据库事务日志读增量"这一能力。
- 具体 connector 深挖(JDBC source/sink、S3 sink)——把框架机制落到某个真实 connector 的参数与限制上。
- Kafka Streams 子教程——补上 Connect 故意不做的部分:有状态计算、join、聚合、开窗,明确两者分工边界。
- ksqlDB——用 SQL 表达流处理,理解它与 Streams、Connect 在同一生态里各占什么位置。