📌 本文是「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
数据格式转换器。
内置转换器:
| 转换器 | 说明 |
|---|---|
| JsonConverter | JSON 格式 |
| AvroConverter | Apache Avro |
| ProtobufConverter | Protocol Buffers |
| StringConverter | 字符串 |
| BytesConverter | 字节数组 |
7.3.5 Single Message Transform (SMT)
消息级别的转换,用于对单条记录进行轻量级修改。
内置转换器详解
Cast - 类型转换
将字段转换为指定类型。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| spec | list | 无 | 字段和类型的映射列表,格式为 field1:type,field2:type |
| replace.null.with.default | boolean | true | 是否用默认值替换 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 - 删除消息头
删除指定的消息头。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| headers | list | 无 | 要删除的消息头名称列表 |
transforms=DropHeaders
transforms.DropHeaders.type=org.apache.kafka.connect.transforms.DropHeaders
transforms.DropHeaders.headers=deprecated-header,sensitive-info
ExtractField - 提取字段
从记录值中提取指定字段作为新的值。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| field | string | 无 | 要提取的字段名称 |
| field.syntax.version | string | V1 | 字段访问语法版本:V1(根级别)或 V2(支持嵌套) |
| replace.null.with.default | boolean | true | 是否用默认值替换 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 展平为单层结构。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| delimiter | string | . | 展平时用于连接字段名的分隔符 |
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 - 字段转消息头
将记录中的字段复制或移动到消息头。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| fields | list | 无 | 要复制或移动的字段名称列表 |
| headers | list | 无 | 消息头名称列表(与 fields 对应) |
| operation | string | 无 | 操作类型:move(移动)或 copy(复制) |
| replace.null.with.default | boolean | true | 是否用默认值替换 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 - 提升字段
将整个值包装到单个字段中。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| field | string | 无 | 将创建的字段名称 |
transforms=HoistField
transforms.HoistField.type=org.apache.kafka.connect.transforms.HoistField$Value
transforms.HoistField.field=line
InsertField - 插入字段
插入静态或元数据字段到记录中。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| offset.field | string | null | Kafka 偏移量字段名(仅 Sink 连接器) |
| partition.field | string | null | Kafka 分区字段名 |
| static.field | string | null | 静态数据字段名 |
| static.value | string | null | 静态字段值 |
| timestamp.field | string | null | 记录时间戳字段名 |
| topic.field | string | null | Kafka 主题字段名 |
字段名后可添加后缀:
!- 必需字段?- 可选字段(默认)
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 - 插入消息头
插入静态值到消息头。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| header | string | 无 | 消息头名称 |
| value.literal | string | 无 | 要设置的静态值 |
transforms=InsertHeader
transforms.InsertHeader.type=org.apache.kafka.connect.transforms.InsertHeader
transforms.InsertHeader.header=x-source
transforms.InsertHeader.value.literal=file-connector
MaskField - 掩码字段
用指定值或占位符替换字段值,用于敏感数据脱敏。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| fields | list | 无 | 要掩码的字段名称列表 |
| replacement | string | null | 自定义替换值 |
| replace.null.with.default | boolean | true | 是否用默认值替换 null |
transforms=MaskField
transforms.MaskField.type=org.apache.kafka.connect.transforms.MaskField$Value
transforms.MaskField.fields=credit_card,ssn,password
transforms.MaskField.replacement=****
RegexRouter - 正则路由
根据正则表达式重写主题名称。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| regex | string | 无 | 用于匹配的 Java 正则表达式 |
| replacement | string | 无 | 替换字符串 |
transforms=RegexRouter
transforms.RegexRouter.type=org.apache.kafka.connect.transforms.RegexRouter
transforms.RegexRouter.regex=^(.*)$
transforms.RegexRouter.replacement=prefix-$1
ReplaceField - 替换字段
重命名字段或包含/排除特定字段。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| exclude | list | "" | 要排除的字段(优先级高于 include) |
| include | list | "" | 要包含的字段(指定后只使用这些字段) |
| renames | list | "" | 字段重命名映射,格式 old:new |
| replace.null.with.default | boolean | true | 是否用默认值替换 null |
| blacklist | list | null | 已弃用,使用 exclude 代替 |
| whitelist | list | null | 已弃用,使用 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.name | string | null | Schema 名称 |
| schema.version | int | null | Schema 版本 |
| replace.null.with.default | boolean | true | 是否用默认值替换 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.type | string | 无 | 目标类型:string, unix, Date, Time, Timestamp |
| field | string | "" | 包含时间戳的字段,空则整个值是时间戳 |
| format | string | "" | SimpleDateFormat 兼容的格式字符串 |
| unix.precision | string | milliseconds | Unix 精度:nanoseconds, microseconds, milliseconds, seconds |
| replace.null.with.default | boolean | true | 是否用默认值替换 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.format | string | yyyyMMdd | SimpleDateFormat 兼容的格式字符串 |
| topic.format | string | ${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 - 值转键
从记录值中提取字段作为键。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| fields | list | 无 | 要提取为键的字段名称列表 |
| replace.null.with.default | boolean | true | 是否用默认值替换 null |
transforms=ValueToKey
transforms.ValueToKey.type=org.apache.kafka.connect.transforms.ValueToKey
transforms.ValueToKey.fields=user_id
内置谓词(Predicates)
谓词用于与转换器配合,实现条件性转换。
TopicNameMatches - 主题名称匹配
根据正则表达式匹配主题名称。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| pattern | string | 无 | 用于匹配主题名称的 Java 正则表达式 |
predicates=IsImportant
predicates.IsImportant.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches
predicates.IsImportant.pattern=important-.*
HasHeaderKey - 存在消息头
检查记录是否包含指定的消息头。
配置参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| name | string | 无 | 要检查的消息头名称 |
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
下一步学习
- 第八章:Kafka Streams - 流处理库
- 实践:连接器开发实战
- 官方连接器列表