Kafka Connect 自定义 Connector 开发实战

深入讲解 Kafka Connect 自定义 Connector 的开发路径:Worker、Connector、Task、Converter、Transform 五层扩展点、standalone 与 distributed 部署差异、SourceTask 与 SinkTask 的接口契约、poll 与 put 批处理语义、offset 存储模型、tasks.max 与分区级并行度取舍、offset 提交时机与 Exactly-Once 语义、Sink 端幂等与事务写入、SMT 与 Predicate 组合、插件打包与 REST API 运维,并给出生产参数表、监控指标与踩坑清单

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-totalSource 拉取记录总数长期为 0 需排查
source-record-write-totalSource 写入 Kafka 总数与 poll 差值过大
sink-record-send-totalSink 写出记录总数停滞需排查
offset-commit-success-percentage位点提交成功率低于 100 告警
rebalance-totalTask 重分配次数突增需排查
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 之前先把这三件事想清楚,代码只是最后的落地。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. Kubernetes 上的 Kafka:Strimzi Operator 生产实践
  2. ksqlDB 流式 SQL:流表模型、窗口聚合与生产运维
  3. Kafka 消费延迟诊断:Lag 定位、分区倾斜与治理