主题
第 16 章 Kafka Connect:让数据像水一样在系统间流动
目标读者:每次要把 MySQL 的数据搬到 Kafka 里,就自己写个轮询脚本;要把 Kafka 的数据落到 ES 或者 ClickHouse 里,就自己写个 Consumer + 重试 + 死信,结果项目里堆了几十个「胶水脚本」、各自有 Bug、各自的 Offset、各自的死信、出问题没人修的同学。
学完你会:能说清「Kafka Connect 是什么、为什么不是『再写一个 Producer/Consumer』就行」;能用 Distributed 模式的 REST API 部署 / 监控 / 重启一个 JDBC Source 或 Debezium CDC Source;能说出 Converter、SMT、DLQ、Exactly Once Source 各自解决什么痛点;能在面试中讲清楚 Connect 与 Flink CDC、自研 Canal 之间的取舍。
0. 一个生活类比:自来水管接头
把 Kafka 想成「自来水主管道」,主管道两头要接很多支管:
- 上游:水库(MySQL)、河流(Web 日志)、井(IoT 设备)—— 都要把水接进主管。
- 下游:洗衣机(ES)、淋浴(数仓)、绿化喷头(监控)—— 都要从主管引水出去。
如果你每接一处都自己用 PVC 焊一段管子(自己写 Producer/Consumer),就会出现:
- 接头各种规格不统一、漏水点没人维护;
- 一处堵了不知道从哪儿排查;
- 想换水源(MySQL → PostgreSQL)要重新焊;
- 加新的喷头要等开发排期。
Kafka Connect 就是一套「标准化的水管接头」:
- 每种数据源 / 落地系统都有现成的 Connector(Source 或 Sink);
- 部署在 Connect Cluster 这个统一的工人小组里,自动负载均衡、自动接管故障;
- 配置统一走 REST API,不需要改业务代码;
- 出错的「漏水」可以通过 DLQ 单独接走,不污染主管道。
记住这个类比,下面所有概念都好理解。
1. 为什么不能「自己写 Producer/Consumer」?
很多读者第一反应:「Source / Sink 不就是 Consumer / Producer 吗,自己写不就行了?」
把所有「自己写」会遇到的坑列一遍,就知道为什么需要 Connect:
| 痛点 | 自己写 | Kafka Connect |
|---|---|---|
| 增量抓取(数据库表只读新数据) | 自己维护 watermark / 时间戳,写到自己的表 | Connector 自带,状态存 connect-offsets Topic |
| 失败重启续传 | 自己写 checkpoint,每次升级都担心 Offset 丢 | 框架统一管 Offset |
| 错误消息处理 | 写一遍死信、自己监控、自己重试 | DLQ 配置一行:errors.deadletterqueue.topic.name |
| Schema 变更 | 业务代码改字段映射、跑迁移脚本 | Schema Registry + Avro 自动演进 |
| 集群伸缩 | 自己写选举 / 心跳 / 任务再平衡 | Distributed Worker 自动 Rebalance |
| 部署 / 监控 | 一个项目一套 Dockerfile / 一份监控 | 统一的 REST API + JMX |
| Exactly Once | 自己实现两阶段提交(很难做对) | Connect 3.3+ 提供 EOS Source |
| 团队扩展 | 来个新需求开个新项目,没人愿意维护「胶水代码」 | 改一份 JSON 配置即可 |
一句话:写 Producer/Consumer 是「一次性焊管」,用 Connect 是「装可换接头」。规模一上去,差别天壤之别。
📌 横向对比:
- Flink CDC:Source 端比 Connect 强(CDC 增量更细),但 Sink 端不如 Connect 生态丰富,且需要 Flink 集群。
- Apache NiFi:可视化拖拽,灵活度极高但学习成本陡,社区比 Connect 小。
- 自研 Canal / Maxwell:只解决 MySQL CDC 这一项,没生态。
- Pulsar IO:Pulsar 的等价物,函数化、轻量,但生态远小于 Connect。
2. 核心概念:Connector / Task / Worker
Connect 的对象模型只有 4 个名词,但它们之间的关系是 Connect 调度的全部秘密:
- Worker:一个 JVM 进程,是 Connect 的运行容器。3 个 Worker = 一个 Connect 集群。
- Connector:一个逻辑抽象,对应一份 JSON 配置(「我要把 MySQL 的 orders 表导到 Kafka」)。它的唯一职责是「告诉 Connect 我要做什么 + 我可以拆成几个并行 Task」。
- Task:Connector 真正干活的子任务。Task 数受
tasks.max限制,也受数据源本身能切的并行度限制(比如 JDBC Source 的 Task 数 ≤ 表的分区数 / 表数)。 - REST API:管理面唯一入口。所有创建 / 暂停 / 重启 / 删除都走 HTTP。
关键不变量:
- 一个 Worker 可以同时跑 N 个 Task(Worker:Task = 1:多)。
- 一个 Task 任意时刻只属于一个 Worker(Task 不会被同时运行)。
- Worker 宕机后,它身上的 Task 会被 Rebalance 到剩下的 Worker(自动 Failover)。
- Connector 配置改了,会触发 Task 重启。
3. Source vs Sink
| 角色 | 数据流向 | 等价于自己写的 | 例子 |
|---|---|---|---|
| Source Connector | 外部系统 → Kafka | Producer | JDBC Source、Debezium MySQL/PG/Mongo CDC、File Source |
| Sink Connector | Kafka → 外部系统 | Consumer | Elasticsearch Sink、S3 Sink、HDFS Sink、JDBC Sink、MongoDB Sink |
简单的一条线:
text
[ MySQL ] -- Debezium Source --> [ Kafka Topic ] -- ES Sink --> [ Elasticsearch ]中间那段 Kafka Topic 是事件总线 + 缓冲,让上下游解耦:MySQL 端短暂宕机不影响 ES 端,ES 升级不影响 MySQL 端。
4. 部署模式:Standalone vs Distributed
4.1 Standalone
单进程跑 Connect,配置写在本地 properties 文件里,Offset 存在本地文件。
bash
bin/connect-standalone.sh \
config/connect-standalone.properties \
config/source-1.properties \
config/source-2.properties特点:
- 简单,适合本地开发 / 单机集成;
- 没有 HA,进程挂了任务全停;
- 配置改了要手动重启进程;
- Offset 存本地文件,迁移机器麻烦。
适用:开发调试、写在边缘节点的轻量同步(IoT 网关)、教学。
4.2 Distributed(生产推荐)
多个 Worker 进程组成一个集群,配置 / Offset / 状态都存在 Kafka 内部 Topic 里。
bash
bin/connect-distributed.sh config/connect-distributed.properties特点:
- 高可用,单 Worker 挂掉自动 Rebalance;
- 配置 / Offset / 状态存在 Kafka,迁移机器只需起新 Worker;
- 通过 REST API 增删 Connector,热更新无需重启。
4.3 Distributed 的三大内部 Topic
Distributed 模式靠这三个 Topic 维持「集群状态」:
| Topic | 默认名 | 作用 | Compaction |
|---|---|---|---|
| 配置 Topic | connect-configs | 存 Connector / Task 配置 | ✅(每个 connector name 取最新) |
| Offset Topic | connect-offsets | 存 Source Connector 已读到的位置(Sink 走 __consumer_offsets) | ✅(每个 partition key 取最新) |
| 状态 Topic | connect-status | 存 Connector / Task 当前状态(RUNNING / PAUSED / FAILED) | ✅ |
connect-distributed.properties 关键项:
properties
bootstrap.servers=broker1:9092,broker2:9092,broker3:9092
group.id=connect-cluster # ★ 同 group.id 的 Worker 自动组成集群
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-status
config.storage.replication.factor=3
offset.storage.replication.factor=3
status.storage.replication.factor=3
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false
plugin.path=/usr/share/java,/opt/connectors # Connector jar 包目录⚠️ 三个内部 Topic 的副本数必须 ≥ 2(生产建议 3),否则 Worker 之间状态不一致后会很惨。 ⚠️ 不要让多个不同业务的 Connect 集群共用
group.id,否则配置 / Offset 会撞车。
5. REST API 操作:Connect 的「方向盘」
Distributed 模式下,所有运维都通过 REST API(默认端口 8083)。
5.1 列出已注册的 Connector
bash
curl -s http://localhost:8083/connectors | jq
# ["jdbc-source-orders","es-sink-orders"]5.2 创建 Connector
bash
curl -X POST http://localhost:8083/connectors \
-H 'Content-Type: application/json' \
-d @jdbc_source_config.jsonjdbc_source_config.json 的最外层是:
json
{
"name": "jdbc-source-orders",
"config": { "connector.class": "...", "tasks.max": "3", ... }
}5.3 查看状态
bash
curl -s http://localhost:8083/connectors/jdbc-source-orders/status | jq返回:
json
{
"name": "jdbc-source-orders",
"connector": {"state": "RUNNING", "worker_id": "10.0.0.1:8083"},
"tasks": [
{"id": 0, "state": "RUNNING", "worker_id": "10.0.0.2:8083"},
{"id": 1, "state": "FAILED", "worker_id": "10.0.0.3:8083",
"trace": "java.sql.SQLException: ..."}
]
}5.4 暂停 / 恢复
bash
curl -X PUT http://localhost:8083/connectors/jdbc-source-orders/pause
curl -X PUT http://localhost:8083/connectors/jdbc-source-orders/resume「暂停」不会丢配置 / Offset,只是停止拉新数据。常用于对源系统做迁移、临时降级。
5.5 重启
bash
# 重启整个 connector + 所有 task
curl -X POST http://localhost:8083/connectors/jdbc-source-orders/restart
# 重启单个 task
curl -X POST http://localhost:8083/connectors/jdbc-source-orders/tasks/1/restart
# 一次重启所有 FAILED task
curl -X POST 'http://localhost:8083/connectors/jdbc-source-orders/restart?includeTasks=true&onlyFailed=true'5.6 修改配置
bash
curl -X PUT http://localhost:8083/connectors/jdbc-source-orders/config \
-H 'Content-Type: application/json' \
-d @new_config.json注意:用 PUT,body 只是 config 部分(不带最外层 name)。
5.7 删除
bash
curl -X DELETE http://localhost:8083/connectors/jdbc-source-orders删除会清掉配置和状态,但不会清掉 connect-offsets 里这个 connector 的 Offset(除非你显式调 /connectors/{name}/offsets 删除)。这点很重要:删了重建会从原 Offset 继续,避免重复消费。
5.8 查看 / 编辑 Offset(3.6+)
bash
curl -s http://localhost:8083/connectors/jdbc-source-orders/offsets | jq
curl -X DELETE http://localhost:8083/connectors/jdbc-source-orders/offsets # 重置完整 REST API 调用集合见
16_connect/code/connect_rest_demo.sh。
6. 常见 Connector 一览
Connector jar 包不内置在 Kafka,需要
plugin.path下放置。Confluent Hub 是最大的 Connector 仓库。
6.1 JDBC Source / Sink(confluentinc/kafka-connect-jdbc)
最常用的「关系数据库 ↔ Kafka」桥梁。
Source 三种模式:
| 模式 | 说明 | 适用 |
|---|---|---|
bulk | 每次都全表扫描 | 字典表、维度表、不大的表 |
incrementing | 按自增 ID 拉新(WHERE id > ?) | 只插入不更新的表 |
timestamp | 按时间戳列拉新(WHERE updated_at > ?) | 有时间戳的表,能感知更新 |
timestamp+incrementing | 复合,最稳 | 推荐:时间戳定位 + ID 兜底 |
最小配置:
json
{
"name": "jdbc-source-orders",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"tasks.max": "3",
"connection.url": "jdbc:mysql://mysql:3306/shop",
"connection.user": "kafka",
"connection.password": "kafka-secret",
"table.whitelist": "orders,users",
"mode": "timestamp+incrementing",
"timestamp.column.name": "updated_at",
"incrementing.column.name": "id",
"topic.prefix": "mysql.shop.",
"poll.interval.ms": "5000"
}
}JDBC Source 的局限:感知不到 DELETE(删了的行你 WHERE > ? 永远扫不到),所以对增删改要求完整的场景,必须上 Debezium。
6.2 Debezium CDC(debezium/debezium-connector-mysql/postgres/mongodb)
「Change Data Capture」—— 直接读数据库的 binlog / WAL / oplog,把每一条增删改作为事件发到 Kafka。
为什么需要 Debezium:
- JDBC Source 感知不到 DELETE;
- 高频更新场景下 JDBC 全表扫描会拖死库;
- 业务字段没有
updated_at列; - 需要严格的「事件次序」(binlog 是写入顺序)。
MySQL 配置(关键):
json
{
"name": "debezium-mysql-shop",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"tasks.max": "1",
"database.hostname": "mysql",
"database.port": "3306",
"database.user": "debezium",
"database.password": "debezium-secret",
"database.server.id": "184054",
"topic.prefix": "dbz.shop",
"database.include.list": "shop",
"table.include.list": "shop.orders,shop.users",
"schema.history.internal.kafka.bootstrap.servers": "broker:9092",
"schema.history.internal.kafka.topic": "dbz.shop.history",
"snapshot.mode": "initial"
}
}MySQL 端要求:
binlog_format = ROWbinlog_row_image = FULL- 给 debezium 用户开
REPLICATION CLIENT、REPLICATION SLAVE、SELECT、RELOAD权限
输出消息样例(dbz.shop.shop.orders):
json
{
"before": null,
"after": {"id": 100, "amount": 99.5, "status": "PAID", "updated_at": "..."},
"source": {"db": "shop", "table": "orders", "ts_ms": 1700000000, ...},
"op": "c", // c=create, u=update, d=delete, r=read(snapshot)
"ts_ms": 1700000000123
}6.3 Elasticsearch Sink
json
{
"name": "es-sink-orders",
"config": {
"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
"tasks.max": "3",
"topics": "dbz.shop.shop.orders",
"connection.url": "http://es-node:9200",
"type.name": "_doc",
"key.ignore": "false",
"schema.ignore": "true",
"behavior.on.malformed.documents": "warn",
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "dlq.es-sink-orders",
"errors.deadletterqueue.context.headers.enable": "true"
}
}注意点:
key.ignore=false:用 Kafka message key 当 ES 文档 ID,实现「同主键覆盖」幂等;schema.ignore=true:不要求 message 带 Connect schema(搭配 JsonConverter + schemas.enable=false)。
6.4 S3 Sink
把 Kafka 数据按时间分区写到 S3,常用作「冷数据归档 / 数仓 ODS 层」。
json
{
"connector.class": "io.confluent.connect.s3.S3SinkConnector",
"tasks.max": "6",
"topics": "dbz.shop.shop.orders",
"s3.region": "us-east-1",
"s3.bucket.name": "data-lake-prod",
"flush.size": "100000",
"rotate.interval.ms": "600000",
"storage.class": "io.confluent.connect.s3.storage.S3Storage",
"format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
"partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
"path.format": "'year'=YYYY/'month'=MM/'day'=dd",
"partition.duration.ms": "86400000",
"timestamp.extractor": "Record"
}6.5 HDFS Sink / MongoDB Sink
类似 S3 Sink,把 Kafka 数据落到 HDFS / Mongo。配置思路一致:连接信息 + Topic 列表 + 写入策略 + DLQ。
7. Converter:消息「翻译机」
Connect 把消息 in/out 时要做 序列化 / 反序列化,由 Converter 负责。
| Converter | 输出格式 | 优点 | 缺点 | 与 Schema Registry |
|---|---|---|---|---|
StringConverter | 字节即字符串 | 最简单 | 没 schema | ❌ |
JsonConverter | JSON | 人类可读、调试友好 | 体积大、字段类型靠猜 | 可选(schemas.enable) |
AvroConverter | Avro | 强 schema、向后兼容、紧凑 | 需 Schema Registry | ✅ 必须 |
ProtobufConverter | Protobuf | 强 schema、最紧凑 | gRPC 友好但工具支持稍弱 | ✅ 必须 |
JsonSchemaConverter | JSON + schema | 折中 | 体积仍较大 | ✅ |
配置区分 key 和 value:
properties
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://schema-registry:8081生产推荐:
- 跨团队 / 强 schema 场景:AvroConverter + Schema Registry;
- 内部小工具 / 调试:JsonConverter(带
schemas.enable=false减小体积); - 高性能要求 / gRPC 生态:ProtobufConverter。
8. SMT(Single Message Transforms)
Connect 在「Source 写出 → Topic」和「Topic → Sink 读入」中间,可以串一组单消息转换器。它的特点是「单条消息无状态变换」(不能跨消息聚合,那是 Streams 的活)。
8.1 常用 SMT 一览
| SMT | 干什么 | 用例 |
|---|---|---|
InsertField | 插入字段(如 source topic name、时间戳) | 加审计字段 |
ReplaceField | 改名 / 删掉字段 | 适配下游 schema |
MaskField | 把字段值置 null 或 0 | 脱敏(手机号、身份证) |
ValueToKey | 把某些 value 字段提到 key | 让 Sink 用业务主键去重 |
ExtractField | 从结构里抽一个子字段当整条 value | 把 Debezium 的 after 提出来 |
Cast | 类型强转 | int → string |
TimestampRouter | 按时间戳路由到不同 topic | 按天分 topic |
RegexRouter | 按正则改 topic 名 | mysql.shop.orders → orders |
Filter | 满足条件的丢弃 | 只保留 op=u(update) |
Flatten | 把嵌套结构拍平 | 嵌套 JSON 转扁平表 |
8.2 SMT 配置示例
把 Debezium 的 after 字段抽出来 + 重命名 topic + 加审计字段:
json
"transforms": "unwrap,route,addTs",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.drop.tombstones": "false",
"transforms.unwrap.delete.handling.mode": "rewrite",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "dbz\\.shop\\.shop\\.(.*)",
"transforms.route.replacement": "shop.$1",
"transforms.addTs.type": "org.apache.kafka.connect.transforms.InsertField$Value",
"transforms.addTs.timestamp.field": "_kafka_ts"效果:原本是 {before, after, source, op, ts_ms} 的 Debezium 消息,被「展平」成业务字段 + _kafka_ts,topic 名也从 dbz.shop.shop.orders 改成了 shop.orders。
SMT 链很方便,但不要塞太多。复杂逻辑请上 Kafka Streams / ksqlDB / Flink,否则故障排查会哭。
9. Dead Letter Queue(DLQ):错误消息的隔离区
问题:Sink 端遇到「无法处理的消息」(schema 不匹配、外部系统报 4xx、SMT 失败……)该怎么办?
- 默认
errors.tolerance=none—— 整个 Task FAILED 卡住,影响后续所有消息。 - 改为
errors.tolerance=all+ 配 DLQ:把烂消息扔到一个单独的 topic,主流量继续跑。
json
"errors.tolerance": "all",
"errors.log.enable": "true",
"errors.log.include.messages": "true",
"errors.deadletterqueue.topic.name": "dlq.es-sink-orders",
"errors.deadletterqueue.context.headers.enable": "true",
"errors.deadletterqueue.topic.replication.factor": "3"DLQ 消息会带 headers,包含原 topic、partition、offset、错误类、错误栈,方便人工排查。
最佳实践:
- 每个生产 Connector 都配 DLQ(哪怕一开始没消息);
- 监控 DLQ 的写入速率,> 0 就告警;
- 写一个「DLQ 重放工具」:人工修复后把消息重新发回主 topic。
10. Exactly Once Source(Connect 3.3+)
历史上 Connect 的语义是 At Least Once:Worker 崩溃 / Rebalance 时,可能重复发送一段消息。3.3 起,Source 端支持 Exactly Once。
10.1 开启方式
connect-distributed.properties:
properties
exactly.once.source.support=enabled # disabled / preparing / enabled切换流程必须先 preparing 后 enabled,让所有 Worker 升级状态:
text
disabled --> preparing --> enabled10.2 实现原理
Connect 把每条 Source 记录放进 事务 Producer,并把「我已经把这批数据写到了 Topic」这个事实和 Source 偏移量 (source offset) 一起塞进同一个事务,commit 到 __transaction_state。Worker 崩溃 → 接管的 Worker 拿到上次提交的 source offset,从下条开始读。
内部其实就是把第 13 章 EOS 的事务 Producer 封装到框架里,业务方写 Connector 不用感知。
10.3 限制
- 仅 Source 端,Sink 端的 EOS 仍要靠下游系统的幂等(如 ES 用 docId 覆盖);
- 性能比 At Least Once 略低(事务有 commit 开销);
- Connector 必须实现
SourceConnector接口的exactlyOnceSupport()方法(社区主流 Connector 都已支持)。
11. 端到端示例:MySQL → Kafka → Elasticsearch
来一个完整链路,把前面的概念串起来。
11.1 拓扑
11.2 部署步骤
bash
# 1. 启动 Connect Cluster
./start_connect.sh
# 2. 注册 Debezium Source
curl -X POST http://localhost:8083/connectors \
-H 'Content-Type: application/json' \
-d @debezium_mysql_config.json
# 3. 注册 ES Sink
curl -X POST http://localhost:8083/connectors \
-H 'Content-Type: application/json' \
-d @elastic_sink_config.json
# 4. 验证
curl -s http://localhost:8083/connectors?expand=status | jq11.3 链路特性
| 环节 | 故障表现 | 自愈机制 |
|---|---|---|
| Debezium Worker 挂 | Source Task FAILED | 另一个 Worker 接管,从 binlog position 续读 |
| Kafka Broker 挂 | Source/Sink 重试 | 副本切 Leader,连自动重连 |
| ES 短暂宕机 | Sink 重试,Lag 积压 | ES 起来后自动追平 |
| 某条消息 schema 坏 | 进 DLQ,Task 不挂 | 人工修复后重放 |
这套链路的写量级:单 Connect Worker(4 核 8G)+ 3 partition 的 ES Sink,可稳定支撑 ~10 万 TPS 的「PAID 订单实时入 ES」场景。
12. 与其他 MQ 集成方案对比
| 方案 | 优点 | 缺点 | 适用 |
|---|---|---|---|
| Kafka Connect | 生态丰富、运维统一、HA 自带 | JVM 依赖、内存大 | Kafka 生态全栈 |
| Flink CDC | 流处理 + CDC 一体、Exactly Once | 需 Flink 集群、Sink 不如 Connect 多 | 实时数仓 |
| Apache NiFi | 拖拽 UI、灵活度极高 | 学习陡、社区小 | 企业 ETL |
| Logstash / Fluentd | 轻量、日志生态强 | 不适合事务性数据 | 日志收集 |
| Airbyte | 现代 UI、SaaS 模型 | 实时性弱(更像批) | 数据仓库 ELT |
| Debezium 直接当 lib | 无 Connect 框架开销 | 自己写部署 / HA | 嵌入式 CDC |
13. 生产踩坑
坑 1:tasks.max 设太大
JDBC Source 表只有 1 张,tasks.max=10 也只会起 1 个 Task(受源切片能力限制)。但 Sink 端 tasks.max 不能 > 订阅 Topic 总分区数,否则多余 Task 空转。
坑 2:connect-offsets Topic 不要手动改
Source Offset 错乱 → 数据重复 / 丢失。改之前必须先 stop connector,并用官方 REST /offsets 接口。
坑 3:JDBC Source 感知不到 DELETE
如开头提到,要严格 CDC 必须 Debezium。
坑 4:Schema 演进破坏 Sink
下游 ES 已经按旧字段建好 mapping,Source 端突然多字段 → ES 写入报错。用 Schema Registry + BACKWARD 策略卡住(详见第 18 章)。
坑 5:DLQ 用同一个 connector 自己消费
DLQ 出问题再进 DLQ → 死循环。DLQ 需要独立 Connector(甚至独立 Connect 集群)处理。
坑 6:plugin.path 路径问题
Connector jar 必须放在 plugin.path 指定目录的子目录里(每个 connector 一个子目录),不是直接平铺。错误的目录结构会让 ClassNotFound 或 Connector 加载不到。
坑 7:Distributed 集群里 group.id 撞车
两个不同业务的 Connect 集群用了相同 group.id → 配置 / Offset 互相覆盖,灾难现场。group.id 必须每集群唯一。
坑 8:内部 Topic 副本数 1
config.storage.replication.factor=1 —— Worker 重启时配置丢一半。必须 ≥ 3。
14. 小结
- Connect = Kafka 与外部系统之间的「标准化水管接头」,避免每个集成点都自己写 Producer/Consumer。
- Source / Sink 两类 Connector,Connector / Task / Worker 三级对象模型。
- Standalone vs Distributed:生产必选 Distributed(HA + REST API + 热更新)。
- 三大内部 Topic:
connect-configs/connect-offsets/connect-status,副本数 ≥ 3。 - REST API 是唯一管理入口:创建 / 暂停 / 恢复 / 重启 / 删除 / 改配置。
- Converter(JSON / Avro / Protobuf)+ SMT(InsertField / ReplaceField / MaskField / Filter / Router)+ DLQ(errors.tolerance=all + deadletterqueue.topic.name),是 Connect 三大「易用性利器」。
- Connect 3.3+ 提供 Exactly Once Source,业务方无感知。
- 生产部署口诀:「Distributed + 副本 3 + DLQ + Schema Registry + 监控 task state」。
🎯 面试高频题
Q1:Kafka Connect 是什么?为什么不直接写 Producer/Consumer?
考察点:架构理解、工程权衡。
答案:
- 定位:Kafka 与外部系统的「标准化数据桥」,把「连什么、怎么连、状态在哪、错了怎么办」抽象成 Connector + 配置。
- 解决的痛点:增量抓取、Offset 管理、错误处理(DLQ)、Schema 演进、集群 HA、统一监控、Exactly Once,每一项自己写都是大工程。
- 核心对象:Connector(逻辑配置)/ Task(实际执行单元)/ Worker(JVM 进程)。
- 部署模式:Standalone(单机调试)/ Distributed(生产标配,HA + REST API + 配置热更新)。
- 替代方案对比:Flink CDC(强 Source 弱 Sink、需 Flink 集群);NiFi(UI 灵活、学习陡);自研(重复造轮子)。
- 加分项:提到 Distributed 集群靠
connect-configs/connect-offsets/connect-status三个内部 Topic 维持状态;提到 Connect 3.3+ 的 Exactly Once Source。
Q2:Source Connector 和 Sink Connector 的区别?两端的 Offset 分别存在哪?
考察点:模型本质、Offset 管理。
答案:
- Source:外部系统 → Kafka。它「等价于一个 Producer」。Offset = 「我从源系统读到了哪里」(如 binlog position、JDBC 时间戳),存在
connect-offsets这个内部 Topic(compaction)。 - Sink:Kafka → 外部系统。它「等价于一个 Consumer」。Offset = 「我从 Kafka 消费到了哪里」,存在
__consumer_offsets,和普通 Consumer 走完全相同的提交路径。 - 不要混淆:Sink 的 Offset 跟 Source 完全是两个概念,存储位置都不同。
- 影响:删除 Source Connector 后再创建同名的,会从
connect-offsets续读;删除 Sink Connector 再创建同名的,会从__consumer_offsets续读(按 group.id)。 - 加分项:提到 3.6+ 的
/connectors/{name}/offsetsREST API 可以查看 / 修改 / 删除 Offset;提到 EOS Source 利用事务把 source offset 和 produce 写入捆绑提交。
Q3:SMT 是什么?能干什么不能干什么?举几个常用 SMT。
考察点:Connect 灵活性、流处理边界。
答案:
- SMT = Single Message Transform,Source 写出后 / Sink 读入前的「单消息无状态变换」流水线。
- 能干:改字段名(ReplaceField)、加字段(InsertField)、删字段、脱敏(MaskField)、按时间戳路由 topic(TimestampRouter)、按正则换 topic 名(RegexRouter)、过滤条件(Filter)、提取嵌套字段(ExtractField)、类型强转(Cast)、Debezium 解包(ExtractNewRecordState)。
- 不能干:跨消息聚合 / Join / 窗口(这些是 Kafka Streams / ksqlDB 的活);调用外部 API(也不该在 SMT 里做)。
- 典型 Debezium SMT 链:
unwrap(提 after)→RegexRouter(重命名 topic)→InsertField(加审计字段)。 - 加分项:提到 SMT 越多越难维护,复杂转换上 Streams / ksqlDB;提到 SMT 顺序很重要(Filter 放越早越省钱);提到自定义 SMT 只需实现
Transformation<R>接口。
Q4:Connect 的 DLQ 怎么配?为什么生产必须配?
考察点:错误处理、生产可用性。
答案:
- 配置:json
"errors.tolerance": "all", "errors.deadletterqueue.topic.name": "dlq.es-sink-orders", "errors.deadletterqueue.context.headers.enable": "true", "errors.deadletterqueue.topic.replication.factor": "3" - 不配的后果:默认
errors.tolerance=none,遇到一条 schema 错的消息就把整个 Task 弄成 FAILED,整条管道卡住,所有后续消息堆积。 - DLQ headers:Kafka header 里包含原 topic / partition / offset / 错误类 / 错误栈,方便人工排查。
- 生产范式:
- 每个 Sink Connector 都配 DLQ;
- 监控 DLQ 写入速率,> 0 告警;
- 提供「DLQ 重放工具」:修复后回灌主 topic;
- DLQ 的处理 Connector / Consumer 不能再写自己 → 死循环。
- 加分项:提到 Source Connector 也可以配
errors.tolerance,但 Source 的「无法转换」一般是 schema 库的问题,更要谨慎;提到 OPA / Sentinel 可以做更复杂的「策略式错误处理」。
Q5:Connect 的 Exactly Once 是怎么实现的?跟普通 Producer 的事务有什么关系?
考察点:EOS 协议、Connect 内部实现。
答案:
- 范围:3.3+ 仅支持 Source EOS,Sink EOS 仍依赖下游系统幂等(如 ES 用 docId 覆盖、JDBC 用 PK 主键 ON DUPLICATE KEY)。
- 实现:Worker 把每个 poll 周期产出的消息放进事务 Producer,并用
sendOffsetsToTransaction把 source offset(写到connect-offsets)和业务消息绑定提交。 - 效果:commit 成功 → 消息 + offset 同时可见;中途崩溃 → 事务被 abort,下次 Worker 接管时从上次 commit 的 offset 重新开始,业务 Topic 看不到中间未提交的消息。
- 配置:
exactly.once.source.support=enabled,必须先preparing再enabled。 - 限制:性能略低(事务 commit 有开销);Connector 必须实现
exactlyOnceSupport()(主流社区 Connector 都已支持);fencing 机制依赖 transactional.id,Worker 数与 connector 配置变更会触发 epoch++。 - 加分项:提到这就是把第 13 章「事务 Producer + read_committed」的能力封装到框架里;提到下游消费者必须
isolation.level=read_committed才能看不到 abort 的消息。
Q6:你们生产的 MySQL → Kafka → ES 链路是怎么搭的?遇到什么坑?
考察点:实战经验、问题排查。
答案:
- 架构:Debezium MySQL Source → Kafka Topic → SMT (unwrap + RegexRouter + InsertField) → ES Sink + DLQ。
- 关键配置:
- MySQL 端
binlog_format=ROW、binlog_row_image=FULL; - Debezium snapshot.mode=initial;
- ES Sink
key.ignore=false(用 PK 当 docId 实现幂等覆盖); - Schema Registry + Avro 卡住兼容性。
- MySQL 端
- 常见坑:
- 表无主键 → ES 文档 ID 错乱 → 加
ValueToKeySMT; - 大事务一次推几万行 → ES 写入限速 → Sink 端
linger.ms+batch.size调优; - 某些字段 NULL / 类型变更 → 进 DLQ → 人工修 mapping 后重放;
- Worker OOM → 调
-Xmx并把 Sink batch 调小; - Schema 突然变 → ES mapping 拒收 → 通过 Schema Registry 卡住兼容性。
- 表无主键 → ES 文档 ID 错乱 → 加
- 观测:JMX 监控
connector-metrics(state)、task-metrics(offset-commit-failure-percentage)、sink-task-metrics(sink-record-send-rate);DLQ 速率监控。 - 加分项:提到对比 Flink CDC 的取舍(Connect 部署轻、Flink 流处理强);提到一致性等级(Source EOS + Sink At Least Once + ES 幂等覆盖 = 端到端 Effectively Once)。
本章配套:
16_connect/demo.html:Connect 集群拓扑 + Connector 状态机 + DLQ 流转动画 + SMT 前后对比。16_connect/code/start_connect.sh:本地启动 Distributed Connect 脚本。16_connect/code/jdbc_source_config.json:JDBC Source 配置。16_connect/code/debezium_mysql_config.json:Debezium MySQL Source 配置。16_connect/code/elastic_sink_config.json:ES Sink 配置。16_connect/code/connect_rest_demo.sh:REST API 操作集合。
🔗 延伸阅读
- 第 13 章 幂等与事务 —— Connect EOS Source 复用的事务 Producer 协议。
- 第 15 章 安全与多租户 —— Connect Worker 的 SuperUser / SASL 配置。
- 第 17 章 Streams 与 ksqlDB —— SMT 不够用时升级到流处理。
- 第 18 章 Schema Registry —— Avro Converter 的兼容性策略。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
bash
#!/usr/bin/env bash
# ====================================================================
# 第 16 章 - Kafka Connect
# connect_rest_demo.sh - REST API 操作集合(curl + jq)
# --------------------------------------------------------------------
# 默认 Connect REST 端口 8083;如需鉴权,加 -u user:pass
# ====================================================================
set -euo pipefail
CONNECT=${CONNECT:-http://localhost:8083}
NAME=${NAME:-jdbc-source-orders}
hr() { printf '\n\033[36m%s\033[0m\n' "==> $*"; }
hr "0. 集群信息"
curl -s "$CONNECT/" | jq
hr "0.1 列出所有可用 Connector Plugin(看你装了哪些 Connector jar)"
curl -s "$CONNECT/connector-plugins" | jq '.[] | {class, type, version}'
hr "1. 列出所有已注册的 Connector"
curl -s "$CONNECT/connectors" | jq
hr "1.1 列出 Connector + 状态(一行命令看全集群)"
curl -s "$CONNECT/connectors?expand=info&expand=status" | jq
hr "2. 创建 Connector(从 jdbc_source_config.json 加载)"
echo "(示例:curl -X POST $CONNECT/connectors -H 'Content-Type: application/json' -d @jdbc_source_config.json)"
hr "3. 查看某 Connector 的配置"
curl -s "$CONNECT/connectors/$NAME/config" | jq
hr "4. 查看状态(重点看每个 task 是否 RUNNING)"
curl -s "$CONNECT/connectors/$NAME/status" | jq
hr "5. 查看 Source Offset(3.6+ 才有)"
curl -s "$CONNECT/connectors/$NAME/offsets" | jq || echo "(若 404 说明 Connect 版本 < 3.6)"
hr "6. 暂停"
curl -X PUT "$CONNECT/connectors/$NAME/pause"
hr "7. 恢复"
curl -X PUT "$CONNECT/connectors/$NAME/resume"
hr "8. 重启整个 Connector + 所有 task"
curl -X POST "$CONNECT/connectors/$NAME/restart?includeTasks=true&onlyFailed=false" | jq
hr "8.1 仅重启 FAILED 的 task"
curl -X POST "$CONNECT/connectors/$NAME/restart?includeTasks=true&onlyFailed=true" | jq
hr "8.2 重启单个 task(id=0)"
curl -X POST "$CONNECT/connectors/$NAME/tasks/0/restart"
hr "9. 修改配置(PUT,body 只是 config 部分)"
echo "示例:"
cat <<EOF
curl -X PUT $CONNECT/connectors/$NAME/config \\
-H 'Content-Type: application/json' \\
-d '{
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"tasks.max": "5",
"...其它原有配置": "..."
}'
EOF
hr "10. 查看 task 列表 + 配置"
curl -s "$CONNECT/connectors/$NAME/tasks" | jq
hr "11. 验证配置(不真正提交,仅校验)"
echo "示例:"
cat <<EOF
curl -X PUT $CONNECT/connector-plugins/io.confluent.connect.jdbc.JdbcSourceConnector/config/validate \\
-H 'Content-Type: application/json' \\
-d @jdbc_source_config.json
EOF
hr "12. 删除 Connector(仅删配置 + 状态,Offset 保留在 connect-offsets)"
echo "(危险!) curl -X DELETE $CONNECT/connectors/$NAME"
hr "12.1 重置 Source Offset(3.6+,先删后建可让 connector 从头读)"
echo "curl -X DELETE $CONNECT/connectors/$NAME/offsets"
hr "13. 集群 logger 级别(在线改 log4j 等级,无需重启)"
curl -s "$CONNECT/admin/loggers" | jq
echo "改某 package 等级:"
echo "curl -X PUT $CONNECT/admin/loggers/io.debezium -H 'Content-Type: application/json' -d '{\"level\":\"DEBUG\"}'"
# ====================================================================
# 一行命令快速排障:所有 connector 哪个 task 在 FAILED 状态?
# ====================================================================
echo
hr "🔧 Bonus:一行排障 - 列出所有 FAILED task"
curl -s "$CONNECT/connectors?expand=status" | \
jq -r 'to_entries[] |
.key as $name |
.value.status.tasks[] |
select(.state == "FAILED") |
"\($name)/task-\(.id) worker=\(.worker_id)\n trace: \((.trace // "")[0:200])"'json
{
"name": "debezium-mysql-shop",
"config": {
"_doc": "Debezium MySQL Source - 直接读 binlog,捕获 INSERT/UPDATE/DELETE 全增量",
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"tasks.max": "1",
"_db_doc": "MySQL 必须开启 binlog_format=ROW + binlog_row_image=FULL",
"database.hostname": "mysql",
"database.port": "3306",
"database.user": "debezium",
"database.password": "debezium-secret-CHANGE-ME",
"database.server.id": "184054",
"database.connectionTimeZone": "UTC",
"_topic_doc": "Debezium 会发到 ${topic.prefix}.${db}.${table}",
"topic.prefix": "dbz.shop",
"_filter_doc": "只采集 shop 库的 orders/users 表",
"database.include.list": "shop",
"table.include.list": "shop.orders,shop.users",
"_schema_history_doc": "DDL 变更历史存到独立 Kafka Topic(不能和业务 topic 混)",
"schema.history.internal.kafka.bootstrap.servers": "broker:9092",
"schema.history.internal.kafka.topic": "dbz.shop.history",
"_snapshot_doc": "首次启动做快照(initial);后续仅读 binlog(schema_only / never)",
"snapshot.mode": "initial",
"snapshot.locking.mode": "minimal",
"_tombstone_doc": "DELETE 时除了 op=d 的消息,还发一条 value=null 的墓碑(compaction 友好)",
"tombstones.on.delete": "true",
"_smt_doc": "解 Debezium 包:把 envelope (before/after/op/source) 平展成业务 schema",
"transforms": "unwrap,route,addTs",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.drop.tombstones": "false",
"transforms.unwrap.delete.handling.mode": "rewrite",
"transforms.unwrap.add.fields": "op,table,source.ts_ms",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "dbz\\.shop\\.shop\\.(.*)",
"transforms.route.replacement": "shop.$1",
"transforms.addTs.type": "org.apache.kafka.connect.transforms.InsertField$Value",
"transforms.addTs.timestamp.field": "_kafka_ts",
"_converter_doc": "强 schema 场景推荐 Avro + Schema Registry,演示这里仍用 JSON",
"key.converter": "org.apache.kafka.connect.json.JsonConverter",
"key.converter.schemas.enable": "false",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"_error_doc": "对 schema 异常的消息进入 DLQ,主管道不受影响",
"errors.tolerance": "all",
"errors.log.enable": "true",
"errors.log.include.messages": "true",
"errors.deadletterqueue.topic.name": "dlq.debezium-mysql-shop",
"errors.deadletterqueue.topic.replication.factor": "3"
}
}json
{
"name": "es-sink-orders",
"config": {
"_doc": "Elasticsearch Sink - 把 shop.orders Topic 的消息以 PK 幂等覆盖到 ES",
"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
"tasks.max": "3",
"_topic_doc": "订阅一个或多个 topic(topic == ES index 名,自动小写)",
"topics": "shop.orders",
"_es_doc": "ES 集群连接,可配置多节点用逗号隔开",
"connection.url": "http://es-node:9200",
"connection.username": "elastic",
"connection.password": "elastic-secret-CHANGE-ME",
"_doc_id_doc": "用 message key 作 ES 文档 ID,实现『同主键覆盖』幂等",
"key.ignore": "false",
"schema.ignore": "true",
"_write_mode_doc": "INSERT vs UPSERT - 一般 ES 用 upsert(不存在就新建,存在就部分更新)",
"write.method": "upsert",
"_bulk_doc": "ES bulk 写参数 - 调大 batch 提吞吐",
"batch.size": "2000",
"max.in.flight.requests": "5",
"linger.ms": "1000",
"flush.timeout.ms": "30000",
"max.buffered.records": "20000",
"max.retries": "5",
"retry.backoff.ms": "1000",
"_delete_doc": "value=null 的墓碑消息 → ES DELETE",
"behavior.on.null.values": "delete",
"behavior.on.malformed.documents": "warn",
"_converter_doc": "和 Source 的 Converter 对齐",
"key.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"key.converter.schemas.enable": "false",
"value.converter.schemas.enable": "false",
"_error_doc": "宽松错误:所有 ES 4xx 单行错误丢 DLQ;只有 5xx 类系统错才 fail",
"errors.tolerance": "all",
"errors.log.enable": "true",
"errors.log.include.messages": "true",
"errors.deadletterqueue.topic.name": "dlq.es-sink-orders",
"errors.deadletterqueue.topic.replication.factor": "3",
"errors.deadletterqueue.context.headers.enable": "true",
"_consumer_doc": "consumer 端的 isolation.level 必须 read_committed 才能配合 Source EOS",
"consumer.override.isolation.level": "read_committed",
"consumer.override.max.poll.records": "500"
}
}json
{
"name": "jdbc-source-orders",
"config": {
"_doc": "JDBC Source - 增量按 (timestamp + incrementing) 拉取 MySQL shop 库",
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"tasks.max": "3",
"_connection_doc": "JDBC URL / 用户 / 密码(生产建议放 Vault 用 ConfigProvider 注入)",
"connection.url": "jdbc:mysql://mysql:3306/shop?useSSL=false&serverTimezone=UTC",
"connection.user": "kafka",
"connection.password": "kafka-secret-CHANGE-ME",
"_table_doc": "白名单:仅同步 orders 和 users 两张表",
"table.whitelist": "orders,users",
"catalog.pattern": "shop",
"_mode_doc": "复合模式:先按 updated_at 时间戳拉新,相同 updated_at 时按 id 兜底",
"mode": "timestamp+incrementing",
"timestamp.column.name": "updated_at",
"incrementing.column.name": "id",
"_topic_doc": "topic 名 = topic.prefix + 表名,如 mysql.shop.orders",
"topic.prefix": "mysql.shop.",
"_pace_doc": "每 5s 轮询一次源数据库;单批最多 1000 行",
"poll.interval.ms": "5000",
"batch.max.rows": "1000",
"_dq_doc": "Schema 模式(建议 +Schema Registry + AvroConverter)",
"numeric.mapping": "best_fit",
"validate.non.null": "false",
"_converter_doc": "本 Connector 单独覆盖 Worker 默认 Converter",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"_smt_doc": "用 ValueToKey 把 id 字段提到 key,让下游 Sink 按主键幂等",
"transforms": "createKey,extractInt",
"transforms.createKey.type": "org.apache.kafka.connect.transforms.ValueToKey",
"transforms.createKey.fields": "id",
"transforms.extractInt.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
"transforms.extractInt.field": "id",
"_error_doc": "至少容忍 schema 错误,进 DLQ;其它错误(如 SQL 异常)仍 fail",
"errors.tolerance": "all",
"errors.log.enable": "true",
"errors.log.include.messages": "true",
"errors.deadletterqueue.topic.name": "dlq.jdbc-source-orders",
"errors.deadletterqueue.topic.replication.factor": "3",
"errors.deadletterqueue.context.headers.enable": "true"
}
}bash
#!/usr/bin/env bash
# ====================================================================
# 第 16 章 - Kafka Connect
# start_connect.sh - 本地启动 Distributed Connect Worker
# --------------------------------------------------------------------
# 适用:开发 / 调试 / 教学。生产请用 systemd / k8s 管理 worker 进程。
# ====================================================================
set -euo pipefail
KAFKA_HOME=${KAFKA_HOME:-/opt/kafka}
PLUGIN_PATH=${PLUGIN_PATH:-/opt/kafka/connectors}
BOOTSTRAP=${BOOTSTRAP:-localhost:9092}
WORKER_REST_PORT=${WORKER_REST_PORT:-8083}
CONFIG_FILE=$(mktemp /tmp/connect-distributed.XXXXXX.properties)
trap 'rm -f "$CONFIG_FILE"' EXIT
cat > "$CONFIG_FILE" <<EOF
# ===== Kafka 集群 =====
bootstrap.servers=${BOOTSTRAP}
# ===== Connect Cluster =====
group.id=connect-cluster-learn
rest.host.name=0.0.0.0
rest.port=${WORKER_REST_PORT}
rest.advertised.host.name=$(hostname -i 2>/dev/null || echo localhost)
rest.advertised.port=${WORKER_REST_PORT}
# ===== 三大内部 Topic(生产副本数 >= 3) =====
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-status
config.storage.replication.factor=1
offset.storage.replication.factor=1
status.storage.replication.factor=1
offset.storage.partitions=25
status.storage.partitions=5
# ===== Converter 默认 JSON(不带 schema,节省体积) =====
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false
# 内部 Converter(Worker 之间通信,必须 JSON)
internal.key.converter=org.apache.kafka.connect.json.JsonConverter
internal.value.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter.schemas.enable=false
internal.value.converter.schemas.enable=false
# ===== Connector jar 包目录(每个 connector 一个子目录) =====
plugin.path=${PLUGIN_PATH}
# ===== Source EOS(Connect 3.3+,可选) =====
# exactly.once.source.support=enabled
# ===== 默认 producer / consumer 配置(可被 connector override) =====
producer.acks=all
producer.compression.type=zstd
producer.linger.ms=10
consumer.fetch.min.bytes=1
consumer.fetch.max.wait.ms=500
EOF
echo "[*] Connect distributed worker 启动中..."
echo "[*] REST API: http://localhost:${WORKER_REST_PORT}"
echo "[*] plugin.path = ${PLUGIN_PATH}"
echo "[*] 临时 config: ${CONFIG_FILE}"
echo
exec "${KAFKA_HOME}/bin/connect-distributed.sh" "${CONFIG_FILE}"
# ====================================================================
# 启动后用以下命令验证:
# curl -s http://localhost:8083/ | jq
# curl -s http://localhost:8083/connector-plugins | jq # 列出可用 plugin
# curl -s http://localhost:8083/connectors # 当前已注册 connector
# ====================================================================connect_rest_demo.sh ↗ · debezium_mysql_config.json ↗ · elastic_sink_config.json ↗ · jdbc_source_config.json ↗ · start_connect.sh ↗