Kafka Connect 自带 JDBC、S3、Elasticsearch 等大量官方连接器,但真正落到企业环境,几乎总会撞上「官方连接器覆盖不到」的场景:某个自研存储、某套私有协议、某种带状态的分页拉取。此时唯一的路就是自己写 Connector。本文不讲怎么用现成连接器,而是把自定义 Connector 的开发契约、并行度设计、offset 语义与打包运维完整拆一遍。
1. Connect 架构与扩展点
1.1 五层扩展点
Kafka Connect 的扩展点可以抽象成五层,理解这五层的职责边界,是写对 Connector 的前提:
| 层 | 职责 | 开发者是否常改 |
|---|---|---|
| Worker | 进程与线程池,负责调度 Task、管理 offset、提供 REST 接口 | 基本不改,只做配置 |
| Connector | 逻辑单元,负责切分 Task 与校验配置 | 必写 |
| Task | 真正干活的单元,一个 Task 对应一份并行度 | 必写 |
| Converter | 在 byte 与 Connect 内部结构之间转换 | 少数场景自写 |
| Transform | 单条消息的无状态改写 | 按需自写 |
Worker 是运行容器,Connector 与 Task 是业务代码,Converter 与 Transform 是可选增强。绝大多数自定义 Connector 只需要实现 Connector 与 Task 两个类。
1.2 standalone 与 distributed
两种部署模式在开发阶段就要选定,因为 offset 存储位置不同:
# standalone:单进程,配置与 offset 都落本地文件
# offset.storage.file.filename=/data/connect/offset.dat
# 适合本地调试、单机采集,无高可用
# distributed:多进程集群,offset/config/status 三个内部主题
# offset.storage.topic=connect-offsets
# config.storage.topic=connect-configs
# status.storage.topic=connect-status
生产一律用 distributed。它把 offset、配置、状态持久化到 Kafka 内部主题,Worker 挂了之后 Task 会被其他 Worker 接管,从上次提交的 offset 继续。standalone 只在本地验证代码逻辑时使用。
1.3 插件目录与 classpath 隔离
插件(plugin)是 Connector 及其依赖 jar 的集合,放在 plugin.path 指定的目录下,每个插件一个子目录:
# 目录结构:每个插件独立目录,避免依赖冲突
# /opt/kafka/plugins/
# ├── my-source-connector/my-source-connector-1.0.0.jar
# └── my-sink-connector/my-sink-connector-1.0.0.jar
# plugin.path=/opt/kafka/plugins
Connect 为每个插件目录创建独立的类加载器(PluginClassLoader),实现依赖隔离。若把多个 Connector 的 jar 混在同一目录,极易出现「A 的依赖版本覆盖 B」的诡异问题。一个插件一个目录是硬纪律。
2. Source Connector 开发
2.1 接口与生命周期
Source 负责把外部系统的数据拉进 Kafka。生命周期固定为:start 校验配置并计算 Task 数 → taskConfigs 生成每个 Task 的配置 → stop 释放资源。
public class MySourceConnector extends SourceConnector {
private Map<String, String> configProps;
@Override
public void start(Map<String, String> props) {
// 校验配置,尽早失败
this.configProps = props;
if (props.get(MySourceConfig.ENDPOINT) == null) {
throw new ConnectException("endpoint 不能为空");
}
}
@Override
public Class<? extends Task> taskClass() {
return MySourceTask.class;
}
@Override
public List<Map<String, String>> taskConfigs(int maxTasks) {
// 按外部系统的分片数切分,最多 maxTasks 个
List<Map<String, String>> configs = new ArrayList<>();
int shardCount = Math.min(maxTasks, queryShardCount());
for (int i = 0; i < shardCount; i++) {
Map<String, String> cfg = new HashMap<>(configProps);
cfg.put(MySourceConfig.SHARD_ID, String.valueOf(i));
configs.add(cfg);
}
return configs;
}
@Override
public ConfigDef config() {
return MySourceConfig.CONFIG_DEF;
}
@Override
public String version() {
return "1.0.0";
}
}
关键点:taskConfigs 返回的列表长度就是 Task 数量,每个元素是该 Task 的独立配置。外部系统有多少个可并行分片,就切多少个 Task。
2.2 SourceTask 与 poll 契约
Task 的核心只有一个方法:poll()。Worker 在一个循环里反复调用它,每次返回一批 SourceRecord。
public class MySourceTask extends SourceTask {
private HttpClient client;
private String shardId;
private long lastOffset;
@Override
public void start(Map<String, String> props) {
this.shardId = props.get(MySourceConfig.SHARD_ID);
this.client = new HttpClient(props.get(MySourceConfig.ENDPOINT));
// 从 offset 存储恢复断点,null 表示全新 Task
Map<String, Object> offset = context.offsetStorageReader()
.offset(Collections.singletonMap("shard", shardId));
this.lastOffset = offset == null ? 0L : (Long) offset.get("position");
}
@Override
public List<SourceRecord> poll() throws InterruptedException {
List<Row> rows = client.fetch(shardId, lastOffset, BATCH_SIZE);
if (rows.isEmpty()) {
Thread.sleep(POLL_INTERVAL_MS); // 空轮询必须退避
return Collections.emptyList();
}
List<SourceRecord> records = new ArrayList<>(rows.size());
for (Row row : rows) {
Map<String, Object> partition = Collections.singletonMap("shard", shardId);
Map<String, Object> offset = new HashMap<>();
offset.put("position", row.position);
records.add(new SourceRecord(partition, offset, topic,
Schema.STRING_SCHEMA, row.key, valueSchema, row.value));
lastOffset = row.position;
}
return records;
}
}
2.3 partition 与 offset 的语义
SourceRecord 的 sourcePartition 与 sourceOffset 是 Source 端最容易被误解的两个字段:
- sourcePartition 标识数据来源的「分区」,对 Source 而言是外部系统的分片,比如数据库表、文件、队列。同一 sourcePartition 内的 offset 必须单调可比较。
- sourceOffset 是「下一个要读的位置」还是「刚读过的位置」,由开发者自行约定,但必须自洽——Connect 只负责持久化,不理解其含义。
# offset 提交时机:Connect 在每次 poll 之后、下一批之前提交
# 提交的是上一批最后一条 SourceRecord 的 sourceOffset
# 因此 offset 语义天然是「至少一次」:崩溃会重放最后一批
这也是为什么 Source 端必须做下游幂等:offset 提交与数据落 Kafka 不是原子的,崩溃时最后一批必然重复。
3. Sink Connector 开发
3.1 接口与生命周期
Sink 把 Kafka 数据写到外部系统。生命周期:start 校验 → open 分配分区 → put 批量写 → flush 落盘 → close 释放。
public class MySinkTask extends SinkTask {
private Writer writer;
private int remainingRetries;
@Override
public void start(Map<String, String> props) {
this.writer = new Writer(props.get(MySinkConfig.ENDPOINT));
this.remainingRetries = Integer.parseInt(
props.getOrDefault(SinkConnectorConfig.MAX_RETRIES_CONFIG, "10"));
}
@Override
public void put(Collection<SinkRecord> records) {
if (records.isEmpty()) {
return;
}
for (SinkRecord record : records) {
writer.stage(record.topic(), record.kafkaPartition(),
record.kafkaOffset(), record.value());
}
// 整批一次性提交,失败整批重试
retry(() -> writer.commitBatch());
}
@Override
public void flush(Map<TopicPartition, OffsetAndMetadata> currentOffsets) {
// Worker 在提交 offset 前调用,确保数据已真正落外部系统
writer.flush();
}
}
3.2 put 的批量语义与幂等
put 收到的是一批记录,不是单条,这是 Sink 性能的关键:逐条写外部系统会被网络往返拖死,必须攒批。同时 put 也天然适合做幂等去重,用 topic + partition + offset 三元组作为外部系统的唯一键,重复批次直接忽略。
# 幂等写入的三种常见手段
# 1) 唯一键 upsert:topic+partition+offset 作为主键,重复写覆盖
# 2) 事务:外部系统支持事务时,整批写入 + offset 提交同一事务
# 3) 版本号/时间戳:只接受比已存版本更新的数据
3.3 错误处理与死信队列
不是所有错误都该让整个 Task 挂掉。Connect 提供错误容忍配置:
{
"errors.tolerance": "all",
"errors.log.enable": true,
"errors.log.include.messages": true,
"errors.deadletterqueue.topic.name": "dlq-my-sink",
"errors.deadletterqueue.context.headers.enable": true
}
errors.tolerance=all 表示遇到单条转换/写入错误时跳过并继续,同时把失败记录投递到 DLQ 主题。但容忍错误会破坏 EOS 语义,且 DLQ 里堆积的记录需要人工处理。生产建议:对「数据格式错误」容忍(进 DLQ),对「外部系统不可用」不容忍(重试直至失败暂停)。
4. Task 切分与并行度
4.1 tasks.max 与真实 Task 数
tasks.max 是上限,不是目标值。真正运行的 Task 数由 Connector 决定:
- Sink Connector:Task 数等于订阅主题的分区总数(超过
tasks.max则截断)。因为 Sink 的并行单位是分区,一个分区只能被一个 Task 消费。 - Source Connector:Task 数由
taskConfigs返回的列表长度决定,完全由开发者按外部系统分片控制。
# Sink:tasks.max=8,订阅主题共 12 个分区 → 实际 8 个 Task
# 剩余 4 个分区怎么办?没有 Task 消费 → 数据积压
# 结论:Sink 的 tasks.max 必须 >= 订阅主题分区总数,否则必然积压
4.2 task 分配与 rebalance
distributed 模式下,Worker 组成一个组,Task 在 Worker 间分配。Worker 增减会触发 Task 重新分配(Connect 层的 rebalance,与消费组 rebalance 独立):
- Worker 宕机:它名下的 Task 被其他 Worker 接管,从 offset 主题恢复。
- Worker 新增:新 Worker 分走部分 Task,触发整体重分配。
- 重分配期间 Task 会 stop 再 start,因此
stop/start必须可重入、可快速恢复。
4.3 分区级并行与单 Task 取舍
| 场景 | 建议 |
|---|---|
| Sink 写外部数据库 | 按分区并行,但注意外部系统连接数上限 |
| Sink 写顺序敏感的下游 | 减少 Task 数,保证同 key 同分区同 Task |
| Source 拉分页 API | 单 Task 即可,分页本身串行 |
| Source 拉分片存储 | 每分片一个 Task,最大化吞吐 |
并行度不是越高越好:每个 Task 都有独立的连接、缓冲区与线程,Task 数超过外部系统承载能力反而会拖垮下游。
5. Offset 管理与 Exactly-Once
5.1 offset 提交时机
Connect 的 offset 提交由 Worker 统一负责,不在 Task 内直接操作。时序是:Worker 调用 Task.poll() 拿到一批 SourceRecord,把记录写入 Kafka,再把该批最后一条的 sourceOffset 写入 offset 主题,如此循环。若在第 3 步之前崩溃就会重复,在第 3 步之后崩溃则不丢,因此默认语义是 at-least-once。
5.2 at-least-once 与 EOS 的边界
Connect 的 EOS 能力是局部的:
- Source 端:offset 写入与数据写入 Kafka 在同一事务内,可做到精确一次。
- Sink 端:数据写入外部系统与 offset 提交无法共用一个事务(外部系统不参与 Kafka 事务),因此只能做到 at-least-once + 下游幂等。
# Source EOS 配置
# exactly.once.support=required
# 要求 Worker 侧启用事务生产者,offset 与数据原子提交
# Sink 无法真 EOS,只能靠幂等补偿
# 结论:跨系统端到端 EOS = 上游 EOS + 下游幂等
5.3 Sink 端幂等与事务写入
Sink 的幂等实现取决于外部系统能力:
// 幂等 upsert:重复批次不会产生重复行
// INSERT INTO t (topic, partition, offset, payload)
// VALUES (?, ?, ?, ?)
// ON CONFLICT (topic, partition, offset) DO NOTHING;
// 外部系统支持事务时,整批写入与位点记录同一事务
writer.beginTransaction();
try {
for (SinkRecord r : batch) { writer.write(r); }
writer.recordOffsets(batch); // 把位点与数据放同一事务
writer.commit();
} catch (Exception e) {
writer.rollback();
throw e;
}
5.4 Source 端断点续传
Source 的断点续传依赖两点:offset 存储可用、start 时正确恢复。常见坑是把 offset 存在外部系统本地内存里,Worker 重启后从头拉,导致全量重复。正确做法是 start 时通过 offsetStorageReader 读取上次位点,返回 null 才说明是全新 Task,从初始位点开始。
6. 单消息转换链 SMT
6.1 Transform 接口
SMT(Single Message Transform)在 Converter 之前对每条记录做无状态改写,接口极简:
public class MaskFieldTransform<R extends ConnectRecord<R>> implements Transformation<R> {
private List<String> fields;
@Override
public void configure(Map<String, ?> configs) {
this.fields = Arrays.asList(
String.valueOf(configs.get("fields")).split(","));
}
@Override
public R apply(R record) {
if (record.value() == null) {
return record;
}
Struct value = (Struct) record.value();
for (String f : fields) {
if (value.schema().field(f) != null) {
value.put(f, "***");
}
}
return record.newRecord(record.topic(), record.kafkaPartition(),
record.keySchema(), record.key(),
value.schema(), value, record.timestamp());
}
@Override
public ConfigDef config() {
return new ConfigDef().define("fields", ConfigDef.Type.STRING,
ConfigDef.Importance.HIGH, "需要脱敏的字段");
}
}
6.2 常用内置 SMT
| SMT | 作用 |
|---|---|
| MaskField | 字段脱敏 |
| ReplaceField | 字段改名或删除 |
| InsertField | 插入静态字段或元数据 |
| ExtractField | 从 Struct 中抽出单个字段 |
| Flatten | 嵌套结构拍平 |
| Cast | 字段类型转换 |
| TimestampRouter | 按时间戳路由到不同主题 |
| RegexRouter | 按正则改写主题名 |
6.3 链式顺序与 Predicate
多个 SMT 按 transforms 数组顺序执行,前一个的输出是后一个的输入。Predicate 可以为每个 SMT 加条件:
{
"transforms": "maskPwd,route",
"transforms.maskPwd.type": "com.example.MaskFieldTransform",
"transforms.maskPwd.fields": "password,token",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "orders-(.*)",
"transforms.route.replacement": "orders-clean-$1",
"predicates": "isOrder",
"predicates.isOrder.type": "org.apache.kafka.connect.transforms.predicates.TopicNameMatches",
"predicates.isOrder.pattern": "orders-.*",
"transforms.maskPwd.predicate": "isOrder"
}
Predicate 让 SMT 只作用于匹配的记录,未匹配的记录原样通过。顺序很关键:先脱敏再路由,与先路由再脱敏,结果完全不同。
6.4 自定义 SMT 与 Schema 影响
SMT 会改变记录结构,进而影响下游 Schema 兼容性。若启用了 Schema Registry,改写后的结构必须能通过兼容性检查,否则写入会被拒绝。建议把 SMT 的 schema 变更纳入版本管理,并在测试环境先验证兼容性。
7. 打包分发与运维
7.1 插件打包与目录结构
打包时把 Connector 主 jar 与它私有的依赖打进同一个插件目录,但不要打包 Connect 自身已提供的 API(如 connect-api),否则会出现类加载冲突:
# 推荐目录:my-connector/my-connector-1.0.0.jar(主 jar + 私有依赖 shade)
# 打包时排除 connect-api、kafka-clients 等 Worker 已提供的依赖
7.2 REST API 管理
distributed 模式通过 REST 接口管理 Connector 生命周期:
# 创建 Connector
curl -X POST http://connect:8083/connectors \
-H "Content-Type: application/json" \
-d @my-connector.json
# 查看状态与 Task 分布
curl http://connect:8083/connectors/my-connector/status
# 暂停 / 恢复 / 删除
curl -X PUT http://connect:8083/connectors/my-connector/pause
curl -X DELETE http://connect:8083/connectors/my-connector
# 更新配置(注意:会触发 Task 重启)
curl -X PUT http://connect:8083/connectors/my-connector/config \
-H "Content-Type: application/json" -d @new-config.json
7.3 监控指标
Connect 暴露大量 JMX 指标,生产必看这几个:
| 指标 | 含义 | 告警阈值 |
|---|---|---|
| source-record-poll-total | Source 拉取记录总数 | 长期为 0 需排查 |
| source-record-write-total | Source 写入 Kafka 总数 | 与 poll 差值过大 |
| sink-record-send-total | Sink 写出记录总数 | 停滞需排查 |
| offset-commit-success-percentage | 位点提交成功率 | 低于 100 告警 |
| rebalance-total | Task 重分配次数 | 突增需排查 |
| task-count | 运行中的 Task 数 | 低于预期需排查 |
同时用 connector-status 的 state 字段监控(RUNNING / PAUSED / FAILED),FAILED 必须告警。
7.4 升级与常见坑
- 配置更新触发重启:修改任何配置都会重启 Task,务必在低峰期操作,并确认
stop能快速释放资源。 - offset 主题分区数:
offset.storage.topic建议用高分区数(如 25)且cleanup.policy=compact,分区数一旦定下不可随意改。 - Task 数超过分区数:Sink 端 Task 数超过订阅主题分区总数时,多余 Task 空转,浪费资源。
- 依赖冲突:多个 Connector 共用插件目录导致类加载冲突,表现为莫名的
NoSuchMethodError,务必隔离目录。 - 空轮询不退避:Source 的
poll在没有数据时若立即返回,会造成 CPU 空转,必须Thread.sleep退避。 - DLQ 无限堆积:容忍错误后 DLQ 无人消费,问题被掩盖,需为 DLQ 配置独立消费与告警。
Kafka Connect 的部署与内部主题细节可以对照 https://plumephp.com/kafka-connect/,与外部系统的集成模式参考 https://plumephp.com/kafka-connect-integration/;投递语义的取舍与 https://plumephp.com/kafka-delivery-semantics/ 中的端到端分析一致,事务细节见 https://plumephp.com/kafka-transactions/。
8. 总结
自定义 Connector 的开发难点不在接口本身,而在语义设计:Source 的 taskConfigs 决定了并行度上限,poll 的 offset 语义决定了重复窗口,Sink 的 put 批量与幂等决定了端到端一致性。三条主线要记住:并行度上,Sink 按分区切、Source 按分片切,tasks.max 是上限不是目标;一致性上,Source 端可做 EOS,Sink 端只能 at-least-once 加幂等,端到端精确一次必须上下游配合;运维上,插件目录隔离、offset 主题高分区、指标与 DLQ 监控缺一不可。写 Connector 之前先把这三件事想清楚,代码只是最后的落地。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。