跳到正文
SL Blog 技术探索 · 工程实践 · AI 时代思考
返回
📨 Kafka 3.9 教程 · 8 / 10 查看系列简介 →

Kafka 3.9 教程(七):Kafka Connect

文章目录
  1. 目录
  2. 7.1 架构概述
  3. 7.1.1 什么是 Kafka Connect?
  4. 7.1.2 连接器类型
  5. 7.1.3 Kafka Connect 生态系统
  6. 7.2 运行模式
  7. 7.2.1 独立模式(Standalone)
  8. 7.2.2 分布式模式(Distributed)
  9. 7.2.3 模式对比
  10. 7.3 核心概念
  11. 7.3.1 Worker
  12. 7.3.2 Connector
  13. 7.3.3 Task
  14. 7.3.4 Converter
  15. 7.3.5 Single Message Transform (SMT)
  16. 7.4 用户指南
  17. 7.4.1 快速开始
  18. 7.4.2 File Source 示例
  19. 7.4.3 JDBC Source 示例
  20. 7.4.4 JDBC Sink 示例
  21. 7.5 连接器开发
  22. 7.5.1 Maven 依赖
  23. 7.5.2 Source Connector 开发
  24. 7.5.3 Sink Connector 开发
  25. 7.6 管理运维
  26. 7.6.1 REST API
  27. 7.6.2 连接器状态
  28. 7.6.3 监控
  29. 7.6.4 故障排查
  30. 附录:常用连接器配置
  31. 通用配置
  32. 错误处理
  33. 下一步学习

📌 本文是「Kafka 3.9 教程」系列第 七 篇 · 系列目录

本文档基于 Apache Kafka 3.9 官方文档整理

目录


7.1 架构概述

7.1.1 什么是 Kafka Connect?

Kafka Connect 是 Kafka 官方提供的数据集成框架,用于在 Kafka 与外部系统之间可靠地传输数据。

核心特性:

  • 可扩展:分布式架构,支持水平扩展
  • 可靠:支持Exactly-Once 语义
  • 灵活:支持多种运行模式和配置
  • 易于使用:提供 REST API 和配置文件

7.1.2 连接器类型

┌─────────────────────────────────────────────────────────────┐
│                     Kafka Connect                            │
├─────────────────────────────────────────────────────────────┤
│                                                              │
│   ┌─────────────────────┐      ┌─────────────────────┐      │
│   │   Source Connector  │      │   Sink Connector    │      │
│   │   (数据导入)         │      │   (数据导出)         │      │
│   │                     │      │                     │      │
│   │  外部系统 → Kafka   │      │  Kafka → 外部系统   │      │
│   └─────────────────────┘      └─────────────────────┘      │
│                                                              │
└─────────────────────────────────────────────────────────────┘

Source Connector

将数据从外部系统导入 Kafka。

示例:

  • JDBC Source:从数据库读取数据变更
  • File Source:从文件读取数据
  • Debezium:捕获数据库变更事件
  • HTTP Source:从 HTTP 端点拉取数据

Sink Connector

将数据从 Kafka 导出到外部系统。

示例:

  • JDBC Sink:写入数据库
  • File Sink:写入文件
  • Elasticsearch Sink:写入搜索引擎
  • S3 Sink:写入对象存储

7.1.3 Kafka Connect 生态系统

┌─────────────────────────────────────────────────────────────┐
│                     外部系统                                  │
│                                                              │
│   数据库        文件系统       云服务       消息队列          │
│   ┌─────┐      ┌─────┐      ┌─────┐      ┌─────┐          │
│   │MySQL│      │  HDFS│      │  S3 │      │Redis│          │
│   └─────┘      └─────┘      └─────┘      └─────┘          │
│      │            │            │            │               │
│      └────────────┴─────┬──────┴────────────┘               │
│                         │                                    │
│                   ┌─────┴─────┐                              │
│                   │ 连接器     │                              │
│                   └─────┬─────┘                              │
│                         │                                    │
└─────────────────────────┼────────────────────────────────────┘
                          │
┌─────────────────────────┼────────────────────────────────────┐
│                         ▼                                    │
│                   ┌─────────────┐                            │
│                   │  Kafka      │                            │
│                   │  Connect    │                            │
│                   └─────────────┘                            │
│                         │                                    │
│                         ▼                                    │
│                   ┌─────────────┐                            │
│                   │   Kafka     │                            │
│                   │   Cluster   │                            │
│                   └─────────────┘                            │
└─────────────────────────────────────────────────────────────┘

7.2 运行模式

7.2.1 独立模式(Standalone)

特点:

  • 单进程运行
  • 简单配置,适合开发和测试
  • 不支持容错
  • 偏移量存储在本地文件

适用场景:

  • 开发测试
  • 一次性数据迁移
  • 资源受限环境

配置示例:

# connect-standalone.properties
bootstrap.servers=localhost:9092
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

# 偏移量存储
offset.storage.file.filename=/tmp/connect.offsets

# 插件路径
plugin.path=/opt/kafka/plugins

# 任务数
tasks.max=1

7.2.2 分布式模式(Distributed)

特点:

  • 多节点分布式运行
  • 支持自动负载均衡
  • 支持故障转移
  • 偏移量存储在 Kafka 主题

适用场景:

  • 生产环境
  • 需要高可用
  • 大规模数据集成

配置示例:

# connect-distributed.properties
bootstrap.servers=localhost: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
offset.storage.replication.factor=3

config.storage.topic=connect-configs
config.storage.replication.factor=3

status.storage.topic=connect-status
status.storage.replication.factor=3

# REST API
rest.port=8083
rest.advertised.host.name=connect-host

7.2.3 模式对比

特性独立模式分布式模式
进程数单进程多节点
容错无自动故障转移
扩展性有限水平扩展
偏移量存储本地文件Kafka 主题
适用场景开发/测试生产环境

7.3 核心概念

7.3.1 Worker

运行连接器的进程。

Worker 职责:

  • 管理连接器生命周期
  • 协调任务分配
  • 处理 REST API 请求
  • 管理偏移量

7.3.2 Connector

定义数据如何传输的逻辑配置。

Connector 职责:

  • 定义数据源/目标
  • 创建和管理任务
  • 处理配置变更

7.3.3 Task

实际执行数据传输的工作单元。

Task 职责:

  • 读取/写入数据
  • 数据转换
  • 错误处理

7.3.4 Converter

数据格式转换器。

内置转换器:

转换器说明
JsonConverterJSON 格式
AvroConverterApache Avro
ProtobufConverterProtocol Buffers
StringConverter字符串
BytesConverter字节数组

7.3.5 Single Message Transform (SMT)

消息级别的转换,用于对单条记录进行轻量级修改。

内置转换器详解

Cast - 类型转换

将字段转换为指定类型。

配置参数:

参数类型默认值说明
speclist无字段和类型的映射列表,格式为 field1:type,field2:type
replace.null.with.defaultbooleantrue是否用默认值替换 null

有效类型:int8, int16, int32, int64, float32, float64, boolean, string

transforms=Cast
transforms.Cast.type=org.apache.kafka.connect.transforms.Cast$Value
transforms.Cast.spec=price:float64,quantity:int32,name:string
DropHeaders - 删除消息头

删除指定的消息头。

配置参数:

参数类型默认值说明
headerslist无要删除的消息头名称列表
transforms=DropHeaders
transforms.DropHeaders.type=org.apache.kafka.connect.transforms.DropHeaders
transforms.DropHeaders.headers=deprecated-header,sensitive-info
ExtractField - 提取字段

从记录值中提取指定字段作为新的值。

配置参数:

参数类型默认值说明
fieldstring无要提取的字段名称
field.syntax.versionstringV1字段访问语法版本:V1(根级别)或 V2(支持嵌套)
replace.null.with.defaultbooleantrue是否用默认值替换 null
transforms=ExtractField
transforms.ExtractField.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.ExtractField.field=user.name
transforms.ExtractField.field.syntax.version=V2
Flatten - 展平嵌套结构

将嵌套的 Map 或 Struct 展平为单层结构。

配置参数:

参数类型默认值说明
delimiterstring.展平时用于连接字段名的分隔符
transforms=Flatten
transforms.Flatten.type=org.apache.kafka.connect.transforms.Flatten$Value
transforms.Flatten.delimiter=_
Filter - 过滤记录

根据谓词条件过滤掉不符合条件的记录。

配置参数:

Filter 没有特定参数,通过谓词(predicate)配置过滤条件。

transforms=Filter
transforms.Filter.type=org.apache.kafka.connect.transforms.Filter
transforms.Filter.predicate=IsImportant
HeaderFrom - 字段转消息头

将记录中的字段复制或移动到消息头。

配置参数:

参数类型默认值说明
fieldslist无要复制或移动的字段名称列表
headerslist无消息头名称列表(与 fields 对应)
operationstring无操作类型:move(移动)或 copy(复制)
replace.null.with.defaultbooleantrue是否用默认值替换 null
transforms=HeaderFrom
transforms.HeaderFrom.type=org.apache.kafka.connect.transforms.HeaderFrom$Value
transforms.HeaderFrom.fields=source,timestamp
transforms.HeaderFrom.headers=x-source,x-timestamp
transforms.HeaderFrom.operation=copy
HoistField - 提升字段

将整个值包装到单个字段中。

配置参数:

参数类型默认值说明
fieldstring无将创建的字段名称
transforms=HoistField
transforms.HoistField.type=org.apache.kafka.connect.transforms.HoistField$Value
transforms.HoistField.field=line
InsertField - 插入字段

插入静态或元数据字段到记录中。

配置参数:

参数类型默认值说明
offset.fieldstringnullKafka 偏移量字段名(仅 Sink 连接器)
partition.fieldstringnullKafka 分区字段名
static.fieldstringnull静态数据字段名
static.valuestringnull静态字段值
timestamp.fieldstringnull记录时间戳字段名
topic.fieldstringnullKafka 主题字段名

字段名后可添加后缀:

  • ! - 必需字段
  • ? - 可选字段(默认)
transforms=InsertField
transforms.InsertField.type=org.apache.kafka.connect.transforms.InsertField$Value
transforms.InsertField.static.field=data_source
transforms.InsertField.static.value=test-file-source
transforms.InsertField.timestamp.field=processed_at
InsertHeader - 插入消息头

插入静态值到消息头。

配置参数:

参数类型默认值说明
headerstring无消息头名称
value.literalstring无要设置的静态值
transforms=InsertHeader
transforms.InsertHeader.type=org.apache.kafka.connect.transforms.InsertHeader
transforms.InsertHeader.header=x-source
transforms.InsertHeader.value.literal=file-connector
MaskField - 掩码字段

用指定值或占位符替换字段值,用于敏感数据脱敏。

配置参数:

参数类型默认值说明
fieldslist无要掩码的字段名称列表
replacementstringnull自定义替换值
replace.null.with.defaultbooleantrue是否用默认值替换 null
transforms=MaskField
transforms.MaskField.type=org.apache.kafka.connect.transforms.MaskField$Value
transforms.MaskField.fields=credit_card,ssn,password
transforms.MaskField.replacement=****
RegexRouter - 正则路由

根据正则表达式重写主题名称。

配置参数:

参数类型默认值说明
regexstring无用于匹配的 Java 正则表达式
replacementstring无替换字符串
transforms=RegexRouter
transforms.RegexRouter.type=org.apache.kafka.connect.transforms.RegexRouter
transforms.RegexRouter.regex=^(.*)$
transforms.RegexRouter.replacement=prefix-$1
ReplaceField - 替换字段

重命名字段或包含/排除特定字段。

配置参数:

参数类型默认值说明
excludelist""要排除的字段(优先级高于 include)
includelist""要包含的字段(指定后只使用这些字段)
renameslist""字段重命名映射,格式 old:new
replace.null.with.defaultbooleantrue是否用默认值替换 null
blacklistlistnull已弃用,使用 exclude 代替
whitelistlistnull已弃用,使用 include 代替
transforms=ReplaceField
transforms.ReplaceField.type=org.apache.kafka.connect.transforms.ReplaceField$Value
transforms.ReplaceField.exclude=deprecated_field
transforms.ReplaceField.renames=old_name:new_name,user_id:userId
SetSchemaMetadata - 设置 Schema 元数据

设置 Schema 的名称和版本。

配置参数:

参数类型默认值说明
schema.namestringnullSchema 名称
schema.versionintnullSchema 版本
replace.null.with.defaultbooleantrue是否用默认值替换 null
transforms=SetSchemaMetadata
transforms.SetSchemaMetadata.type=org.apache.kafka.connect.transforms.SetSchemaMetadata$Value
transforms.SetSchemaMetadata.schema.name=io.example.record
transforms.SetSchemaMetadata.schema.version=2
TimestampConverter - 时间戳转换

在时间戳的不同表示形式之间转换。

配置参数:

参数类型默认值说明
target.typestring无目标类型:string, unix, Date, Time, Timestamp
fieldstring""包含时间戳的字段,空则整个值是时间戳
formatstring""SimpleDateFormat 兼容的格式字符串
unix.precisionstringmillisecondsUnix 精度:nanoseconds, microseconds, milliseconds, seconds
replace.null.with.defaultbooleantrue是否用默认值替换 null
transforms=TimestampConverter
transforms.TimestampConverter.type=org.apache.kafka.connect.transforms.TimestampConverter$Value
transforms.TimestampConverter.target.type=string
transforms.TimestampConverter.field=timestamp
transforms.TimestampConverter.format=yyyy-MM-dd HH:mm:ss
TimestampRouter - 时间戳路由

根据时间戳重写主题名称。

配置参数:

参数类型默认值说明
timestamp.formatstringyyyyMMddSimpleDateFormat 兼容的格式字符串
topic.formatstring${topic}-${timestamp}主题格式,可包含 ${topic} 和 ${timestamp} 占位符
transforms=TimestampRouter
transforms.TimestampRouter.type=org.apache.kafka.connect.transforms.TimestampRouter
transforms.TimestampRouter.topic.format=${topic}-${timestamp}
transforms.TimestampRouter.timestamp.format=yyyyMMdd
ValueToKey - 值转键

从记录值中提取字段作为键。

配置参数:

参数类型默认值说明
fieldslist无要提取为键的字段名称列表
replace.null.with.defaultbooleantrue是否用默认值替换 null
transforms=ValueToKey
transforms.ValueToKey.type=org.apache.kafka.connect.transforms.ValueToKey
transforms.ValueToKey.fields=user_id

内置谓词(Predicates)

谓词用于与转换器配合,实现条件性转换。

TopicNameMatches - 主题名称匹配

根据正则表达式匹配主题名称。

配置参数:

参数类型默认值说明
patternstring无用于匹配主题名称的 Java 正则表达式
predicates=IsImportant
predicates.IsImportant.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.IsImportant.pattern=important-.*
HasHeaderKey - 存在消息头

检查记录是否包含指定的消息头。

配置参数:

参数类型默认值说明
namestring无要检查的消息头名称
predicates=HasHeader
predicates.HasHeader.type=org.apache.kafka.connect.transforms.predicates.HasHeaderKey
predicates.HasHeader.name=x-priority
RecordIsTombstone - 记录是墓碑

检查记录是否为墓碑记录(null 值)。

predicates=IsTombstone
predicates.IsTombstone.type=org.apache.kafka.connect.transforms.predicates.RecordIsTombstone

7.4 用户指南

7.4.1 快速开始

1. 配置连接器

# connect-file-source.properties
name=file-source
connector.class=FileStreamSourceConnector
file=/tmp/test.txt
topic=connect-test

# 转换配置
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter

2. 启动连接器

# 独立模式
bin/connect-standalone.sh \
  config/connect-standalone.properties \
  config/connect-file-source.properties

# 分布式模式(通过 REST API)
curl -X POST http://localhost:8083/connectors \
  -H "Content-Type: application/json" \
  -d '{
    "name": "file-source",
    "config": {
      "connector.class": "FileStreamSourceConnector",
      "file": "/tmp/test.txt",
      "topic": "connect-test"
    }
  }'

7.4.2 File Source 示例

# connect-file-source.properties
name=file-source
connector.class=FileStreamSourceConnector
file=/tmp/data/input.txt
topic=file-topic

# 任务配置
tasks.max=3

# 转换配置
transforms=makeSingleMessage
transforms.makeSingleMessage.type=org.apache.kafka.connect.transforms.RegexRouter
transforms.makeSingleMessage.regex=".*"
transforms.makeSingleMessage.replacement="one-message-per-line"

7.4.3 JDBC Source 示例

# connect-jdbc-source.properties
name=jdbc-source
connector.class=io.confluent.connect.jdbc.JdbcSourceConnector
connection.url=jdbc:mysql://localhost:3306/mydb
connection.user=root
connection.password=password

# 查询配置
topic.prefix=mydb_
mode=incrementing
incrementing.column.name=id

# 白名单
table.whitelist=users,orders,products

# 批量大小
batch.max.rows=100

7.4.4 JDBC Sink 示例

# connect-jdbc-sink.properties
name=jdbc-sink
connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
connection.url=jdbc:mysql://localhost:3306/mydb
connection.user=root
connection.password=password

topics=orders-topic
table.name.format=orders

# 插入模式
insert.mode=upsert
pk.mode=record_key

# 主键字段
pk.fields=order_id

7.5 连接器开发

7.5.1 Maven 依赖

<dependencies>
    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>connect-api</artifactId>
        <version>3.9.1</version>
    </dependency>
</dependencies>

7.5.2 Source Connector 开发

import org.apache.kafka.common.config.ConfigDef;
import org.apache.kafka.connect.connector.Task;
import org.apache.kafka.connect.source.SourceConnector;
import org.apache.kafka.connect.source.SourceRecord;

import java.util.ArrayList;
import java.util.List;
import java.util.Map;

public class MySourceConnector extends SourceConnector {

    private Map<String, String> config;

    @Override
    public String version() {
        return "1.0.0";
    }

    @Override
    public void start(Map<String, String> props) {
        this.config = props;
        // 初始化连接
    }

    @Override
    public Class<? extends Task> taskClass() {
        return MySourceTask.class;
    }

    @Override
    public List<Map<String, String>> taskConfigs(int maxTasks) {
        List<Map<String, String>> configs = new ArrayList<>();
        for (int i = 0; i < maxTasks; i++) {
            configs.add(config);
        }
        return configs;
    }

    @Override
    public void stop() {
        // 清理资源
    }

    @Override
    public ConfigDef config() {
        return new ConfigDef()
            .define("topic", ConfigDef.Type.STRING, ConfigDef.NO_DEFAULT_VALUE,
                ConfigDef.Importance.HIGH, "Target topic");
    }
}
import org.apache.kafka.connect.data.Schema;
import org.apache.kafka.connect.source.SourceRecord;
import org.apache.kafka.connect.source.SourceTask;

import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.List;
import java.util.Map;
import java.util.ArrayList;

public class MySourceTask extends SourceTask {

    private String topic;
    private String filePath;

    @Override
    public void start(Map<String, String> props) {
        topic = props.get("topic");
        filePath = props.get("file");
    }

    @Override
    public List<SourceRecord> poll() {
        List<SourceRecord> records = new ArrayList<>();

        try {
            List<String> lines = Files.readAllLines(Paths.get(filePath));
            for (String line : lines) {
                SourceRecord record = new SourceRecord(
                    null, null,  // source partition and offset
                    topic,
                    null,  // key
                    Schema.STRING_SCHEMA,
                    line,  // value
                    System.currentTimeMillis()
                );
                records.add(record);
            }
        } catch (Exception e) {
            throw new RuntimeException(e);
        }

        return records;
    }

    @Override
    public void stop() {
        // 清理资源
    }

    @Override
    public String version() {
        return "1.0.0";
    }
}

7.5.3 Sink Connector 开发

import org.apache.kafka.common.config.ConfigDef;
import org.apache.kafka.connect.connector.Task;
import org.apache.kafka.connect.sink.SinkConnector;
import org.apache.kafka.connect.sink.SinkRecord;

import java.util.ArrayList;
import java.util.List;
import java.util.Map;

public class MySinkConnector extends SinkConnector {

    private Map<String, String> config;

    @Override
    public String version() {
        return "1.0.0";
    }

    @Override
    public void start(Map<String, String> props) {
        this.config = props;
        // 初始化连接
    }

    @Override
    public Class<? extends Task> taskClass() {
        return MySinkTask.class;
    }

    @Override
    public List<Map<String, String>> taskConfigs(int maxTasks) {
        List<Map<String, String>> configs = new ArrayList<>();
        for (int i = 0; i < maxTasks; i++) {
            configs.add(config);
        }
        return configs;
    }

    @Override
    public void stop() {
        // 清理资源
    }

    @Override
    public ConfigDef config() {
        return new ConfigDef()
            .define("destination", ConfigDef.Type.STRING, ConfigDef.NO_DEFAULT_VALUE,
                ConfigDef.Importance.HIGH, "Destination path");
    }
}
import org.apache.kafka.connect.sink.SinkRecord;
import org.apache.kafka.connect.sink.SinkTask;

import java.io.FileWriter;
import java.io.IOException;
import java.util.Collection;
import java.util.Map;

public class MySinkTask extends SinkTask {

    private String destination;

    @Override
    public void start(Map<String, String> props) {
        destination = props.get("destination");
    }

    @Override
    public void put(Collection<SinkRecord> records) {
        try (FileWriter writer = new FileWriter(destination, true)) {
            for (SinkRecord record : records) {
                writer.write(record.value().toString() + "\n");
            }
        } catch (IOException e) {
            throw new RuntimeException(e);
        }
    }

    @Override
    public void stop() {
        // 清理资源
    }

    @Override
    public String version() {
        return "1.0.0";
    }
}

7.6 管理运维

7.6.1 REST API

方法端点说明
GET/connectors列出所有连接器
GET/connectors/{name}获取连接器详情
POST/connectors创建连接器
PUT/connectors/{name}/config更新连接器配置
DELETE/connectors/{name}删除连接器
GET/connectors/{name}/status获取连接器状态
POST/connectors/{name}/restart重启连接器
GET/connectors/{name}/tasks获取任务列表

7.6.2 连接器状态

状态说明
UNASSIGNED任务未分配
RUNNING运行中
PAUSED已暂停
FAILED失败
DESTROYED已销毁

7.6.3 监控

关键指标:

指标说明
connector-total-task-count连接器任务总数
connector-failed-task-count失败任务数
connector-paused-task-count暂停任务数
task-startup-success-total启动成功次数
task-startup-failure-total启动失败次数

7.6.4 故障排查

# 查看连接器状态
curl -s http://localhost:8083/connectors/file-source/status

# 查看任务配置
curl -s http://localhost:8083/connectors/file-source/tasks

# 查看 Worker 日志
tail -f /opt/kafka/logs/connect.log

# 常见问题:
# 1. 连接器启动失败:检查配置和权限
# 2. 任务卡住:检查外部系统响应
# 3. 数据丢失:检查偏移量配置

附录:常用连接器配置

通用配置

# 基础配置
name=my-connector
connector.class=com.example.MyConnector

# 任务配置
tasks.max=3

# 转换配置
transforms=transform1,transform2
transforms.transform1.type=org.apache.kafka.connect.transforms.InsertField$Value
transforms.transform1.timestamp.field=@timestamp

错误处理

# 死信队列
errors.tolerance=all
errors.deadletterqueue.topic.name=dlq-topic
errors.deadletterqueue.context.headers.enable=true

# 重试配置
errors.retry.timeout=300000
errors.retry.delay.initial.ms=60000

下一步学习


RAG 智能问答

针对本文继续提问:《Kafka 3.9 教程(七):Kafka Connect》