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...)| 概念 | 角色 |
|---|---|
| Worker | JVM 进程,提供运行环境与 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 模式
| 维度 | Standalone | Distributed |
|---|---|---|
| 进程 | 单进程 | 多 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-sinkerrors.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- converter 不匹配:Source 用 Avro 写入,Sink 配 JsonConverter 读,直接反序列化失败。全链路 converter 与 Schema Registry 配置必须一致。
- tasks.max 不等于实际并行度:Sink 的并行上限还是 topic 分区数;JDBC Source 按表拆分,单表不会再拆。
- Debezium 快照锁表:MySQL 大表 initial 快照可能持有全局读锁,生产库要评估 snapshot.locking.mode 或用增量快照(signal 机制)。
- 内部 topic 被误删:connect-offsets 删了 = 所有 Source 从头重跑。三个内部 topic 要按生产标准(3 副本、compact)保护。
判断标准:搬运逻辑是否"通用且无业务语义"。是(库到库、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 妥善保护,是两条运维铁律
- 用 Distributed 模式起一个 FileStreamSource(读本地文件进 topic)和 FileStreamSink(写出到另一文件),体验全程零代码的数据管道。
- 给 Sink 加一个 InsertField SMT 注入同步时间戳,验证输出文件里多了字段。
- 用 Docker 起 MySQL + Debezium,对一张表做 insert/update/delete,观察 CDC 消息中的 before/after 结构与 op 字段。