上一篇【第71篇】Kafka Connect实战——10分钟搭建MySQL到Elasticsearch数据管道 下一篇【第73篇】Kafka Streams快速上手——用流处理做实时WordCount
摘要
Kafka Connect用起来简单,但内部工作机制相当精妙。一个Connector怎么变成多个并行Task?Task失败了吗谁来重启?offset存在哪?数据格式怎么转换?这些问题如果不搞清楚,一旦出问题就只能干瞪眼。
本文将深入Connect的内部机制:Connector与Task的关系(分布式并行模型)、Worker线程模型、Offset存储机制(__consumer_offsets vs connect-offsets)、Converter的作用(JSON/Avro/Protobuf),以及自定义Connector开发实战。读完这篇,你能对Connect的内部运作了如指掌。
一、Connector与Task的关系——一对多的并行模型
1.1 核心概念图解
【Connector → Task 关系图】
┌─────────────────────────────────────────────┐
│ Kafka Connect Worker │
│ │
│ Connector: mysql-orders-source │
│ (蓝图:定义"要同步什么数据") │
│ │ │
│ ▼ │
│ Task 0: 同步 orders 表 (id 0~1000000) │
│ Task 1: 同步 orders 表 (id 1000000~2000000) │
│ Task 2: 同步 order_items 表 │
│ │
│ → 每个 Task 是独立的线程,并行执行 │
│ → Task 数由 tasks.max 参数控制(上限) │
└─────────────────────────────────────────────┘
类比:
Connector = 包工头(看图纸,分配任务)
Task = 工人(真正干活的)
Worker = 工地(提供水电和工具)
关键规则:
- Connector只负责规划和分配,不直接处理数据
- Task才负责实际数据读写
- tasks.max是上限,实际Task数由Connector的taskConfigs()方法决定
1.2 Task并行度分配算法
// 以JdbcSourceConnector为例(简化逻辑)
@Override
public List<Map<String, String>> taskConfigs(int maxTasks) {
// 1. 获取需要同步的表列表
List<String> tables = config.getStringList("table.whitelist");
// 2. 每个表尽量均匀分配到Task
List<List<String>> taskTableAssignments =
assignTablesToTasks(tables, maxTasks);
// 3. 生成每个Task的配置
List<Map<String, String>> taskConfigs = new ArrayList<>();
for (int i = 0; i < taskTableAssignments.size(); i++) {
Map<String, String> taskConfig = new HashMap<>();
taskConfig.put("tables",
String.join(",", taskTableAssignments.get(i)));
taskConfigs.add(taskConfig);
}
return taskConfigs;
}
【Task分配示例】
tables = [orders, order_items, users, products]
tasks.max = 3
分配结果:
Task 0: [orders, order_items]
Task 1: [users]
Task 2: [products]
→ 如果 tables=1,即使 tasks.max=5,实际也只有1个Task
二、Worker线程模型——Task是怎么跑起来的
2.1 完整线程模型
【Worker线程模型(Distributed模式)】
┌─────────────────────────────────────────────┐
│ Worker Process (JVM) │
│ │
│ ┌──────────────────────────────────┐ │
│ │ REST API Thread (HTTP Server) │ │
│ │ → 处理 Connector 的 CRUD 请求 │ │
│ └──────────────┬───────────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────────────────────────┐ │
│ │ Herder Thread (每个Connector一个) │ │
│ │ → 管理 Connector 生命周期 │ │
│ │ → 调用 taskConfigs() 分配Task │ │
│ │ → 处理 Connector 的暂停/恢复 │ │
│ └──────────────┬───────────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────────────────────────┐ │
│ │ Task Threads (每个Task一个线程) │ │
│ │ Task-0-0: 执行 SourceTask.poll() │ │
│ │ Task-0-1: 执行 SourceTask.poll() │ │
│ │ → 每个Task在独立线程中运行 │ │
│ └──────────────────────────────────┘ │
└─────────────────────────────────────────────┘
注意:
• Source Connector的Task主动"拉取"数据(poll())
• Sink Connector的Task被动"推送"数据(put())
2.2 Task故障重启机制
【Task故障处理流程】
T0: Task-0-0 正常运行,poll() 返回数据
T1: Task-0-0 抛出 RuntimeException
T2: Worker 捕获异常,标记 Task 为 FAILED
T3: Worker 等待 retry.backoff.ms (默认5000ms)
T4: Worker 尝试重启 Task-0-0
T5: 重启成功 → 状态变为 RUNNING
重启失败 → 继续重试(最多 Integer.MAX_VALUE 次)
// Worker 中 Task 重启的核心逻辑(简化)
public class WorkerSinkTask {
private volatile boolean stopRequested = false;
public void run() {
while (!stopRequested && !cancelled) {
try {
// 执行一次 put()(Sink)或 poll()(Source)
iteration();
// 成功后重置重试计数
retryCount = 0;
} catch (Exception e) {
retryCount++;
if (retryCount > maxRetries) {
// 超过最大重试次数,标记 Task 为 FAILED
failTask(e);
break;
}
// 指数退避等待
Thread.sleep(retryBackoffMs * Math.pow(2, retryCount));
}
}
}
}
生产环境调优:
# 控制 Task 重启行为
# 每次重试间隔(毫秒)
retry.backoff.ms=3000
# 最大重试次数(默认 Integer.MAX_VALUE,即无限重试)
# 建议设置为 10~20,配合告警人工介入
max.retries=10
三、Offset存储机制——Connect的"书签"
3.1 Source Connector的Offset
Source Connector需要记录"上次读到哪了",以便重启后不丢数据。
【Source Offset 存储格式】
Key: { "connector": "mysql-orders-source", "partition": {…} }
Value: { "offset": {…} }
示例(JDBC Source):
Key: { "connector": "mysql-orders-source",
"partition": {"table": "orders", "id": 0} }
Value: { "offset": {"incrementing_id": 1005, "timestamp_ns": 1717056000000} }
→ 含义:orders 表上次读到 id=1005 的位置
Offset存储位置(二选一):
【Offset 存储对比】
方案A:Kafka Topic (connect-offsets) ★ 推荐
┌─────────────────────────────────────┐
│ • 持久化到 Kafka Topic │
│ • 多个 Worker 共享(Distributed模式)│
│ • 高可用(依赖 Kafka 副本机制) │
└─────────────────────────────────────┘
方案B:文件系统(Standalone模式)
┌─────────────────────────────────────┐
│ • 存储在本地文件 (offset.storage.file) │
│ • 只适用于单机测试 │
│ • 生产环境不要用! │
└─────────────────────────────────────┘
3.2 Sink Connector的Offset
Sink Connector的offset实际上复用Consumer的offset机制(存在__consumer_offsets中):
【Sink Connector Offset 流转】
Sink Task 消费 Kafka Topic 的数据:
┌─────────────────────────────────────┐
│ Topic: mysql-orders │
│ Partition 0: [msg1][msg2][msg3]…[msg100] │
│ ▲ │
│ │ │
│ 已消费并提交 offset │
│ (= 消费进度) │
└─────────────────────────────────────┘
Sink Connector 内部:
• 使用 KafkaConsumer 消费数据
• offset 提交到 __consumer_offsets
• Connect 框架自动管理,不需要手动配置
四、Converter的作用——数据格式的"翻译官"
4.1 为什么需要Converter
【数据流转中的格式转换】
Source System (MySQL) Kafka Connect Kafka Topic Sink System (ES)
┌────────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐
│ MySQL Row │───►│ Source │──►│ JSON │──►│ ES │
│ (二进制) │ │ Connector │ │ /Avro │ │ Document │
└────────────┘ │ (Java Obj) │ │ /Protobuf │ │ (JSON) │
└──────────┘ └──────────┘ └────────────┘
│ │
▼ ▼
┌──────────┐ ┌──────────┐
│ Converter│ │ Converter│
│ (序列化) │ │ (反序列化)│
└──────────┘ └──────────┘
Converter 负责:
① Source 端:Java Object → byte[](写入 Kafka)
② Sink 端:byte[] → Java Object(从 Kafka 读取)
4.2 内置Converter对比
| JsonConverter | ✅ | 中 | ✅ | 调试、简单场景 |
| AvroConverter | ✅ | 高 | ❌ | 生产环境(Schema Registry) |
| ProtobufConverter | ✅ | 最高 | ❌ | 高性能场景 |
| StringConverter | ❌ | 最高 | ✅ | Key/Value是纯字符串 |
| ByteArrayConverter | ❌ | 最高 | ❌ | 已是二进制格式 |
4.3 Converter配置实战
# ===== Source 端(写入 Kafka 时的格式)=====
# 方案A: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
# 方案B:Avro(有Schema,推荐生产环境)
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schemas.enable=true
value.converter.schema.registry.url=http://localhost:8081
# ===== Sink 端(从 Kafka 读取时的格式)=====
# 注意:Sink 端的 Converter 必须与 Source 端匹配!
# 如果 Source 用了 Avro,Sink 也必须用 AvroConverter
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://localhost:8081
五、自定义Connector开发实战
5.1 开发Source Connector(实战案例:文件系统Source)
// 目标:监听指定目录,新文件出现时自动写入 Kafka
// 完整可运行的 SourceConnector 示例
package com.example.connect;
import org.apache.kafka.connect.source.SourceConnector;
import org.apache.kafka.connect.source.SourceRecord;
import org.apache.kafka.connect.source.SourceTask;
import java.util.*;
public class FileMonitorSourceConnector extends SourceConnector {
private String monitorDir;
private String topic;
@Override
public void start(Map<String, String> props) {
monitorDir = props.get("monitor.dir");
topic = props.get("topic");
}
@Override
public Class<? extends SourceTask> taskClass() {
return FileMonitorSourceTask.class;
}
@Override
public List<Map<String, String>> taskConfigs(int maxTasks) {
// 将目录下的文件均匀分配到各 Task
List<Map<String, String>> configs = new ArrayList<>();
File dir = new File(monitorDir);
File[] files = dir.listFiles();
for (int i = 0; i < Math.min(maxTasks, files.length); i++) {
Map<String, String> taskConfig = new HashMap<>();
taskConfig.put("monitor.dir", monitorDir);
taskConfig.put("file.path", files[i].getAbsolutePath());
taskConfig.put("topic", topic);
configs.add(taskConfig);
}
return configs;
}
@Override
public void stop() {}
@Override
public ConfigDef config() {
return new ConfigDef()
.define("monitor.dir", ConfigDef.Type.STRING, "重要性:HIGH")
.define("topic", ConfigDef.Type.STRING, "重要性:HIGH");
}
}
// 配套 Task 实现(核心:poll() 方法)
public class FileMonitorSourceTask extends SourceTask {
private String filePath;
private String topic;
private BufferedReader reader;
private long offset = 0;
@Override
public void start(Map<String, String> props) {
filePath = props.get("file.path");
topic = props.get("topic");
try {
reader = new BufferedReader(new FileReader(filePath));
// 从 offset 存储中恢复读取位置
offset = context.offsetStorageReader()
.offset(Collections.singletonMap("file", filePath))
.get("position");
if (offset > 0) reader.skip(offset);
} catch (IOException e) {
throw new RuntimeException(e);
}
}
@Override
public List<SourceRecord> poll() throws InterruptedException {
List<SourceRecord> records = new ArrayList<>();
try {
String line;
while ((line = reader.readLine()) != null) {
// 构建 SourceRecord(写入 Kafka)
records.add(new SourceRecord(
Collections.singletonMap("file", filePath),
Collections.singletonMap("position", offset),
topic,
line
));
offset += line.getBytes().length + 1;
}
} catch (IOException e) {
// 文件读取完毕或出错
}
return records.isEmpty() ? null : records;
}
@Override
public void stop() {
try { reader.close(); } catch (IOException e) {}
}
}
5.2 打包与部署
# Step 1: 打包成 uber-jar(包含所有依赖)
mvn clean package
# Step 2: 将 jar 放到 Connect 的 plugin.path 目录
cp target/file-monitor-connector-1.0.0.jar /opt/kafka/plugins/
# Step 3: 重启所有 Worker(加载新插件)
bin/connect-distributed.sh config/connect-distributed.properties &
# Step 4: 创建 Connector 实例
curl -X POST http://localhost:8083/connectors \\
-H "Content-Type: application/json" \\
-d '{
"name": "file-monitor-source",
"config": {
"connector.class": "com.example.connect.FileMonitorSourceConnector",
"monitor.dir": "/var/log/app",
"topic": "app-logs",
"tasks.max": "3"
}
}'
六、生产环境最佳实践
6.1 Connector配置Checklist
【生产环境 Connector 检查清单】
□ 高可用
• Distributed 模式(非 Standalone)
• Worker 数量 ≥ 2
• 内部 Topic 副本数 = 3
□ Offset 管理
• offset.storage.replication.factor = 3
• 定期检查 offset lag(防止重复消费)
□ 并行度
• tasks.max 设置合理(不过大,不浪费资源)
• 观察 Task 状态:curl localhost:8083/connectors/{name}/tasks
□ 错误处理
• 配置了 connector.class 的错误重试参数
• 有告警监控 Task 失败事件
□ 数据格式
• 生产环境使用 Avro + Schema Registry(有 Schema 演进支持)
• 或使用 Protobuf(高性能场景)
6.2 监控关键指标
# Prometheus 监控规则
groups:
– name: kafka_connect
rules:
# 告警1:Connector 状态异常
– alert: ConnectorUnhealthy
expr: kafka_connect_connector_state != 1
for: 1m
labels:
severity: critical
annotations:
summary: "Connector {{ $labels.connector }} 状态异常"
description: "当前状态: {{ $value }} (1=RUNNING)"
# 告警2:Task 数量不足
– alert: ConnectorTaskShortage
expr: kafka_connect_connector_task_count < kafka_connect_connector_task_config_count
for: 5m
labels:
severity: warning
annotations:
summary: "Connector {{ $labels.connector }} Task 数量不足"
description: "预期: {{ $value }} 个,实际: {{ $value }} 个"
本篇小结
今天我们深入了Kafka Connect的内部工作机制:
核心要点:Connector的开发90%的难度在错误处理和Offset管理上——网络抖动怎么办?文件被删了怎么办?重启后从哪里继续?想清楚这些,才能写出一个生产级Connector。
上一篇【第71篇】Kafka Connect实战——10分钟搭建MySQL到Elasticsearch数据管道 下一篇【第73篇】Kafka Streams快速上手——用流处理做实时WordCount







