Skip to content

第 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。

关键不变量

  1. 一个 Worker 可以同时跑 N 个 Task(Worker:Task = 1:多)。
  2. 一个 Task 任意时刻只属于一个 Worker(Task 不会被同时运行)。
  3. Worker 宕机后,它身上的 Task 会被 Rebalance 到剩下的 Worker(自动 Failover)。
  4. Connector 配置改了,会触发 Task 重启。

3. Source vs Sink

角色数据流向等价于自己写的例子
Source Connector外部系统 → KafkaProducerJDBC Source、Debezium MySQL/PG/Mongo CDC、File Source
Sink ConnectorKafka → 外部系统ConsumerElasticsearch 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
配置 Topicconnect-configs存 Connector / Task 配置✅(每个 connector name 取最新)
Offset Topicconnect-offsets存 Source Connector 已读到的位置(Sink 走 __consumer_offsets✅(每个 partition key 取最新)
状态 Topicconnect-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.json

jdbc_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 = ROW
  • binlog_row_image = FULL
  • 给 debezium 用户开 REPLICATION CLIENTREPLICATION SLAVESELECTRELOAD 权限

输出消息样例(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
JsonConverterJSON人类可读、调试友好体积大、字段类型靠猜可选(schemas.enable)
AvroConverterAvro强 schema、向后兼容、紧凑需 Schema Registry✅ 必须
ProtobufConverterProtobuf强 schema、最紧凑gRPC 友好但工具支持稍弱✅ 必须
JsonSchemaConverterJSON + 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.ordersorders
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、错误类、错误栈,方便人工排查。

最佳实践

  1. 每个生产 Connector 都配 DLQ(哪怕一开始没消息);
  2. 监控 DLQ 的写入速率,> 0 就告警;
  3. 写一个「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 --> enabled

10.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 | jq

11.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?

考察点:架构理解、工程权衡。

答案

  1. 定位:Kafka 与外部系统的「标准化数据桥」,把「连什么、怎么连、状态在哪、错了怎么办」抽象成 Connector + 配置。
  2. 解决的痛点:增量抓取、Offset 管理、错误处理(DLQ)、Schema 演进、集群 HA、统一监控、Exactly Once,每一项自己写都是大工程。
  3. 核心对象:Connector(逻辑配置)/ Task(实际执行单元)/ Worker(JVM 进程)。
  4. 部署模式:Standalone(单机调试)/ Distributed(生产标配,HA + REST API + 配置热更新)。
  5. 替代方案对比:Flink CDC(强 Source 弱 Sink、需 Flink 集群);NiFi(UI 灵活、学习陡);自研(重复造轮子)。
  6. 加分项:提到 Distributed 集群靠 connect-configs / connect-offsets / connect-status 三个内部 Topic 维持状态;提到 Connect 3.3+ 的 Exactly Once Source。

Q2:Source Connector 和 Sink Connector 的区别?两端的 Offset 分别存在哪?

考察点:模型本质、Offset 管理。

答案

  1. Source:外部系统 → Kafka。它「等价于一个 Producer」。Offset = 「我从源系统读到了哪里」(如 binlog position、JDBC 时间戳),存在 connect-offsets 这个内部 Topic(compaction)。
  2. Sink:Kafka → 外部系统。它「等价于一个 Consumer」。Offset = 「我从 Kafka 消费到了哪里」,存在 __consumer_offsets,和普通 Consumer 走完全相同的提交路径。
  3. 不要混淆:Sink 的 Offset 跟 Source 完全是两个概念,存储位置都不同。
  4. 影响:删除 Source Connector 后再创建同名的,会从 connect-offsets 续读;删除 Sink Connector 再创建同名的,会从 __consumer_offsets 续读(按 group.id)。
  5. 加分项:提到 3.6+ 的 /connectors/{name}/offsets REST API 可以查看 / 修改 / 删除 Offset;提到 EOS Source 利用事务把 source offset 和 produce 写入捆绑提交。

Q3:SMT 是什么?能干什么不能干什么?举几个常用 SMT。

考察点:Connect 灵活性、流处理边界。

答案

  1. SMT = Single Message Transform,Source 写出后 / Sink 读入前的「单消息无状态变换」流水线。
  2. 能干:改字段名(ReplaceField)、加字段(InsertField)、删字段、脱敏(MaskField)、按时间戳路由 topic(TimestampRouter)、按正则换 topic 名(RegexRouter)、过滤条件(Filter)、提取嵌套字段(ExtractField)、类型强转(Cast)、Debezium 解包(ExtractNewRecordState)。
  3. 不能干:跨消息聚合 / Join / 窗口(这些是 Kafka Streams / ksqlDB 的活);调用外部 API(也不该在 SMT 里做)。
  4. 典型 Debezium SMT 链unwrap(提 after)→ RegexRouter(重命名 topic)→ InsertField(加审计字段)。
  5. 加分项:提到 SMT 越多越难维护,复杂转换上 Streams / ksqlDB;提到 SMT 顺序很重要(Filter 放越早越省钱);提到自定义 SMT 只需实现 Transformation<R> 接口。

Q4:Connect 的 DLQ 怎么配?为什么生产必须配?

考察点:错误处理、生产可用性。

答案

  1. 配置
    json
    "errors.tolerance": "all",
    "errors.deadletterqueue.topic.name": "dlq.es-sink-orders",
    "errors.deadletterqueue.context.headers.enable": "true",
    "errors.deadletterqueue.topic.replication.factor": "3"
  2. 不配的后果:默认 errors.tolerance=none,遇到一条 schema 错的消息就把整个 Task 弄成 FAILED,整条管道卡住,所有后续消息堆积。
  3. DLQ headers:Kafka header 里包含原 topic / partition / offset / 错误类 / 错误栈,方便人工排查。
  4. 生产范式
    • 每个 Sink Connector 都配 DLQ;
    • 监控 DLQ 写入速率,> 0 告警;
    • 提供「DLQ 重放工具」:修复后回灌主 topic;
    • DLQ 的处理 Connector / Consumer 不能再写自己 → 死循环。
  5. 加分项:提到 Source Connector 也可以配 errors.tolerance,但 Source 的「无法转换」一般是 schema 库的问题,更要谨慎;提到 OPA / Sentinel 可以做更复杂的「策略式错误处理」。

Q5:Connect 的 Exactly Once 是怎么实现的?跟普通 Producer 的事务有什么关系?

考察点:EOS 协议、Connect 内部实现。

答案

  1. 范围:3.3+ 仅支持 Source EOS,Sink EOS 仍依赖下游系统幂等(如 ES 用 docId 覆盖、JDBC 用 PK 主键 ON DUPLICATE KEY)。
  2. 实现:Worker 把每个 poll 周期产出的消息放进事务 Producer,并用 sendOffsetsToTransaction 把 source offset(写到 connect-offsets)和业务消息绑定提交。
  3. 效果:commit 成功 → 消息 + offset 同时可见;中途崩溃 → 事务被 abort,下次 Worker 接管时从上次 commit 的 offset 重新开始,业务 Topic 看不到中间未提交的消息。
  4. 配置exactly.once.source.support=enabled,必须先 preparingenabled
  5. 限制:性能略低(事务 commit 有开销);Connector 必须实现 exactlyOnceSupport()(主流社区 Connector 都已支持);fencing 机制依赖 transactional.id,Worker 数与 connector 配置变更会触发 epoch++。
  6. 加分项:提到这就是把第 13 章「事务 Producer + read_committed」的能力封装到框架里;提到下游消费者必须 isolation.level=read_committed 才能看不到 abort 的消息。

Q6:你们生产的 MySQL → Kafka → ES 链路是怎么搭的?遇到什么坑?

考察点:实战经验、问题排查。

答案

  1. 架构:Debezium MySQL Source → Kafka Topic → SMT (unwrap + RegexRouter + InsertField) → ES Sink + DLQ。
  2. 关键配置
    • MySQL 端 binlog_format=ROWbinlog_row_image=FULL
    • Debezium snapshot.mode=initial;
    • ES Sink key.ignore=false(用 PK 当 docId 实现幂等覆盖);
    • Schema Registry + Avro 卡住兼容性。
  3. 常见坑
    • 表无主键 → ES 文档 ID 错乱 → 加 ValueToKey SMT;
    • 大事务一次推几万行 → ES 写入限速 → Sink 端 linger.ms + batch.size 调优;
    • 某些字段 NULL / 类型变更 → 进 DLQ → 人工修 mapping 后重放;
    • Worker OOM → 调 -Xmx 并把 Sink batch 调小;
    • Schema 突然变 → ES mapping 拒收 → 通过 Schema Registry 卡住兼容性。
  4. 观测:JMX 监控 connector-metrics (state)、task-metrics (offset-commit-failure-percentage)、sink-task-metrics (sink-record-send-rate);DLQ 速率监控。
  5. 加分项:提到对比 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 操作集合。

🔗 延伸阅读

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

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 ↗