Article

日志采集 Debezium

更新于:2026-07-13

第一章:Debezium 入门概述

1.1 什么是 Debezium

概念名称说明注意事项
Debezium一个开源的分布式平台,用于捕获数据库的变更数据(Change Data Capture, CDC),基于数据库的事务日志(如 MySQL binlog、PostgreSQL WAL)实时流式输出数据变更事件。Debezium 本身不处理数据消费逻辑,而是通过 Kafka Connect 或 Debezium Server 将变更事件发送到消息系统(如 Kafka)。
CDC(变更数据捕获)一种技术,用于捕获数据库中行级别的 INSERT、UPDATE、DELETE 操作,并将这些变更作为事件流发布。CDC 不依赖于应用层代码,是数据库层面的”被动监听”,对业务侵入性低。
实时数据同步Debezium 可实现数据库到其他系统(如数据仓库、缓存、搜索引擎)的低延迟同步。需确保目标系统具备处理流式数据的能力。
事件驱动架构(EDA)Debezium 是构建事件驱动系统的关键组件,数据库变更即事件,驱动后续服务响应。需设计良好的事件消费逻辑,避免重复处理或丢失。

1.2 Debezium 的核心原理与优势

概念名称说明注意事项
基于日志的 CDCDebezium 直接读取数据库的事务日志(如 MySQL 的 binlog、PostgreSQL 的 WAL),无需修改表结构或触发器。数据库必须启用对应的日志功能(如 binlog_format=ROW)。
低延迟变更发生后,通常在毫秒级内即可捕获并发布事件。网络、Kafka 吞吐量、消费者处理速度会影响端到端延迟。
高可靠性支持断点续传(通过 offset 机制),即使连接器重启也能从上次位置继续读取。必须正确配置 offset 存储(如 Kafka 或文件),避免数据重复或丢失。
支持多种数据库提供对 MySQL、PostgreSQL、MongoDB、SQL Server、Oracle、Db2 等主流数据库的支持。不同数据库的配置和权限要求不同,需分别准备。
与 Kafka 生态无缝集成原生基于 Kafka Connect 构建,可轻松与 Kafka Streams、Flink、Spark 等流处理框架集成。若不使用 Kafka,可选择 Debezium Server 模式。
快照机制(Snapshot)首次启动时可对现有数据进行快照,确保历史数据不丢失,之后继续增量捕获。快照期间可能对数据库产生负载,建议在低峰期执行。

1.3 典型应用场景与架构角色

应用场景说明注意事项
实时数据仓库同步将业务数据库(OLTP)的变更实时同步到数据仓库(OLAP),用于实时分析。需处理 schema 变更和数据类型映射问题。
缓存失效与更新当数据库记录更新时,自动清除或更新 Redis 缓存,保证一致性。可结合事件中的 key 信息精准失效缓存。
搜索引擎索引更新将 MySQL 数据变更同步到 Elasticsearch,实现搜索数据的实时更新。需处理 DELETE 事件对应的索引删除操作。
微服务间数据共享避免服务间直接数据库耦合,通过事件流共享数据变更。应定义清晰的事件格式与版本控制策略。
数据备份与容灾将主库变更实时复制到备用系统,支持快速故障转移。需确保备用系统具备回放能力。
审计与合规记录所有数据变更操作,用于审计追踪。建议保留事件日志较长时间,并加密存储。
架构角色说明注意事项
Source System数据变更的源头数据库(如 MySQL)。必须允许 Debezium 账号读取事务日志。
Debezium Connector负责连接数据库、读取日志、生成变更事件。每个数据库实例通常对应一个连接器。
Kafka / Message Broker作为事件的中转站,实现解耦与缓冲。需合理设计 topic 分区与保留策略。
Consumers消费变更事件的服务,如 Flink 作业、Elasticsearch 同步程序等。消费者应具备幂等性处理能力。

1.4 与其他 CDC 工具的对比

工具名称说明优势劣势注意事项
Debezium开源 CDC 平台,基于 Kafka Connect,支持多数据库。社区活跃,文档完善,与 Kafka 生态深度集成。依赖 Kafka 生态,部署复杂度较高。适合已使用 Kafka 的企业。
Maxwell仅支持 MySQL 的 CDC 工具,输出 JSON 格式事件到 Kafka、Kinesis 等。轻量级,配置简单,易于调试。仅支持 MySQL,功能相对单一。适合 MySQL 单一场景的轻量需求。
Canal阿里开源的 MySQL CDC 工具,基于 Java 开发。在国内应用广泛,支持多种输出方式。主要面向 Java 生态,社区国际化较弱。适合阿里技术栈企业。
AWS DMSAWS 提供的托管式数据迁移与 CDC 服务。全托管,支持异构数据库迁移,图形化界面。成本较高,锁定云厂商。适合云上环境且预算充足的团队。
Flink CDC基于 Flink SQL 的 CDC 方案,可直接在 Flink 中捕获变更。无需 Kafka Connect,端到端流处理一体化。依赖 Flink 运行时,学习成本高。适合已使用 Flink 的流处理场景。

第二章:Debezium 核心架构与组件

2.1 Connectors(连接器)概述

连接器类型支持数据库说明注意事项
MySQL ConnectorMySQL基于 binlog ROW 模式捕获变更,支持快照和增量。需配置 server-idbinlog_format=ROW
PostgreSQL ConnectorPostgreSQL使用逻辑复制(Logical Replication)读取 WAL 日志。需创建复制槽(Replication Slot)和复制用户。
MongoDB ConnectorMongoDB基于变更流(Change Streams)捕获副本集或分片集群的变更。必须使用副本集部署,单节点不支持。
SQL Server ConnectorSQL Server使用变更数据捕获(CDC)或变更跟踪(CT)功能。需启用数据库和表级别的 CDC。
Oracle ConnectorOracle基于 LogMiner 或 XStream API 读取 redo log。配置复杂,需 Oracle 高级权限。
Db2 ConnectorIBM Db2支持 LUW(Linux/Unix/Windows)版本的 CDC。需启用日志归档和 CDC 功能。

2.2 Debezium Server 与 Embedded 模式

模式类型说明适用场景注意事项
Kafka Connect 模式Debezium 作为 Kafka Connect 的插件运行,事件写入 Kafka。已有 Kafka 基础设施,需与其他流处理系统集成。部署依赖 Kafka 和 ZooKeeper。
Debezium Server 模式独立运行的 Java 应用,可将变更事件发送到 Kafka、Kinesis、Pub/Sub、Redis 等。无需 Kafka,希望直接对接云服务或轻量消息系统。需自行管理 Server 实例的高可用。
Embedded 模式将 Debezium 引擎嵌入到自定义 Java 应用中,直接控制连接器生命周期。需定制 CDC 逻辑,或集成到现有服务中。开发复杂度高,需处理 offset、错误恢复等。

2.3 Kafka Connect 与 Debezium 的关系

概念说明注意事项
Kafka ConnectApache Kafka 提供的可扩展、容错的数据集成框架,用于连接 Kafka 与其他系统。支持 Source(数据源)和 Sink(目标)连接器。
Debezium 作为 Source ConnectorDebezium 实现了 Kafka Connect 的 Source Connector 接口,将数据库变更作为源数据输入 Kafka。必须将 Debezium 插件加载到 Kafka Connect Worker。
分布式模式Kafka Connect 可以以集群方式运行,支持高可用和任务负载均衡。多个 Worker 实例需共享配置和 offset 存储。
REST 接口管理可通过 HTTP 请求创建、查询、暂停、删除 Debezium 连接器。默认端口为 8083。
插件路径配置Kafka Connect 启动时需通过 plugin.path 指定 Debezium JAR 包位置。插件目录需包含所有依赖 JAR。

2.4 数据流模型:Source、Task、Offset 管理

概念说明注意事项
Source Connector代表一个数据库连接实例,负责逻辑配置与生命周期管理。一个 Connector 可拆分为多个 Task 并行执行。
Task(任务)实际执行数据捕获的单元,每个 Task 负责一部分数据分区(如一个数据库表或 binlog 文件段)。Task 数量由 tasks.max 控制,但实际分配由 Connect 框架决定。
Offset记录连接器读取数据库日志的位置(如 binlog 文件名 + 位置、LSN、timestamp 等),用于故障恢复。必须持久化存储,避免重复或丢失数据。
Offset 存储机制默认使用 Kafka 的特殊 topic(connect-offsets)存储 offset,也可配置为文件存储。生产环境推荐使用 Kafka 存储,保证高可用。
Offset 提交流程Task 周期性将当前读取位置提交到 offset 存储,频率由 offset.flush.interval.ms 控制。设置过长可能导致重启时重复,过短影响性能。
Schema History记录数据库表结构(DDL)变更历史,默认存储在 Kafka topic 中。用于保证事件 schema 的一致性,不可关闭。

第三章:部署与运行环境准备

3.1 环境依赖:Java、Kafka、ZooKeeper、Kafka Connect

组件版本要求说明注意事项
JavaJDK 8 或 JDK 11(推荐)Kafka 与 Kafka Connect 基于 Java 开发,需安装 JDK。建议使用 OpenJDK 或 Oracle JDK,避免使用 JRE。
Apache Kafka2.8+(推荐 3.x)消息中间件,用于传输 Debezium 产生的变更事件。需提前下载并解压 Kafka 发行包。
ZooKeeper3.5+(Kafka 自带)Kafka 依赖 ZooKeeper 管理集群元数据(仅 Kafka < 3.0 需要)。Kafka 3.0+ 支持 KRaft 模式,可不依赖 ZooKeeper。
Kafka Connect内置于 Kafka分布式数据集成框架,Debezium 作为其 Source Connector 运行。需以分布式模式启动 connect-distributed.sh
网络与端口9092(Kafka)、2181(ZooKeeper)、8083(Connect REST API)各组件间需网络互通,防火墙开放对应端口。建议在单机测试时使用 localhost,生产环境配置正确主机名。

3.2 安装 Debezium Kafka Connect 插件

步骤操作说明示例命令/配置注意事项
下载 Debezium 插件从 Debezium 官网下载对应版本的 MySQL Connector 插件包。wget https://repo1.maven.org/maven2/io/debezium/debezium-connector-mysql/2.5.0.Final/debezium-connector-mysql-2.5.0.Final-plugin.tar.gz选择与 Kafka Connect 版本兼容的 Debezium 版本。
解压插件包将插件解压到指定目录,作为 Kafka Connect 的插件路径。tar -xzf debezium-connector-mysql-2.5.0.Final-plugin.tar.gz -C /opt/kafka/plugins/debezium插件目录应包含 debezium-api.jarmysql-connector-java.jar 等。
配置 plugin.path在 Kafka Connect 配置文件中指定插件路径。plugin.path=/opt/kafka/plugins/debezium,/opt/kafka/libs多路径用逗号分隔,确保包含 Kafka 自带的连接器。
添加 MySQL JDBC 驱动确保插件目录中包含 mysql-connector-java-x.x.x.jar可从 Maven 仓库下载并放入插件目录。若使用 MariaDB,可使用 mariadb-java-client

3.3 启动 Kafka Connect 分布式模式

配置项说明示例值注意事项
bootstrap.serversKafka 集群地址localhost:9092确保 Kafka 已启动。
group.idConnect 集群标识connect-cluster同一集群所有 Worker 使用相同 ID。
key.converterKey 序列化方式org.apache.kafka.connect.json.JsonConverter建议生产环境使用 Avro + Schema Registry。
value.converterValue 序列化方式org.apache.kafka.connect.json.JsonConverter同上。
key.converter.schemas.enable是否在序列化中包含 schematrue便于解析结构化数据。
value.converter.schemas.enable是否在 value 中包含 schematrue必须开启以支持 Debezium 事件格式。
config.storage.topic存储连接器配置的 topicconnect-configs自动创建,需设置 replication.factor ≥ 3(生产)。
offset.storage.topic存储 offset 的 topicconnect-offsets同上。
status.storage.topic存储连接器状态的 topicconnect-status同上。
rest.portREST API 端口8083用于提交、查询连接器配置。
plugin.path插件目录路径/opt/kafka/plugins/debezium必须包含 Debezium 插件 JAR 包。

启动命令:

bin/connect-distributed.sh config/connect-distributed.properties

注意事项:

  • 确保 config/connect-distributed.properties 文件已正确配置上述参数。
  • 启动后可通过 http://localhost:8083/ 访问 REST API。
  • 日志位于 logs/connect.log,用于排查插件加载问题。

3.4 验证 Debezium 连接器是否加载成功

方法操作说明示例请求/响应注意事项
REST API 查询连接器类型发送 GET 请求到 /connector-plugins,查看已加载的插件。curl -s http://localhost:8083/connector-plugins响应应包含 "type": "source""class": "io.debezium.connector.mysql.MySqlConnector"
检查特定连接器是否存在查看返回列表中是否包含 MySqlConnectorPostgresConnector 等。响应示例见下方。若未出现,检查 plugin.path 和日志。
日志验证查看 connect.log 是否有 Debezium 插件加载日志。INFO Registered loader for plugin 'debezium-mysql' ...若报 ClassNotFoundException,说明 JAR 包缺失。
创建连接器测试尝试创建一个 MySQL 连接器,观察是否报错。使用 4.3 节方法提交配置。成功创建说明插件加载正常。

响应示例:

[
  {
    "class": "io.debezium.connector.mysql.MySqlConnector",
    "type": "source",
    "version": "2.5.0.Final"
  }
]

第四章:MySQL CDC 实践(基础篇)

4.1 MySQL 环境准备与 Binlog 配置要求

配置项说明推荐值注意事项
binlog_formatBinlog 格式必须为 ROW 模式ROWmy.cnf 中设置 binlog-format=ROW
binlog_row_image行日志内容完整性FULL确保 UPDATE 前后镜像都能捕获。
server-id每个 MySQL 实例必须有唯一 ID12若为 0,Debezium 将无法读取 binlog。
log_bin启用 binlogmysql-bin.log必须开启。
expire_logs_daysbinlog 保留天数7避免日志过早清理导致 Debezium 断流。
用户权限Debezium 连接用户所需权限SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT可使用如下 SQL 创建用户。
数据库引擎表必须使用 InnoDBInnoDBMyISAM 不支持行级变更。

创建用户 SQL:

CREATE USER 'debezium'@'%' IDENTIFIED BY 'dbz';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'debezium'@'%';
FLUSH PRIVILEGES;

4.2 创建 MySQL 连接器配置(Connector Configuration)

参数名语法用途示例值注意事项
name字符串连接器唯一名称mysql-connector-inventory不可重复。
connector.class字符串指定连接器实现类io.debezium.connector.mysql.MySqlConnector固定值。
tasks.max整数最大任务数1MySQL 通常设为 1(单点写入)。
database.hostname字符串MySQL 主机名localhost支持 IP 或域名。
database.port整数MySQL 端口3306默认 3306。
database.user字符串连接用户名debezium需具备复制权限。
database.password字符串连接密码dbz明文存储,生产建议使用密钥管理。
database.server.id整数MySQL server-id223344必须唯一,避免与主从冲突。
database.server.name字符串逻辑服务器名,用于 Kafka topic 命名dbserver1topic 将以 dbserver1.database.table 命名。
database.include.list字符串包含的数据库列表inventory多个用逗号分隔。
table.include.list字符串包含的表列表inventory.customers,inventory.orders可选,过滤特定表。
database.history.kafka.bootstrap.servers字符串Kafka 地址,用于存储 schema 历史localhost:9092必须配置。
database.history.kafka.topic字符串schema 历史存储的 topic 名schema-changes.inventory自动创建。
snapshot.mode字符串快照模式initialinitial: 首次全量 + 增量;never: 仅增量。
include.schema.changes布尔值是否包含 DDL 事件true推荐开启。

4.3 启动 MySQL 连接器(REST API 方式)

操作说明示例请求注意事项
提交连接器配置使用 POST 请求创建连接器见下方示例JSON 格式必须正确,字段名与值注意引号。
查看连接器状态使用 GET 请求查询状态curl -s http://localhost:8083/connectors/mysql-connector-inventory/status响应中 "state": "RUNNING" 表示成功。
停止连接器使用 DELETE 请求curl -X DELETE localhost:8083/connectors/mysql-connector-inventory停止后任务中断,重启需从 offset 恢复。
暂停连接器使用 PUT 请求暂停curl -X PUT localhost:8083/connectors/mysql-connector-inventory/pause可随时恢复。
恢复连接器恢复运行curl -X PUT localhost:8083/connectors/mysql-connector-inventory/resume恢复后继续捕获变更。

提交连接器配置示例:

curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/ -d '{
  "name": "mysql-connector-inventory",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "tasks.max": "1",
    "database.hostname": "localhost",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "dbz",
    "database.server.id": "223344",
    "database.server.name": "dbserver1",
    "database.include.list": "inventory",
    "table.include.list": "inventory.customers",
    "database.history.kafka.bootstrap.servers": "localhost:9092",
    "database.history.kafka.topic": "schema-changes.inventory",
    "snapshot.mode": "initial"
  }
}'

4.4 查看变更事件数据格式(Change Data Event Structure)

字段名说明示例值注意事项
before更新或删除前的行数据{ "id": 1001, "name": "Alice" }INSERT 事件中为 null。
after插入或更新后的行数据{ "id": 1001, "name": "Bob" }DELETE 事件中为 null。
source源数据库元信息见下方说明包含数据库名、表名、binlog 位置等。
op操作类型c(create/INSERT)、r(read/snapshot)、u(update)、d(delete)用于判断事件类型。
ts_ms事件在 Debezium 中生成的时间戳(毫秒)1634567890123不是数据库时间。
transaction事务相关信息(可选){ "id": "123", "total_order": 5, "data_collection_order": 1 }需开启事务元数据配置。

source 字段示例:

{
  "version": "2.5.0.Final",
  "connector": "mysql",
  "name": "dbserver1",
  "ts_ms": 1634567890000,
  "db": "inventory",
  "table": "customers",
  "server_id": 223344,
  "event": "INSERT",
  "row": 0
}

事件结构完整示例:

{
  "before": null,
  "after": {
    "id": 1001,
    "name": "Alice"
  },
  "source": {
    "version": "2.5.0.Final",
    "connector": "mysql",
    "name": "dbserver1",
    "ts_ms": 1634567890000,
    "db": "inventory",
    "table": "customers",
    "server_id": 223344,
    "event": "INSERT"
  },
  "op": "c",
  "ts_ms": 1634567890123
}

4.5 监听 INSERT、UPDATE、DELETE 事件

事件类型op 值beforeafter识别方式处理建议
INSERTcnull包含数据op == 'c'将 after 数据插入目标系统。
UPDATEu包含旧值包含新值op == 'u'使用 after 更新目标,可用 before 做条件更新。
DELETEd包含旧值nullop == 'd'根据 before 中的主键删除目标数据。
Snapshot(初始快照)rnull包含数据op == 'r'通常作为全量同步,需注意去重。

注意事项:

  • 所有事件均通过 Kafka topic 发送,topic 名为 {database.server.name}.{database}.{table},如 dbserver1.inventory.customers
  • 消费者需解析 JSON 或 Avro 格式事件,提取 op 字段判断操作类型。
  • 建议消费者实现幂等性,避免重复处理导致数据错乱。

第五章:PostgreSQL CDC 实践

5.1 PostgreSQL 的逻辑复制(Logical Replication)配置

配置项说明推荐值/操作注意事项
wal_levelWAL 日志级别,必须为 logical 才支持逻辑复制logicalpostgresql.conf 中设置 wal_level = logical,需重启数据库生效。
max_wal_senders最大 WAL 发送进程数5 或更高每个复制连接占用一个 sender,建议 ≥3。
max_replication_slots最大复制槽数量5 或更高Debezium 使用 replication slot 记录位置,不可少于连接器数。
创建复制用户为 Debezium 创建专用用户,并授予复制权限见下方 SQL用户必须具有 REPLICATION 属性。
shared_preload_libraries必须加载 pgoutput 插件确保 shared_preload_libraries 包含 pgoutputpgoutput 是 PostgreSQL 内置的逻辑解码插件,无需额外安装。
监听地址允许远程连接listen_addresses = '*'生产环境应限制 IP 范围。
pg_hba.conf 配置允许复制用户通过流复制连接见下方配置位于 pg_hba.conf 文件末尾,重启或重载配置生效。

创建复制用户:

CREATE USER debezium WITH REPLICATION LOGIN PASSWORD 'dbz';

pg_hba.conf 配置:

host replication debezium 0.0.0.0/0 md5

验证命令:

-- 查看当前 WAL 设置
SHOW wal_level;
-- 查看可用的复制槽(初始为空)
SELECT * FROM pg_replication_slots;

5.2 创建 PostgreSQL 连接器配置

参数名语法用途示例值注意事项
name字符串连接器唯一名称postgres-connector-inventory不可重复。
connector.class字符串指定连接器类io.debezium.connector.postgresql.PostgresConnector固定值。
tasks.max整数最大任务数1PostgreSQL 通常设为 1。
database.hostname字符串PostgreSQL 主机名localhost支持 IP 或域名。
database.port整数PostgreSQL 端口5432默认端口。
database.user字符串连接用户名debezium需具备复制权限。
database.password字符串连接密码dbz明文存储,生产建议使用密钥管理工具。
database.dbname字符串要连接的数据库名inventory必填。
database.server.name字符串逻辑服务器名,用于 Kafka topic 命名pgserver1topic 格式:pgserver1.public.customers
plugin.name字符串逻辑解码插件名pgoutputPostgreSQL 官方插件,Debezium 默认使用。
slot.name字符串复制槽名称debezium_slot必须唯一,不能包含 -,推荐使用字母数字下划线。
snapshot.mode字符串快照模式initialinitial: 初始快照 + 增量;never: 仅增量。
publication.name字符串发布名称dbz_publication若未指定,Debezium 自动创建名为 dbz_publication 的 publication。
database.include.list字符串包含的数据库列表inventory多租户场景可选。
schema.include.list字符串包含的 schema 列表public,sales多个用逗号分隔。
table.include.list字符串包含的表列表public.customers,public.orders可选,用于过滤特定表。
include.schema.changes布尔值是否输出 DDL 事件true推荐开启。
lsn.commit.interval.ms整数LSN 提交间隔(毫秒)10000控制 offset 更新频率,影响恢复点。

5.3 启动与管理 PostgreSQL 连接器

操作说明示例请求注意事项
创建连接器使用 REST API 提交配置见下方示例JSON 必须格式正确,避免缺少引号或逗号。
查询连接器状态获取当前运行状态curl -s http://localhost:8083/connectors/postgres-connector-inventory/status响应中 "state": "RUNNING" 表示正常。
暂停连接器暂停捕获,不提交 offsetcurl -X PUT localhost:8083/connectors/postgres-connector-inventory/pause可减少资源占用,但 WAL 可能堆积。
恢复连接器继续捕获curl -X PUT localhost:8083/connectors/postgres-connector-inventory/resume从上次 offset 继续读取。
删除连接器彻底移除连接器curl -X DELETE localhost:8083/connectors/postgres-connector-inventory注意:默认不会自动删除 replication slot,需手动清理(见 5.4 节)。

创建连接器示例:

curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/ -d '{
  "name": "postgres-connector-inventory",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "tasks.max": "1",
    "database.hostname": "localhost",
    "database.port": "5432",
    "database.user": "debezium",
    "database.password": "dbz",
    "database.dbname": "inventory",
    "database.server.name": "pgserver1",
    "plugin.name": "pgoutput",
    "slot.name": "debezium_slot",
    "publication.name": "dbz_publication",
    "database.include.list": "inventory",
    "schema.include.list": "public",
    "table.include.list": "public.customers",
    "snapshot.mode": "initial"
  }
}'

5.4 处理 LSN 与 Slot 管理

概念说明操作/查询命令注意事项
LSN (Log Sequence Number)WAL 日志的唯一位置标识,Debezium 使用它记录读取进度。由 PostgreSQL 自动生成,如 16/B37E4A8不可人为修改。
Replication Slot一种机制,用于保留 WAL 日志直到消费者处理完毕,防止日志被过早清理。Debezium 自动创建(若不存在),名称由 slot.name 指定。若连接器被删除而 slot 未清理,WAL 将持续堆积,导致磁盘爆满!
查看复制槽查询当前存在的 slot见下方 SQL确认 active = true 表示正在使用。
手动创建复制槽可预先创建见下方 SQL非必需,Debezium 可自动创建。
删除复制槽连接器删除后必须手动清理见下方 SQL关键操作:避免 WAL 日志无限增长。
Slot 清理策略建议在删除连接器后立即执行流程:1. DELETE REST API 删除连接器 2. SQL 执行 pg_drop_replication_slot可结合监控告警,定期检查 inactive slots。
LSN 与 Offset 关系Debezium 将 LSN 存储在 Kafka 的 offset 中,用于故障恢复存储格式:"lsn": 234567890重启后从该 LSN 继续读取。

查看复制槽:

SELECT slot_name, plugin, slot_type, active, restart_lsn FROM pg_replication_slots;

手动创建复制槽:

SELECT pg_create_logical_replication_slot('debezium_slot', 'pgoutput');

删除复制槽:

SELECT pg_drop_replication_slot('debezium_slot');

第六章:MongoDB CDC 实践

6.1 MongoDB 副本集(Replica Set)配置要求

配置项说明要求注意事项
部署模式必须为副本集(Replica Set)或分片集群(Sharded Cluster)单节点不支持变更流变更流(Change Streams)依赖 oplog,仅副本集及以上支持。
副本集初始化至少一个主节点和一个从节点使用 rs.initiate() 初始化即使单机多实例也需配置为副本集。
Oplog 大小操作日志文件大小默认 1GB,建议根据写入量调整若 oplog 太小,可能因延迟导致”oplog gap”错误。
用户权限Debezium 用户需具备 clusterMonitorread 权限见下方创建用户命令clusterMonitor 用于访问集群状态,read 用于读取数据。
认证机制支持 SCRAM-SHA-1 或 SCRAM-SHA-256配置 authenticationMechanism默认自动检测。
SSL/TLS可选加密连接配置 ssl.enabled=true生产环境建议启用。
Storage Engine存储引擎必须为 WiredTigerMMAPv1 不支持变更流。

创建用户:

use admin
db.createUser({
  user: "debezium",
  pwd: "dbz",
  roles: ["clusterMonitor", "read"]
})

验证副本集状态:

// 进入 mongo shell
rs.status()
// 应看到成员状态为 PRIMARY / SECONDARY

6.2 配置 MongoDB 连接器

参数名语法用途示例值注意事项
name字符串连接器名称mongodb-connector-sales唯一。
connector.class字符串连接器类io.debezium.connector.mongodb.MongoDbConnector固定值。
tasks.max整数任务数1MongoDB Connector 通常为 1。
mongodb.hosts字符串MongoDB 主机地址rs0/localhost:27017格式:<replica_set_name>/host:port
mongodb.user字符串用户名debezium具备 clusterMonitorread 权限。
mongodb.password字符串密码dbz明文。
mongodb.authsource字符串认证数据库admin通常是 admin
database.include.list字符串包含的数据库sales,users多个用逗号分隔。
collection.include.list字符串包含的集合sales.orders,sales.customers可选,用于过滤。
database.history.kafka.bootstrap.servers字符串Kafka 地址localhost:9092必填。
database.history.kafka.topic字符串schema 历史 topicschema-changes.sales自动创建。
snapshot.mode字符串快照模式initialinitial: 先全量再增量;never: 仅增量。
heartbeat.interval.ms整数心跳间隔(毫秒)10000保持连接活跃,避免超时。
poll.interval.ms整数轮询变更流间隔1000控制延迟与负载。
ssl.enabled布尔值是否启用 SSLfalsetrue生产建议启用。

6.3 变更流(Change Streams)事件解析

字段名说明示例值注意事项
operationType操作类型insert, update, delete, replace, drop, rename对应 MongoDB 操作。
fullDocumentINSERT/REPLACE 时的完整文档{ "_id": "abc", "name": "Alice" }UPDATE/DELETE 为 null。
ns.db数据库名sales源数据库。
ns.coll集合名customers源集合。
documentKey文档主键(_id){ "_id": "abc" }所有事件都包含。
updateDescriptionUPDATE 操作的修改详情{ "updatedFields": { "name": "Bob" }, "removedFields": [] }仅 UPDATE 事件存在。
clusterTime逻辑时间戳Timestamp(1634567890, 1)用于排序和一致性。
txnNumber事务编号(如有)123支持事务性写入。
lsid会话 ID{ "id": { "$binary": "..." } }用于跟踪会话。

事件结构示例(INSERT):

{
  "operationType": "insert",
  "fullDocument": {
    "_id": "1001",
    "name": "Alice"
  },
  "ns": {
    "db": "sales",
    "coll": "customers"
  },
  "documentKey": {
    "_id": "1001"
  }
}

6.4 处理集合过滤与文档前置映像

功能配置参数说明示例值注意事项
集合过滤collection.include.list指定要监听的集合sales.orders,sales.customers不配置则监听所有集合。
集合排除无直接 exclude 参数需通过 include 显式列出建议明确包含所需集合。
文档前置映像(Pre-image)capture.mode控制是否捕获更新前的文档change_streams_update_full:捕获 before;change_streams:不捕获需 MongoDB 6.0+ 支持。
获取 before 数据结合 SMT 或自定义逻辑在 UPDATE 事件中提取旧值需启用 capture.mode=change_streams_update_full用于审计、对比等场景。
输出格式控制value.converter可转换为 Avro、JSON Schema 等推荐使用 Avro + Schema Registry更高效且结构清晰。

前置映像配置示例:

{
  "name": "mongodb-connector-with-preimage",
  "config": {
    "connector.class": "io.debezium.connector.mongodb.MongoDbConnector",
    "capture.mode": "change_streams_update_full"
  }
}

第七章:连接器通用配置详解

7.1 基础配置参数(name, connector.class, tasks.max 等)

参数名语法用途示例值注意事项
name字符串连接器唯一标识名称mysql-connector-inventory必须全局唯一,不可重复。
connector.class字符串指定连接器实现类io.debezium.connector.mysql.MySqlConnector必须与数据库类型匹配。
io.debezium.connector.postgresql.PostgresConnector
io.debezium.connector.mongodb.MongoDbConnector
tasks.max整数最大并行任务数1大多数 CDC 连接器设为 1(单点写入),部分支持并行(如 PostgreSQL 分表)。
connector.client.id字符串客户端标识debezium-connector-1用于日志和监控,可选。
errors.log.enable布尔值是否记录错误日志true推荐开启以便排查问题。
errors.log.include.messages布尔值是否在错误日志中包含消息内容true有助于调试,但可能泄露敏感数据。
transforms字符串指定 SMT 转换链名称(多个用逗号分隔)Reroute, RenameField用于后续单消息转换配置。
header.converter字符串Header 序列化方式org.apache.kafka.connect.storage.SimpleHeaderConverter默认值,通常无需修改。

7.2 数据库连接配置(hostname, port, user, password 等)

参数名语法用途示例值注意事项
database.hostname字符串数据库主机名或 IPlocalhost192.168.1.100支持 DNS 解析。
database.port整数数据库端口3306(MySQL)、5432(PostgreSQL)、27017(MongoDB)必须开放且可访问。
database.user字符串登录用户名debezium需具备对应数据库的复制/读取权限。
database.password字符串登录密码dbz明文存储,生产环境建议使用密钥管理工具(如 HashiCorp Vault)。
database.ssl.mode字符串SSL 连接模式disabled, required, verified生产环境建议设为 requiredverified
database.ssl.keystore.location字符串SSL 密钥库路径/path/to/keystore.jks启用 SSL 时需配置。
database.ssl.keystore.password字符串密钥库密码changeit
database.ssl.truststore.location字符串信任库路径/path/to/truststore.jks
database.ssl.truststore.password字符串信任库密码changeit

7.3 数据库对象过滤(include.list, exclude.list 等)

参数名语法用途示例值注意事项
database.include.list字符串包含的数据库列表inventory,sales多个用逗号分隔,不配置则监听所有数据库。
database.exclude.list字符串排除的数据库列表sys,mysql,information_schema常用于跳过系统库。
schema.include.list字符串(PostgreSQL/MongoDB)包含的 schemapublic,app_dataPostgreSQL 适用。
schema.exclude.list字符串排除的 schemapg_catalog,information_schema
table.include.list字符串包含的表/集合inventory.customers,inventory.orders格式:db.tableschema.table
table.exclude.list字符串排除的表/集合%.audit_log,%.temp_%支持通配符 %_
collection.include.list字符串(MongoDB)包含的集合sales.orders,sales.customersMongoDB 专用。
collection.exclude.list字符串(MongoDB)排除的集合sales.logs
topic.regex.replacement字符串自定义 topic 命名规则$1-$2-$3配合 SMT 使用,重写 topic 名称。

注意include.listexclude.list 不能同时用于同一层级(如不能同时设置 table.include.listtable.exclude.list),优先级为 include > exclude。

7.4 事件格式与序列化配置(key.converter, value.converter)

参数名语法用途示例值注意事项
key.converter字符串Key 的序列化器org.apache.kafka.connect.json.JsonConverter推荐生产使用 Avro + Schema Registry。
io.confluent.connect.avro.AvroConverter
value.converter字符串Value 的序列化器同上必须与 key.converter 一致或兼容。
key.converter.schemas.enable布尔值是否在 key 中包含 schematrue建议开启,便于解析结构化数据。
value.converter.schemas.enable布尔值是否在 value 中包含 schematrue必须开启以支持 Debezium 事件结构。
key.converter.schema.registry.url字符串Schema Registry 地址(Avro)http://localhost:8081使用 Avro 时必须配置。
value.converter.schema.registry.url字符串同上http://localhost:8081
internal.key.converter字符串内部元数据 key 序列化器org.apache.kafka.connect.json.JsonConverter用于 offset、config 存储,建议与 key.converter 一致。
internal.value.converter字符串内部元数据 value 序列化器org.apache.kafka.connect.json.JsonConverter同上。

7.5 Offset 与恢复机制配置(offset.storage, offset.flush.interval.ms)

参数名语法用途示例值注意事项
offset.storage字符串offset 存储位置org.apache.kafka.connect.storage.KafkaOffsetBackingStore生产环境必须使用 Kafka 存储以保证高可用。
org.apache.kafka.connect.storage.FileOffsetBackingStore
offset.storage.topic字符串存储 offset 的 Kafka topicconnect-offsets自动创建,需设置 replication.factor ≥ 3
offset.storage.replication.factor整数offset topic 的副本数3生产环境建议 ≥3。
offset.storage.partitions整数offset topic 分区数25通常无需修改。
offset.flush.interval.ms整数offset 提交间隔(毫秒)10000(10秒)设置过短影响性能,过长可能导致重启时重复。
offset.flush.timeout.ms整数offset 提交超时时间30000必须大于 offset.flush.interval.ms
offset.backing.store.file.filename字符串(File 模式)offset 文件路径/tmp/connect.offsets仅测试使用,不支持分布式。

恢复机制说明:

  • Debezium 在启动时从 offset 存储中读取上次位置(如 binlog pos、LSN、clusterTime),继续增量捕获。
  • 若 offset 丢失或无效,将根据 snapshot.mode 决定是否重新快照。

第八章:高级特性与数据处理

8.1 SMT(单消息转换)介绍与常用转换器

转换器名称类名用途注意事项
ExtractFieldorg.apache.kafka.connect.transforms.ExtractField从结构体中提取字段作为 value 或 key支持嵌套字段($.field.subfield)。
ValueToKeyorg.apache.kafka.connect.transforms.ValueToKey将 value 中的某个字段提取为消息 key常用于按主键分区。
ReplaceFieldorg.apache.kafka.connect.transforms.ReplaceField重命名、排除或类型转换字段支持嵌套字段操作。
MaskFieldorg.apache.kafka.connect.transforms.MaskField对字段值进行掩码(如手机号、身份证)用于脱敏场景。
Filterorg.apache.kafka.connect.transforms.Filter过滤掉满足条件的消息可结合 predicate 实现复杂逻辑。
FilterTopicorg.apache.kafka.connect.transforms.FilterTopic根据 topic 名称过滤消息非 record 级别过滤。
InsertFieldorg.apache.kafka.connect.transforms.InsertField添加静态或动态字段(如 timestamp、header)支持 RecordHeadertimestamp 等。
HoistFieldorg.apache.kafka.connect.transforms.HoistField将整个 value 包装为一个字段如将 {name: "Alice"} 变为 {"value": {name: "Alice"}}
Flattenorg.apache.kafka.connect.transforms.Flatten展平嵌套结构(用 . 连接字段名){"addr.city"} 替代 {"addr": {"city": ...}}
TimestampRouterorg.apache.kafka.connect.transforms.TimestampRouter根据时间戳重写 topic 名称如按天分表:topic-${YYYYMMdd}
RegexRouterorg.apache.kafka.connect.transforms.RegexRouter使用正则表达式重写 topic 名称强大灵活,可用于路由。

SMT 链说明:多个 SMT 可串联执行,顺序由配置名决定,如 transforms=A,B,C

8.2 使用 SMT 进行字段过滤、重命名、类型转换

操作配置示例说明注意事项
字段过滤见下方示例排除 credit_cardssn 字段使用 blacklist 排除,whitelist 保留指定字段。
字段重命名见下方示例col1 改为 field1多个用逗号分隔。
类型转换不直接支持类型转换需结合 Schema 或下游处理SMT 不改变字段物理类型,仅逻辑重命名。
提取字段为 Key见下方示例使用 id 字段作为 Kafka 消息 keyschema=false 表示不保留 schema。
展平嵌套结构见下方示例{"name": {"first": "A", "last": "B"}} 变为 {"name_first": "A", "name_last": "B"}减少嵌套层级,便于消费。

字段过滤示例:

transforms=FilterFields
transforms.FilterFields.type=org.apache.kafka.connect.transforms.ReplaceField$Value
transforms.FilterFields.blacklist=credit_card,ssn

字段重命名示例:

transforms=RenameField
transforms.RenameField.type=org.apache.kafka.connect.transforms.ReplaceField$Value
transforms.RenameField.renames=col1:field1,col2:field2

提取字段为 Key 示例:

transforms=KeyFromId
transforms.KeyFromId.type=org.apache.kafka.connect.transforms.ValueToKey
transforms.KeyFromId.fields=id
transforms.KeyFromId.schema=false

展平嵌套结构示例:

transforms=FlattenIt
transforms.FlattenIt.type=org.apache.kafka.connect.transforms.Flatten$Value
transforms.FlattenIt.delimiter=_

8.3 忽略特定操作(如只捕获 INSERT)

操作配置示例说明注意事项
仅捕获 INSERT见下方示例过滤掉非 op=c(INSERT)的事件需定义 predicate 判断 op 字段。
仅捕获 UPDATEpredicates.IsUpdate.value=u将 value 改为 u
仅捕获 DELETEpredicates.IsDelete.value=d将 value 改为 d
忽略 DELETE见下方示例保留非 DELETE 事件negate=true 表示取反。

仅捕获 INSERT 示例:

transforms=FilterUpdatesDeletes
transforms.FilterUpdatesDeletes.type=org.apache.kafka.connect.transforms.Filter
transforms.FilterUpdatesDeletes.predicate=IsInsert

predicates=IsInsert
predicates.IsInsert.type=org.apache.kafka.connect.transforms.predicates.TopicField$Value
predicates.IsInsert.field=op
predicates.IsInsert.value=c

忽略 DELETE 示例:

transforms=FilterDeletes
transforms.FilterDeletes.type=org.apache.kafka.connect.transforms.Filter
transforms.FilterDeletes.if=op == 'd'
transforms.FilterDeletes.negate=true

说明:Debezium 事件中 op 字段表示操作类型:c=INSERT, r=READ(快照), u=UPDATE, d=DELETE。

8.4 快照(Snapshot)机制与配置(snapshot.mode)

snapshot.mode 值说明适用场景注意事项
initial若无 offset,先全量快照再增量;若有 offset,直接增量首次部署或从头开始默认值,最常用。
when_needed类似 initial,但遇到不可恢复错误时自动重新快照容错性要求高可能意外触发全量同步。
never仅增量捕获,不执行快照已有历史数据,只关心新变更若无 offset 或 binlog 已清理,将失败。
schema_only仅读取 schema,不读取数据仅需同步表结构不生成数据事件。
schema_only_recovery用于恢复损坏的 schema history修复场景需配合 database.history.skip.unparseable.ddl
custom使用自定义快照类高级定制需实现 Snapshotter 接口。

相关参数:

参数说明示例值注意事项
snapshot.locking.mode快照期间锁表方式minimal(推荐,最小化锁)、none(无锁,MySQL 8.0+)、extended(长时间锁)minimal 对业务影响小。
snapshot.include.collection.list快照包含的集合(MongoDB)sales.orders,sales.customers
snapshot.delay.ms快照开始前延迟10000避免启动风暴。
snapshot.fetch.size每次读取行数1024控制内存使用。

8.5 处理大事务与延迟问题

问题解决方案配置参数与值说明
大事务导致延迟增加缓冲与超时poll.interval.ms=500poll.timeout.ms=30000max.queue.size=8192max.batch.size=2048提高拉取频率和队列容量。
网络延迟高调整心跳与超时heartbeat.interval.ms=10000heartbeat.action.query=SELECT 1PostgreSQL/MongoDB 适用。
消费者积压(Lag)增加消费者或优化处理逻辑Debezium 本身是 Source,需下游消费者优化。
快照期间负载高分批快照或低峰期执行snapshot.mode=initialsnapshot.delay.ms=60000避免影响线上业务。
Binlog/WAL 清理过快延长日志保留时间MySQL: expire_logs_days=7,PostgreSQL: wal_keep_size=1GB防止 Debezium 断流。
连接中断自动恢复启用错误重试errors.tolerance=allerrors.deadletterqueue.topic.name=dlq-topic容错处理,避免连接器失败。

监控建议:

  • 使用 Kafka 监控工具(如 Kafka Manager、Confluent Control Center)查看 consumer lag。
  • 监控 Debezium 连接器的 status 和 metrics(通过 JMX 或 REST API)。
  • 设置告警:offset lag > 1000、连接器状态非 RUNNING。

第九章:事件格式与数据解析

9.1 Debezium 事件结构(Before、After、Source、Op、Ts_ms)

Debezium 生成的每条变更事件均为结构化 JSON 或 Avro 格式,包含以下核心字段:

字段名类型说明示例值注意事项
before结构体(Object)更新或删除前的旧数据{"id": 1001, "name": "Alice"}INSERT 时为 null。
after结构体(Object)插入或更新后的新数据{"id": 1001, "name": "Bob"}DELETE 时为 null。
source结构体(Object)源数据库元信息见下方示例包含 LSN、binlog pos、ts_sec 等。
op字符串(String)操作类型c(create/INSERT)、r(read/snapshot)、u(update)、d(delete)用于判断变更类型。
ts_ms长整型(Long)事件在 Debezium 中创建的时间戳(毫秒)1727683200000注意:不是数据库操作时间。
transaction结构体(可选)事务上下文信息{"id": "tx-123", "total_order": 5, "data_collection_order": 2}大事务中用于关联多条事件。

source 字段详解(以 PostgreSQL 为例):

{
  "version": "2.3.0.Final",
  "connector": "postgresql",
  "name": "pgserver1",
  "ts_ms": 1727683200000,
  "snapshot": false,
  "db": "inventory",
  "schema": "public",
  "table": "customers",
  "txId": 12345,
  "lsn": 234567890,
  "xmin": 123456
}

不同操作的事件结构对比:

操作beforeafterop
INSERTnull新数据c
UPDATE旧数据新数据u
DELETE旧数据nulld
Snapshot (READ)null当前数据r

9.2 如何解析 JSON 格式的变更事件

步骤操作说明示例代码(Python)注意事项
1. 接收 Kafka 消息从 Debezium 输出的 topic 消费消息见下方代码确保 value_deserializer 正确解析 JSON。
2. 判断操作类型读取 op 字段见下方代码r 表示快照数据,通常可忽略或标记为初始状态。
3. 提取源信息解析 source 字段见下方代码source['lsn']source['file'] + source['pos'] 可用于恢复点定位。
4. 处理嵌套结构JSON 支持嵌套对象和数组直接访问 event['after']['address']['city']注意空值和类型检查。
5. 序列化存储写入目标系统(如 Elasticsearch、数据湖)使用 json.dumps() 或 ORM 工具建议保留原始事件结构以便审计。

Python 解析示例:

from kafka import KafkaConsumer
import json

consumer = KafkaConsumer(
    'pgserver1.public.customers',
    bootstrap_servers='localhost:9092',
    value_deserializer=lambda m: json.loads(m.decode('utf-8'))
)

for msg in consumer:
    event = msg.value
    op = event['op']
    if op == 'c':
        print("INSERT:", event['after'])
    elif op == 'u':
        print("UPDATE from", event['before'], "to", event['after'])
    elif op == 'd':
        print("DELETE:", event['before'])

    source = event['source']
    table = f"{source['schema']}.{source['table']}"
    db_time = source['ts_ms']

建议:生产环境建议使用 Avro + Schema Registry 替代纯 JSON,以提升性能、类型安全与兼容性。

9.3 使用 Avro 与 Schema Registry 的集成

组件作用配置参数示例值
Schema Registry存储和管理 Avro schema独立服务,如 Confluent Schema Registryhttp://localhost:8081
Avro Converter将数据序列化为 Avro 并注册 schema配置在 Kafka Connect 中io.confluent.connect.avro.AvroConverter
key.converterKey 序列化器key.converterio.confluent.connect.avro.AvroConverter
value.converterValue 序列化器value.converterio.confluent.connect.avro.AvroConverter
schema.registry.urlSchema Registry 地址key.converter.schema.registry.urlvalue.converter.schema.registry.urlhttp://localhost:8081

优势:

  • 强类型:字段类型明确(string, int, long, decimal 等)。
  • 高效:二进制格式,体积小,序列化快。
  • 兼容性:支持 schema evolution(backward/forward compatibility)。
  • 集中管理:所有 schema 可视化查看与版本控制。

验证 schema 是否注册:

curl http://localhost:8081/subjects
# 输出所有 subject(通常为 topic 名)

curl http://localhost:8081/subjects/pgserver1.public.customers-value/versions/latest
# 查看最新 schema

9.4 处理复杂类型(DECIMAL、TIMESTAMP、JSON 等)

Debezium 对复杂类型采用**逻辑类型(Logical Types)**映射,确保精度与语义。

数据库类型Debezium 逻辑类型Avro Schema 示例说明注意事项
DECIMAL(p,s)org.apache.kafka.connect.data.Decimal见下方以二进制存储高精度数值消费者需使用 Kafka Connect 的 Decimal 类解析。
TIMESTAMP / DATETIMEorg.apache.kafka.connect.data.Timestamp见下方转换为毫秒级时间戳值为 1970-01-01 起的毫秒数。
DATEorg.apache.kafka.connect.data.Date见下方转换为天数(自 1970-01-01)
TIMEorg.apache.kafka.connect.data.Time见下方毫秒数(当日零点起)
JSON (MySQL/PostgreSQL)org.apache.kafka.connect.data.Json见下方存储为字符串,但标记为 JSON 类型实际值为 JSON 字符串。
UUIDorg.apache.kafka.connect.data.Decimal 或 string通常映射为 stringPostgreSQL 中为 uuid 类型可通过 SMT 转换。
BYTEA / BLOBbytes{ "type": "bytes" }二进制数据建议避免同步大对象。

DECIMAL 类型 Avro Schema 示例:

{
  "type": "bytes",
  "logicalType": "decimal",
  "precision": 10,
  "scale": 2
}

TIMESTAMP 类型 Avro Schema 示例:

{
  "type": "long",
  "logicalType": "timestamp-millis"
}

DATE 类型 Avro Schema 示例:

{
  "type": "int",
  "logicalType": "date"
}

TIME 类型 Avro Schema 示例:

{
  "type": "int",
  "logicalType": "time-millis"
}

JSON 类型 Avro Schema 示例:

{
  "type": "string",
  "logicalType": "json"
}

消费端处理建议(Java 示例):

// 使用 Kafka Connect 数据 API 解析逻辑类型
Schema schema = field.schema();
if (schema.name().equals("org.apache.kafka.connect.data.Decimal")) {
    BigDecimal value = Decimal.fromConnectData(schema, value);
}

第十章:监控、运维与故障排查

10.1 查看连接器状态与任务信息(REST API)

API 端点HTTP 方法说明示例请求响应关键字段
/connectorsGET列出所有连接器curl http://localhost:8083/connectors["mysql-connector", "pg-connector"]
/connectors/{name}GET获取连接器配置curl http://localhost:8083/connectors/mysql-connector返回完整配置 JSON
/connectors/{name}/statusGET获取连接器运行状态curl http://localhost:8083/connectors/mysql-connector/status"state": "RUNNING", "tasks": [{"state": "RUNNING"}]
/connectors/{name}/tasksGET获取任务列表curl http://localhost:8083/connectors/mysql-connector/tasks任务 ID、状态、worker_id
/connectors/{name}/tasks/{taskid}/statusGET获取单个任务状态curl http://localhost:8083/connectors/mysql-connector/tasks/0/status同上,更细粒度
/connectors/{name}/configGET获取配置curl http://localhost:8083/connectors/mysql-connector/config仅返回 config 部分
/connectors/{name}/configPUT更新配置(热更新)curl -X PUT -H "Content-Type:application/json" ...部分参数支持动态更新

响应示例(status):

{
  "name": "mysql-connector",
  "connector": {
    "state": "RUNNING",
    "worker_id": "192.168.1.10:8083"
  },
  "tasks": [
    {
      "id": 0,
      "state": "RUNNING",
      "worker_id": "192.168.1.10:8083"
    }
  ],
  "type": "source"
}

10.2 监控 Lag、Offset 提交情况

指标说明监控方式工具建议
Consumer Lag消费者落后于最新消息的数量kafka-consumer-groups.sh --bootstrap-server ... --describe --group connect-clusterConfluent Control Center、Prometheus + Grafana
Offset 提交频率offset.flush.interval.ms 控制查看 Kafka topic connect-offsets 的消息速率Kafka Manager
Debezium 内部 MetricsJMX 暴露的指标启用 JMX,使用 jconsole 或 Prometheus JMX Exporter关键 MBean:kafka.connect:type=connector-task,connector=...,task=...

关键指标列表:

指标说明
source-record-poll-rate每秒从数据库拉取的记录数,反映捕获速度
message-rate每秒写入 Kafka 的消息数,反映输出吞吐
batch-size-avg每次 poll 返回的记录数,判断批量效率
connection-age数据库连接存活时间,判断是否频繁重连
event-emit-delay-ms事件从数据库到 Kafka 的延迟,核心延迟指标

计算端到端延迟:

  • 数据库操作时间:event['source']['ts_ms']
  • Kafka 接收时间:msg.timestamp
  • 延迟 = msg.timestamp - event['source']['ts_ms']

10.3 常见错误与日志分析

错误现象可能原因日志关键字解决方案
binlog not found / WAL segment removedbinlog/WAL 被清理,连接器断流Could not find first log file name from binlog indexrequested LSN X is not in timeline增加 expire_logs_days(MySQL)或 wal_keep_size(PostgreSQL);重新全量快照。
snapshot lock timeout快照期间锁表超时,影响业务Lock wait timeout exceeded使用 snapshot.locking.mode=minimal;在低峰期执行快照;优化长事务。
Authentication failed用户名密码错误或权限不足Access denied for user检查 database.user / password;确认用户具备 REPLICATION CLIENT, REPLICATION SLAVE(MySQL)或 REPLICATION(PostgreSQL)。
Replication slot not foundPostgreSQL slot 被手动删除replication slot "debezium_slot" does not exist重建 slot 或修改 slot.name;避免手动删除。
Oplog entry not foundMongoDB oplog 太小,记录被覆盖CappedPositionLost增大 oplog 大小(mongod --oplogSize 4096);重启连接器触发快照。
Out of memoryJVM 内存不足java.lang.OutOfMemoryError增加 KAFKA_HEAP_OPTS="-Xms2g -Xmx2g";减少 max.batch.size
Connection refused数据库不可达Connection refused, timeout检查网络、防火墙、数据库是否运行。

开启详细日志:

# log4j.properties
log4j.logger.io.debezium=DEBUG
log4j.logger.org.apache.kafka.connect=DEBUG

10.4 连接器暂停、重启与删除操作

操作REST API 命令说明注意事项
暂停curl -X PUT http://localhost:8083/connectors/{name}/pause停止读取数据库,不提交 offset资源占用低,可快速恢复。
恢复curl -X PUT http://localhost:8083/connectors/{name}/resume继续从上次位置读取适用于临时维护。
重启curl -X POST http://localhost:8083/connectors/{name}/restart重启连接器或任务用于恢复卡住的任务;可指定任务 ID:/tasks/0/restart
删除curl -X DELETE http://localhost:8083/connectors/{name}彻底移除连接器配置注意:不会自动删除 replication slot(PostgreSQL)或 publication,需手动清理!

安全删除流程:

  1. DELETE /connectors/{name}
  2. 检查数据库端资源:
    • PostgreSQL:SELECT * FROM pg_replication_slots;pg_drop_replication_slot('slot_name');
    • MySQL:无额外资源。
    • MongoDB:无额外资源。
  3. (可选)清理 Kafka topic 和 schema registry subject。

第十一章:Debezium Server 模式(无 Kafka)

11.1 Debezium Server 简介与适用场景

项目内容说明
定义Debezium Server 是一个独立运行的轻量级进程,可直接将数据库变更事件发送到消息中间件或云服务,无需 Kafka 集群。基于 Quarkus 构建,支持原生编译。
架构特点单进程部署、内置 CDC 引擎(与 Debezium Connector 相同)、支持多种输出目标(Sink)轻量、快速启动、资源占用低。
适用场景小型或中等规模系统、已使用云消息服务(如 GCP Pub/Sub、AWS SNS/SQS)、不希望维护 Kafka 集群、快速 PoC 或边缘部署适合云原生、Serverless 架构。
不适用场景高吞吐/大规模数据流、需要复杂流处理(如窗口、聚合)、多消费者/消息回溯需求强Kafka 在这些场景更具优势。
支持的连接器MySQL、PostgreSQL、SQL Server、Oracle、MongoDB 等主流数据库与 Kafka Connect 模式功能一致。

11.2 配置 Debezium Server 输出到 Google Cloud Pub/Sub / Amazon SNS / Redis 等

输出目标配置前缀关键参数示例值注意事项
Google Cloud Pub/Subquarkus.google.cloud.pubsub.project-idcredentials-pathtopicproject-id=my-gcp-projectcredentials-path=/path/to/creds.jsontopic=debezium-events需启用 Pub/Sub API,服务账号需有 Publisher 权限。
Amazon SNSquarkus.amazon.sns.access-keysecret-keyregiontopic-arnaccess-key=AKIA...secret-key=xxxxxxregion=us-east-1topic-arn=arn:aws:sns:...SNS 为广播模式,适合通知类场景。
Amazon SQSquarkus.amazon.sqs.access-keysecret-keyregionqueue-urlqueue-url=https://sqs.us-east-1.amazonaws.com/...SQS 为队列模式,支持多个消费者。
Redisquarkus.redis.hostportpasswordstreams-keyhost=localhostport=6379password=mypasswordstreams-key=dbz-stream使用 Redis Streams 存储事件,支持消费者组。
Azure Event Hubsquarkus.azure.eventhubs.connection-stringevent-hub-nameconnection-string=Endpoint=sb://...event-hub-name=debezium-hub需配置 Event Hubs 实例。
Apache Pulsarquarkus.pulsar.service-urltopicservice-url=pulsar://localhost:6650topic=persistent://public/default/debezium支持 Pulsar 的持久化主题。
Infinispanquarkus.infinispan-client.server-listauth-usernameauth-passwordserver-list=127.0.0.1:11222auth-username=adminauth-password=password用于缓存同步场景。

通用 Debezium 配置(application.properties):

# 数据库连接
debezium.source.connector.class=io.debezium.connector.mysql.MySqlConnector
debezium.source.database.hostname=localhost
debezium.source.database.port=3306
debezium.source.database.user=mysqluser
debezium.source.database.password=mysqlpass
debezium.source.database.dbname=inventory
debezium.source.database.server.name=mysql-server

# 快照与过滤
debezium.source.snapshot.mode=initial
debezium.source.table.include.list=inventory.customers,inventory.orders

# 通用配置
debezium.source.offset.storage=file
debezium.source.offset.storage.file.filename=/tmp/offsets.dat
debezium.source.offset.flush.interval.ms=10000

启用特定输出(以 Redis 为例):

# 启用 Redis 输出
quarkus.debezium.sink.type=redis

# Redis 连接配置
quarkus.redis.hosts=redis://localhost:6379
quarkus.redis.password=mypassword
quarkus.redis.streams-key=dbz-events

11.3 启动与运行 Debezium Server

步骤操作命令/说明
1. 下载或构建从 GitHub 获取或自行构建git clone https://github.com/debezium/debezium-servercd debezium-servermvn clean package -Passembly
2. 解压并配置进入目录,编辑 conf/application.propertiestar -xzf debezium-server-dist-*.tar.gzcd debezium-servervim conf/application.properties
3. 启动服务使用启动脚本bin/debezium-server run
4. 后台运行使用 nohup 或 systemdnohup bin/debezium-server run &
5. 查看日志日志文件位置log/server.log
6. 停止服务发送 SIGTERMpkill -f debezium-server

Docker 运行示例:

docker run -it --rm \
  -v $(pwd)/conf:/debezium/conf \
  -v $(pwd)/log:/debezium/log \
  debezium/server:2.3

健康检查端点:

  • http://localhost:8080/q/health → 返回 UP 表示正常
  • http://localhost:8080/q/metrics → Prometheus 格式监控指标

11.4 与 Kafka Connect 模式的对比

对比维度Debezium Server(无 Kafka)Kafka Connect + Debezium Connector
架构复杂度低,单进程高,需 Kafka 集群、ZooKeeper(或 KRaft)、Schema Registry(推荐)
运维成本低,易于部署和监控高,需维护多个组件
吞吐能力中等,适合中小规模高,支持大规模并行处理
消息持久性依赖目标系统(如 Redis Streams、Pub/Sub)Kafka 提供高持久性、多副本、可回溯
多消费者支持取决于目标(如 Pub/Sub 支持,Redis Streams 支持消费者组)Kafka 天然支持多个消费者组
生态系统集成有限,仅支持官方支持的 Sink极强,可通过 Kafka Connect 集成数百种系统
流处理能力无,仅转发可与 Kafka Streams、ksqlDB、Flink 等结合进行复杂处理
部署方式独立进程、Docker、Kubernetes分布式集群,支持弹性扩展
适用场景轻量级、云服务集成、快速部署企业级数据管道、实时数仓、复杂事件处理

选择建议:

  • 若已使用 Kafka:选择 Kafka Connect 模式。
  • 若使用云服务且不想维护 Kafka:选择 Debezium Server。
  • 若需复杂流处理:必须使用 Kafka + Kafka Connect。

第十二章:实战案例与最佳实践

12.1 实时同步 MySQL 到 Elasticsearch

步骤操作配置说明
1. 部署组件MySQL、Debezium(Kafka Connect)、Kafka、Elasticsearch、Kafka Connect Elasticsearch Sink使用 Confluent Platform 或自行部署
2. 配置 Debezium SourceMySQL 连接器见下方 JSON
3. 配置 Elasticsearch Sink将 Kafka topic 数据写入 ES见下方 JSON
4. 验证同步在 MySQL 插入数据,检查 ES 是否更新INSERT INTO inventory.customers VALUES (1002, 'Carol', 'Engineer'); 然后查询 ES:GET /mysql-server.inventory.customers/_search
5. 处理嵌套字段使用 Flatten SMT 展平结构适用于 JSON 列或嵌套对象。

Debezium Source 配置:

{
  "name": "mysql-es-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": "mysql-server",
    "database.include.list": "inventory",
    "table.include.list": "inventory.customers",
    "database.history.kafka.bootstrap.servers": "kafka:9092",
    "database.history.kafka.topic": "schema-changes.inventory"
  }
}

Elasticsearch Sink 配置:

{
  "name": "es-sink",
  "config": {
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "topics": "mysql-server.inventory.customers",
    "connection.url": "http://elasticsearch:9200",
    "type.name": "_doc",
    "key.ignore": false,
    "schema.ignore": true,
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState"
  }
}

说明ExtractNewRecordState 提取 after 字段作为文档内容。

阶段组件作用说明
1. 捕获变更Debezium + MySQL Connector实时捕获 MySQL 增量变更输出为 Kafka topic,Avro 格式
2. 消息中间件Apache Kafka解耦、缓冲、持久化变更流多副本、高吞吐、可回溯
3. 流处理Apache Flink清洗、聚合、关联、维度退化SQL 或 DataStream API 编写作业
4. 目标数仓Snowflake / BigQuery / Redshift存储和分析支持 CDC 流式加载

Flink 作业示例(SQL):

-- 创建 Kafka 源表
CREATE TABLE mysql_customers (
  id INT,
  name STRING,
  op STRING,
  ts_ms BIGINT
) WITH (
  'connector' = 'kafka',
  'topic' = 'mysql-server.inventory.customers',
  'properties.bootstrap.servers' = 'kafka:9092',
  'format' = 'avro'
);

-- 创建数仓结果表(以 BigQuery 为例)
CREATE TABLE dwh_customers (
  customer_id INT,
  customer_name STRING,
  last_updated TIMESTAMP(3)
) WITH (
  'connector' = 'bigquery',
  'project-id' = 'my-project',
  'dataset' = 'dwh',
  'table' = 'customers'
);

-- 写入最新状态(根据 op 更新)
INSERT INTO dwh_customers
SELECT id, name, TO_TIMESTAMP(ts_ms / 1000)
FROM mysql_customers
WHERE op IN ('c', 'u');

优势:近实时(秒级)数据可见性,替代 T+1 批处理。

12.3 微服务间事件驱动架构(Event-Driven Architecture)

模式说明Debezium 角色示例
事件溯源(Event Sourcing)业务状态由事件流推导捕获数据库变更作为”事实”事件用户服务 → 写数据库 → Debezium → Kafka → 订单服务消费 UserCreated 事件
CQRS(命令查询职责分离)命令写主库,查询读物化视图同步主库变更到查询库主库(MySQL)← 写操作 → Debezium → ES/Redis → 查询服务
服务解耦服务间通过事件通信,无需直接调用作为”变更广播”机制支付服务更新订单状态 → Debezium → Kafka → 通知服务发送短信
审计与合规记录所有数据变更提供不可篡改的变更日志所有 op=d(删除)事件写入审计系统

关键设计:

  • 事件命名:<AggregateType><Action>,如 OrderCreatedCustomerDeleted
  • 使用 SMT 提取 after 并重命名字段以符合事件契约
  • 保证事件顺序:按主键分区,确保同一实体变更有序

12.4 安全配置:SSL、权限控制、敏感字段脱敏

安全维度配置项实现方式说明
传输加密SSL/TLS数据库连接:database.ssl.mode=required;Kafka:security.protocol=SSL;Debezium Server:各云服务 SDK 支持 TLS所有链路启用加密
身份认证凭据管理数据库:专用只读 + 复制权限账号;Kafka:SASL/PLAIN 或 SASL/SCRAM;云服务:IAM 角色/密钥遵循最小权限原则
敏感字段脱敏SMT 转换见下方配置在进入 Kafka 前屏蔽敏感信息
审计日志日志记录启用 errors.log.enable=true,集中收集日志(ELK/Splunk)记录所有连接、错误、配置变更
密钥管理外部化密钥使用 HashiCorp Vault、AWS Secrets Manager避免明文密码写入配置文件
网络隔离防火墙数据库仅允许 Connect/Server 主机访问;Kafka 集群内网部署减少攻击面

敏感字段脱敏 SMT 配置:

transforms=mask
transforms.mask.type=org.apache.kafka.connect.transforms.MaskField$Value
transforms.mask.fields=credit_card,ssn
transforms.mask.replacement=***

推荐实践:

  • 生产环境必须启用 SSL 和认证。
  • 敏感字段禁止进入 Kafka,使用 SMT 在源头脱敏。
  • 定期轮换数据库和 Kafka 凭据。