Kafka Connect 深入
概述
Kafka Connect 是 Kafka 生态的数据集成框架:Source Connector 把外部数据导入 Kafka,Sink Connector 把 Kafka 数据导出到外部系统。支持 Standalone(单机)与 Distributed(集群)两种模式,REST API 统一管理。本文讲透 Connector 开发、API 管理与部署模式。
一、Kafka Connect 定位
外部系统 ⇄ Kafka
Source Connector:外部 → Kafka
Sink Connector:Kafka → 外部| 特性 | 说明 |
|---|---|
| 框架 | 分布式集成 |
| 管理 | REST API |
| 扩展 | 插件生态 |
| 可靠 | Offset 管理 |
适用:
数据库同步(Debezium)
HDFS/S3 导出
文件/日志采集二、核心概念
| 概念 | 说明 |
|---|---|
| Connector | 连接器定义(逻辑) |
| Task | 执行单元(物理) |
| Worker | 运行进程 |
| Config | 连接器配置 |
| Offset | 消费位点 |
| Schema | 数据格式 |
关系:
Connector(1)→ Task(1..n)
Worker 集群运行 Task三、Source Connector
3.1 作用
外部数据 → Kafka Topic:
JDBC、文件、MQ、S3...3.2 工作原理
轮询外部数据源:
拉取新数据 → 转为 SinkRecord → 写入 Kafka
Offset 记录位置| 接口 | 说明 |
|---|---|
| start | 初始化 |
| poll | 拉取数据 |
| flush | 提交 |
| stop | 停止 |
3.3 示例:JDBC Source
json
{
"name": "jdbc-source",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"connection.url": "jdbc:mysql://mysql:3306/db",
"connection.user": "root",
"connection.password": "***",
"mode": "incrementing",
"incrementing.column.name": "id",
"topic.prefix": "mysql-"
}
}| 模式 | 说明 |
|---|---|
| incrementing | 自增列 |
| timestamp | 时间戳列 |
| bulk | 全量 |
| 增量组合 | 时间+自增 |
四、Sink Connector
4.1 作用
Kafka Topic → 外部系统:
HDFS、ES、JDBC、S3...4.2 工作原理
消费 Kafka:
按 Topic/分区
批量写入外部
记录 offset(幂等)| 接口 | 说明 |
|---|---|
| start | 初始化 |
| put | 写入批次 |
| flush | 提交 |
| stop | 停止 |
4.3 示例:JDBC Sink
json
{
"name": "jdbc-sink",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
"connection.url": "jdbc:mysql://target:3306/dw",
"connection.user": "root",
"connection.password": "***",
"topics": "orders",
"insert.mode": "upsert",
"pk.mode": "record_key",
"pk.fields": "order_id"
}
}| 配置 | 说明 |
|---|---|
| insert.mode | insert/upsert/update |
| pk.mode | 主键策略 |
| batch.size | 批量 |
| 幂等 | 唯一键 |
五、Connector 开发
5.1 开发流程
实现接口:
Source:SourceConnector + SourceTask
Sink:SinkConnector + SinkTask
打包成插件 → 放 plugins 目录5.2 关键实现
java
// Source Task 核心
public class MySourceTask extends SourceTask {
public List<SourceRecord> poll() throws InterruptedException {
// 拉取外部数据 → SourceRecord(topic, key, value)
}
}
// Sink Task 核心
public class MySinkTask extends SinkTask {
public void put(Collection<SinkRecord> records) {
// 写入外部系统
}
}| 开发要点 | 说明 |
|---|---|
| Offset | Source 记录位点 |
| 幂等 | Sink 可重放 |
| 配置 | ConfigDef |
| 版本 | 元数据 |
六、REST API
6.1 管理接口
| 方法 | 路径 | 说明 |
|---|---|---|
| GET | /connectors | 列表 |
| POST | /connectors | 创建 |
| GET | /connectors/ | 状态 |
| PUT | /connectors/{name}/config | 更新 |
| DELETE | /connectors/ | 删除 |
| POST | /connectors/{name}/pause | 暂停 |
| POST | /connectors/{name}/resume | 恢复 |
| GET | /connectors/{name}/tasks | 任务 |
| GET | /connectors/{name}/status | 详细状态 |
bash
# 查看状态
curl http://connect:8083/connectors/jdbc-source/status6.2 运维操作
| 操作 | 场景 |
|---|---|
| 创建 | 新同步任务 |
| 更新 | 修改配置(热更新) |
| 暂停/恢复 | 维护 |
| 删除 | 下线 |
状态字段:
RUNNING / PAUSED / FAILED
TASK 状态与错误七、Standalone 模式
7.1 特点
单进程运行:
配置本地文件
单 Workerproperties
# connect-standalone.properties
bootstrap.servers=kafka:9092
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
offset.storage.file.filename=/tmp/connect.offsets| 特点 | 说明 |
|---|---|
| 简单 | 单机 |
| 单点 | 无 HA |
| 适用 | 开发/小规模 |
八、Distributed 模式
8.1 特点
集群运行:
多 Worker
任务自动分配/再平衡
Offset 存 Kafkaproperties
# connect-distributed.properties
bootstrap.servers=kafka:9092
group.id=connect-cluster
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
offset.storage.topic=connect-offsets
config.storage.topic=connect-configs
status.storage.topic=connect-status| 特性 | 说明 |
|---|---|
| HA | 节点故障任务迁移 |
| 扩展 | 加节点 |
| 存储 | Kafka Topic |
| 再平衡 | 任务重分配 |
8.2 状态存储
| Topic | 用途 |
|---|---|
| connect-configs | 连接器配置 |
| connect-offsets | 消费位点 |
| connect-status | 状态信息 |
分布式模式生产推荐:
高可用 + 扩展性九、部署与运维
9.1 插件管理
plugins 目录:
打包插件(fat jar)
配置 plugin.path
多版本隔离9.2 监控
| 指标 | 说明 |
|---|---|
| 任务状态 | RUNNING/FAILED |
| 延迟 | 端到端 |
| 吞吐 | 记录/秒 |
| 错误 | 重试/失败 |
监控方案:
JMX / Prometheus
任务状态轮询9.3 常见问题
| 问题 | 处理 |
|---|---|
| 任务失败 | 看日志/状态 |
| 数据重复 | Sink 幂等 |
| 延迟高 | 扩 Task/并行 |
| 插件缺失 | 打包部署 |
| 再平衡频繁 | 检查节点 |
常见问题速查
| 问题 | 要点 |
|---|---|
| Source vs Sink | 进 Kafka vs 出 Kafka |
| 单机还是集群 | Standalone vs Distributed |
| 怎么管理 | REST API |
| 开发 Connector | Connector+Task 接口 |
| 数据不丢 | Offset + 幂等 |