06. Kafka Connect 数据集成

Kafka Connect Source/Sink Connector、Debezium CDC、数据管道构建与 Kafka Connect 生产实践。

1. Kafka Connect 概述

Kafka Connect 是 Kafka 官方提供的数据集成框架,用于在 Kafka 与外部系统(数据库、文件、对象存储、搜索引擎等)之间建立可靠的数据管道。它通过插件化的 Connector 机制,无需编写代码即可实现数据导入导出。

1.1 Connect 架构

┌─────────────────────────────────────────────────────────────┐
│                    Kafka Connect 集群                        │
│                                                             │
│  ┌──────────────┐      ┌──────────────┐      ┌──────────┐  │
│  │ Worker 1     │      │ Worker 2     │      │ Worker 3 │  │
│  │              │      │              │      │          │  │
│  │ [Connector A]│      │ [Connector B]│      │(备用)   │  │
│  │   Task 1     │      │   Task 1     │      │          │  │
│  │   Task 2     │      │   Task 2     │      │          │  │
│  └──────────────┘      └──────────────┘      └──────────┘  │
│         │                     │                             │
└─────────┼─────────────────────┼─────────────────────────────┘
          │                     │
    Source Connector      Sink Connector
          │                     │
  ┌───────▼──────┐       ┌─────▼──────┐
  │   外部系统    │       │  外部系统   │
  │  MySQL/Mongo │       │ ES/S3/HDFS │
  └──────────────┘       └────────────┘

1.2 Connector、Task、Worker 的关系

组件角色说明
Connector逻辑定义定义数据源/目标的配置,进行 Task 拆分
Task执行单元真正执行数据复制的工作单元,无状态
Worker运行进程承载 Task 运行的 JVM 进程,支持集群和分布式调度

2. Source Connector

Source Connector 从外部系统读取数据,写入 Kafka Topic。

2.1 Debezium CDC(变更数据捕获)

Debezium 是最常用的 Source Connector,基于数据库日志(binlog/WAL)实现增量数据捕获。

Debezium 工作原理:

MySQL Binlog               Debezium          Kafka
  INSERT ──→ [捕获] ──→  ──→ Avro/JSON ──→ db.shop.users
  UPDATE ──→ [解析] ──→  ──→ 变更事件   ──→ db.shop.orders
  DELETE ──→ [转换] ──→  ──→           ──→ db.shop.inventory

初始快照(Snapshot):
  首次启动时,全量读取表数据 → 写入 Kafka → 然后切换到 Binlog 增量

2.2 Debezium 事件格式

{
  "before": { "id": 1, "name": "Alice", "age": 30 },
  "after":  { "id": 1, "name": "Alice", "age": 31 },
  "source": {
    "version": "2.5.0.Final",
    "connector": "mysql",
    "name": "db-shop",
    "ts_ms": 1705315200000,
    "db": "shop",
    "table": "users",
    "pos": "binlog.001:1234",
    "row": 0
  },
  "op": "u",
  "ts_ms": 1705315200001
}
字段说明
before更新/删除前的数据
after更新/插入后的数据
source元信息(数据库、表、binlog 位置)
op操作类型:c=create, u=update, d=delete, r=read(快照)
ts_ms处理时间戳

2.3 Debezium 配置

{
  "name": "mysql-source-connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "database.hostname": "mysql",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "dbz",
    "database.server.id": "184054",
    "database.server.name": "db-shop",
    "database.include.list": "shop",
    "table.include.list": "shop.users,shop.orders",
    "snapshot.mode": "initial",
    "topic.prefix": "cdc",
    "schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
    "schema.history.internal.kafka.topic": "schema-changes.shop",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.drop.tombstones": "false",
    "transforms.unwrap.delete.handling.mode": "rewrite"
  }
}

3. Sink Connector

Sink Connector 从 Kafka Topic 读取数据,写入外部系统。

3.1 JDBC Sink Connector

{
  "name": "jdbc-sink-connector",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "tasks.max": "3",
    "topics": "orders",
    "connection.url": "jdbc:postgresql://postgres:5432/analytics",
    "connection.user": "analytics",
    "connection.password": "secret",
    "auto.create": "true",
    "auto.evolve": "true",
    "insert.mode": "upsert",
    "pk.fields": "order_id",
    "pk.mode": "record_key"
  }
}

3.2 Elasticsearch Sink Connector

{
  "name": "es-sink-connector",
  "config": {
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "tasks.max": "3",
    "topics": "orders",
    "connection.url": "http://elasticsearch:9200",
    "type.name": "_doc",
    "key.ignore": "false",
    "schema.ignore": "true",
    "transforms": "extractKey",
    "transforms.extractKey.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
    "transforms.extractKey.field": "order_id"
  }
}

4. 常用 Connector 列表

Connector类型来源用途
DebeziumSourceDebezium 社区MySQL/PostgreSQL/MongoDB CDC
JDBC SourceSourceConfluent轮询查询数据库
FileStreamSourceKafka 官方读取日志文件
S3 SourceSourceConfluent从 S3 读取数据
JDBC SinkSinkConfluent写入关系型数据库
Elasticsearch SinkSinkConfluent写入 ES 索引
S3 SinkSinkConfluent写入 S3(Parquet/JSON)
HDFS SinkSinkConfluent写入 HDFS
Redis SinkSink社区写入 Redis
InfluxDB SinkSink社区写入时序数据库

5. Kafka Connect 运行时

5.1 Standalone vs Distributed

模式适用场景特点
Standalone开发测试、简单任务单进程,配置存储在本地文件
Distributed生产环境多 Worker 集群,配置存储在 Kafka topic(offset.storage.topic)

5.2 部署 Distributed 模式

# connect-distributed.properties
bootstrap.servers=kafka1:9092,kafka2:9092,kafka3:9092
group.id=connect-cluster

# 配置存储topic
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-status

# 启动 Worker
connect-distributed.sh config/connect-distributed.properties

# 注册 Connector(REST API)
curl -X POST http://localhost:8083/connectors \
  -H "Content-Type: application/json" \
  -d @mysql-source-connector.json

# 查看状态
curl http://localhost:8083/connectors/mysql-source-connector/status

6. Transform(数据转换)

Kafka Connect 支持单消息转换(SMT,Single Message Transform),在数据流转过程中转换格式。

{
  "transforms": "maskField,extractValue",
  "transforms.maskField.type": "org.apache.kafka.connect.transforms.MaskField$Value",
  "transforms.maskField.fields": "password,ssn",
  "transforms.extractValue.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
  "transforms.extractValue.field": "id"
}
Transform作用
ExtractField提取嵌套字段
ReplaceField重命名/删除字段
MaskField脱敏(替换为 null 或掩码)
ValueToKey将 value 中的字段设为 key
TimestampConverter时间格式转换
Filter按条件过滤消息

7. 生产实践

7.1 高可用配置

Worker 数量 ≥ 2(防止单点)
tasks.max = min(分区数, Worker 数量 × 单 Worker 承载能力)

监控指标:
- source-record-active-count(活跃记录数)
- put-batch-avg-time-ms(批处理平均时间)
- sink-record-send-rate(发送速率)

7.2 常见问题

问题原因解决
数据延迟Sink 消费慢增加 tasks.max、优化目标数据库写入
字段不匹配Schema 变更启用 auto.evolve、使用 Schema Registry
消息格式错误非 AVro/JSON 格式配置正确的 Converter(key.converter、value.converter)
连接断开网络不稳定配置重试策略、连接池保持

延伸阅读

继续阅读

探索更多技术文章

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

全部文章 返回首页

「kafka」更多文章

  1. 事件驱动架构:Event Sourcing、CQRS 与 Saga 模式
  2. Kafka 运维监控与故障恢复:JMX 指标、Lag 监控与分区重分配
  3. Kafka 详解:分布式日志系统、ISR 与一致性保证