1. Kafka Connect 概述
Kafka Connect 是 Kafka 官方提供的数据集成框架,用于在 Kafka 与外部系统(数据库、文件、对象存储、搜索引擎等)之间建立可靠的数据管道。它通过插件化的 Connector 机制,无需编写代码即可实现数据导入导出。
1.1 Connect 架构
┌─────────────────────────────────────────────────────────────┐
│ Kafka Connect 集群 │
│ │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────┐ │
│ │ Worker 1 │ │ Worker 2 │ │ Worker 3 │ │
│ │ │ │ │ │ │ │
│ │ [Connector A]│ │ [Connector B]│ │(备用) │ │
│ │ Task 1 │ │ Task 1 │ │ │ │
│ │ Task 2 │ │ Task 2 │ │ │ │
│ └──────────────┘ └──────────────┘ └──────────┘ │
│ │ │ │
└─────────┼─────────────────────┼─────────────────────────────┘
│ │
Source Connector Sink Connector
│ │
┌───────▼──────┐ ┌─────▼──────┐
│ 外部系统 │ │ 外部系统 │
│ MySQL/Mongo │ │ ES/S3/HDFS │
└──────────────┘ └────────────┘
1.2 Connector、Task、Worker 的关系
| 组件 | 角色 | 说明 |
|---|---|---|
| Connector | 逻辑定义 | 定义数据源/目标的配置,进行 Task 拆分 |
| Task | 执行单元 | 真正执行数据复制的工作单元,无状态 |
| Worker | 运行进程 | 承载 Task 运行的 JVM 进程,支持集群和分布式调度 |
2. Source Connector
Source Connector 从外部系统读取数据,写入 Kafka Topic。
2.1 Debezium CDC(变更数据捕获)
Debezium 是最常用的 Source Connector,基于数据库日志(binlog/WAL)实现增量数据捕获。
Debezium 工作原理:
MySQL Binlog Debezium Kafka
INSERT ──→ [捕获] ──→ ──→ Avro/JSON ──→ db.shop.users
UPDATE ──→ [解析] ──→ ──→ 变更事件 ──→ db.shop.orders
DELETE ──→ [转换] ──→ ──→ ──→ db.shop.inventory
初始快照(Snapshot):
首次启动时,全量读取表数据 → 写入 Kafka → 然后切换到 Binlog 增量
2.2 Debezium 事件格式
{
"before": { "id": 1, "name": "Alice", "age": 30 },
"after": { "id": 1, "name": "Alice", "age": 31 },
"source": {
"version": "2.5.0.Final",
"connector": "mysql",
"name": "db-shop",
"ts_ms": 1705315200000,
"db": "shop",
"table": "users",
"pos": "binlog.001:1234",
"row": 0
},
"op": "u",
"ts_ms": 1705315200001
}
| 字段 | 说明 |
|---|---|
before | 更新/删除前的数据 |
after | 更新/插入后的数据 |
source | 元信息(数据库、表、binlog 位置) |
op | 操作类型:c=create, u=update, d=delete, r=read(快照) |
ts_ms | 处理时间戳 |
2.3 Debezium 配置
{
"name": "mysql-source-connector",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql",
"database.port": "3306",
"database.user": "debezium",
"database.password": "dbz",
"database.server.id": "184054",
"database.server.name": "db-shop",
"database.include.list": "shop",
"table.include.list": "shop.users,shop.orders",
"snapshot.mode": "initial",
"topic.prefix": "cdc",
"schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
"schema.history.internal.kafka.topic": "schema-changes.shop",
"transforms": "unwrap",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.drop.tombstones": "false",
"transforms.unwrap.delete.handling.mode": "rewrite"
}
}
3. Sink Connector
Sink Connector 从 Kafka Topic 读取数据,写入外部系统。
3.1 JDBC Sink Connector
{
"name": "jdbc-sink-connector",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
"tasks.max": "3",
"topics": "orders",
"connection.url": "jdbc:postgresql://postgres:5432/analytics",
"connection.user": "analytics",
"connection.password": "secret",
"auto.create": "true",
"auto.evolve": "true",
"insert.mode": "upsert",
"pk.fields": "order_id",
"pk.mode": "record_key"
}
}
3.2 Elasticsearch Sink Connector
{
"name": "es-sink-connector",
"config": {
"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
"tasks.max": "3",
"topics": "orders",
"connection.url": "http://elasticsearch:9200",
"type.name": "_doc",
"key.ignore": "false",
"schema.ignore": "true",
"transforms": "extractKey",
"transforms.extractKey.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
"transforms.extractKey.field": "order_id"
}
}
4. 常用 Connector 列表
| Connector | 类型 | 来源 | 用途 |
|---|---|---|---|
| Debezium | Source | Debezium 社区 | MySQL/PostgreSQL/MongoDB CDC |
| JDBC Source | Source | Confluent | 轮询查询数据库 |
| FileStream | Source | Kafka 官方 | 读取日志文件 |
| S3 Source | Source | Confluent | 从 S3 读取数据 |
| JDBC Sink | Sink | Confluent | 写入关系型数据库 |
| Elasticsearch Sink | Sink | Confluent | 写入 ES 索引 |
| S3 Sink | Sink | Confluent | 写入 S3(Parquet/JSON) |
| HDFS Sink | Sink | Confluent | 写入 HDFS |
| Redis Sink | Sink | 社区 | 写入 Redis |
| InfluxDB Sink | Sink | 社区 | 写入时序数据库 |
5. Kafka Connect 运行时
5.1 Standalone vs Distributed
| 模式 | 适用场景 | 特点 |
|---|---|---|
| Standalone | 开发测试、简单任务 | 单进程,配置存储在本地文件 |
| Distributed | 生产环境 | 多 Worker 集群,配置存储在 Kafka topic(offset.storage.topic) |
5.2 部署 Distributed 模式
# connect-distributed.properties
bootstrap.servers=kafka1:9092,kafka2:9092,kafka3:9092
group.id=connect-cluster
# 配置存储topic
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-status
# 启动 Worker
connect-distributed.sh config/connect-distributed.properties
# 注册 Connector(REST API)
curl -X POST http://localhost:8083/connectors \
-H "Content-Type: application/json" \
-d @mysql-source-connector.json
# 查看状态
curl http://localhost:8083/connectors/mysql-source-connector/status
6. Transform(数据转换)
Kafka Connect 支持单消息转换(SMT,Single Message Transform),在数据流转过程中转换格式。
{
"transforms": "maskField,extractValue",
"transforms.maskField.type": "org.apache.kafka.connect.transforms.MaskField$Value",
"transforms.maskField.fields": "password,ssn",
"transforms.extractValue.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
"transforms.extractValue.field": "id"
}
| Transform | 作用 |
|---|---|
ExtractField | 提取嵌套字段 |
ReplaceField | 重命名/删除字段 |
MaskField | 脱敏(替换为 null 或掩码) |
ValueToKey | 将 value 中的字段设为 key |
TimestampConverter | 时间格式转换 |
Filter | 按条件过滤消息 |
7. 生产实践
7.1 高可用配置
Worker 数量 ≥ 2(防止单点)
tasks.max = min(分区数, Worker 数量 × 单 Worker 承载能力)
监控指标:
- source-record-active-count(活跃记录数)
- put-batch-avg-time-ms(批处理平均时间)
- sink-record-send-rate(发送速率)
7.2 常见问题
| 问题 | 原因 | 解决 |
|---|---|---|
| 数据延迟 | Sink 消费慢 | 增加 tasks.max、优化目标数据库写入 |
| 字段不匹配 | Schema 变更 | 启用 auto.evolve、使用 Schema Registry |
| 消息格式错误 | 非 AVro/JSON 格式 | 配置正确的 Converter(key.converter、value.converter) |
| 连接断开 | 网络不稳定 | 配置重试策略、连接池保持 |
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。