Learn
Kafka/14-kafka-connect

Kafka Connect 数据集成

"把 MySQL 的数据同步进 Kafka"、"把 topic 数据写进 Elasticsearch"——这类搬运需求如果都手写 Producer/Consumer,你会重复实现断点续传、失败重试、扩缩容、offset 管理。Kafka Connect 就是把这些通用逻辑框架化:写配置,不写代码,即可搭建可扩展、容错的数据管道。

1. 架构:Worker / Connector / Task

                    Connect 集群
  +---------------------------------------------------+
  | Worker 1            | Worker 2                    |
  |  Connector A (配置)  |                             |
  |  ├ Task A-0         |  ├ Task A-1                 |
  |  └ Task A-2         |  └ Task B-0                 |
  +---------------------------------------------------+
        ^                          |
   Source: 外部系统 -> Kafka    Sink: Kafka -> 外部系统
   (MySQL, MQ, 文件...)        (ES, S3, JDBC, Redis...)
概念角色
WorkerJVM 进程,提供运行环境与 REST API;多个 Worker 组成集群
Connector逻辑作业定义(连什么库、什么表、写哪个 topic),负责拆分任务
Task实际搬数据的工作单元,被分配到各 Worker 上并行执行
  • Source Connector:外部系统 → Kafka(内部就是 Producer)
  • Sink Connector:Kafka → 外部系统(内部就是 Consumer Group,天然支持并行与 Rebalance)
  • 容错:Worker 挂了,它上面的 Task 自动迁移到其他 Worker;进度(source offset / consumer offset)都存在 Kafka 内部 topic,断点续传。

2. Standalone 与 Distributed 模式

维度StandaloneDistributed
进程单进程多 Worker 组集群
配置方式命令行 + properties 文件REST API 提交 JSON
offset 存储本地文件Kafka 内部 topic
容错/扩展无自动任务迁移、加 Worker 即扩容
适用本地测试、单机采集生产环境

Distributed Worker 配置要点:

# connect-distributed.properties
bootstrap.servers=kafka1:9092,kafka2:9092
group.id=connect-cluster-1               # Connect 集群标识
# 三个内部 topic: 存配置/进度/状态, 都要求 compact
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
plugin.path=/opt/connect-plugins         # 连接器 jar 放这里
bin/connect-distributed.sh config/connect-distributed.properties
# 验证
curl localhost:8083/connector-plugins | jq   # 已安装的连接器

3. REST API 管理连接器

Distributed 模式一切通过 REST API(默认 8083):

# 创建 Sink 连接器: 把 orders topic 写入 Elasticsearch
curl -X POST localhost:8083/connectors \
  -H "Content-Type: application/json" -d '{
  "name": "es-orders-sink",
  "config": {
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "topics": "orders",
    "connection.url": "http://es:9200",
    "tasks.max": "3",
    "key.ignore": "false",
    "errors.tolerance": "all",
    "errors.deadletterqueue.topic.name": "orders-sink-dlq",
    "errors.deadletterqueue.context.headers.enable": "true"
  }}'
 
# 常用运维命令
curl localhost:8083/connectors                          # 列表
curl localhost:8083/connectors/es-orders-sink/status    # 状态 (RUNNING/FAILED)
curl -X POST localhost:8083/connectors/es-orders-sink/tasks/0/restart
curl -X PUT  localhost:8083/connectors/es-orders-sink/pause
curl -X PUT  localhost:8083/connectors/es-orders-sink/resume
curl -X DELETE localhost:8083/connectors/es-orders-sink

errors.tolerance=all + 死信队列(DLQ)是 Sink 连接器的标配:个别坏消息进 DLQ,不阻塞整条管道。

4. SMT:单消息转换

SMT(Single Message Transform)在消息进出 Kafka 的路上做轻量转换,纯配置:

{
  "transforms": "mask,addTs,route",
  "transforms.mask.type": "org.apache.kafka.connect.transforms.MaskField$Value",
  "transforms.mask.fields": "phone,id_card",
 
  "transforms.addTs.type": "org.apache.kafka.connect.transforms.InsertField$Value",
  "transforms.addTs.timestamp.field": "sync_ts",
 
  "transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
  "transforms.route.regex": "db1\\.shop\\.(.*)",
  "transforms.route.replacement": "cdc-$1"
}

常用 SMT:字段屏蔽(MaskField)、插入字段(InsertField)、改 topic 名(RegexRouter)、抽取 key(ValueToKey + ExtractField)、扁平化嵌套(Flatten)、类型转换(Cast)。复杂逻辑别硬塞 SMT——那是 Streams/Flink 的活。

5. CDC 与 Debezium

CDC(Change Data Capture)通过解析数据库事务日志(MySQL binlog、PostgreSQL WAL)捕获行级变更,不侵入业务、不轮询、近实时。Debezium 是最流行的开源 CDC 连接器族:

{
  "name": "mysql-shop-cdc",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "database.hostname": "mysql", "database.port": "3306",
    "database.user": "cdc_user", "database.password": "******",
    "database.server.id": "5701",
    "topic.prefix": "shop",
    "database.include.list": "shop",
    "table.include.list": "shop.orders,shop.order_items",
    "schema.history.internal.kafka.bootstrap.servers": "kafka1:9092",
    "schema.history.internal.kafka.topic": "schema-history.shop",
    "snapshot.mode": "initial"
  }}
  • 首次启动做全量快照(snapshot.mode=initial),之后无缝切到 binlog 增量。
  • 每张表一个 topic(shop.shop.orders),消息含 before/after 镜像与操作类型 c/u/d。
  • 下游配 compact topic + ES/数仓 Sink,就是一条完整的实时同步链路:
MySQL --binlog--> Debezium Source --> Kafka --> Sink --> ES / ClickHouse / S3
⚠️Connect 常见坑
  1. converter 不匹配:Source 用 Avro 写入,Sink 配 JsonConverter 读,直接反序列化失败。全链路 converter 与 Schema Registry 配置必须一致。
  2. tasks.max 不等于实际并行度:Sink 的并行上限还是 topic 分区数;JDBC Source 按表拆分,单表不会再拆。
  3. Debezium 快照锁表:MySQL 大表 initial 快照可能持有全局读锁,生产库要评估 snapshot.locking.mode 或用增量快照(signal 机制)。
  4. 内部 topic 被误删:connect-offsets 删了 = 所有 Source 从头重跑。三个内部 topic 要按生产标准(3 副本、compact)保护。
💡选 Connect 还是自己写?

判断标准:搬运逻辑是否"通用且无业务语义"。是(库到库、topic 到 ES)→ Connect,省掉 80% 工程量;否(要调用业务规则、多流 join、复杂聚合)→ 自己写消费者或用 Streams/Flink。两者常常组合:Debezium 进 Kafka,Streams 加工,再由 Sink Connect 出去。

小结

  • Connect 把"搬数据"框架化:Worker 提供运行时,Connector 定义作业,Task 并行执行
  • 生产用 Distributed 模式:REST API 管理、进度存内部 topic、Worker 故障自动迁移任务
  • Sink 标配 errors.tolerance=all + 死信队列;SMT 做轻量转换
  • Debezium 基于 binlog 的 CDC 是数据库到 Kafka 的标准方案:全量快照 + 增量日志
  • converter 全链路一致、内部 topic 妥善保护,是两条运维铁律
🎯练习
  1. 用 Distributed 模式起一个 FileStreamSource(读本地文件进 topic)和 FileStreamSink(写出到另一文件),体验全程零代码的数据管道。
  2. 给 Sink 加一个 InsertField SMT 注入同步时间戳,验证输出文件里多了字段。
  3. 用 Docker 起 MySQL + Debezium,对一张表做 insert/update/delete,观察 CDC 消息中的 before/after 结构与 op 字段。