第一章: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 的核心原理与优势
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| 基于日志的 CDC | Debezium 直接读取数据库的事务日志(如 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 DMS | AWS 提供的托管式数据迁移与 CDC 服务。 | 全托管,支持异构数据库迁移,图形化界面。 | 成本较高,锁定云厂商。 | 适合云上环境且预算充足的团队。 |
| Flink CDC | 基于 Flink SQL 的 CDC 方案,可直接在 Flink 中捕获变更。 | 无需 Kafka Connect,端到端流处理一体化。 | 依赖 Flink 运行时,学习成本高。 | 适合已使用 Flink 的流处理场景。 |
第二章:Debezium 核心架构与组件
2.1 Connectors(连接器)概述
| 连接器类型 | 支持数据库 | 说明 | 注意事项 |
|---|---|---|---|
| MySQL Connector | MySQL | 基于 binlog ROW 模式捕获变更,支持快照和增量。 | 需配置 server-id 和 binlog_format=ROW。 |
| PostgreSQL Connector | PostgreSQL | 使用逻辑复制(Logical Replication)读取 WAL 日志。 | 需创建复制槽(Replication Slot)和复制用户。 |
| MongoDB Connector | MongoDB | 基于变更流(Change Streams)捕获副本集或分片集群的变更。 | 必须使用副本集部署,单节点不支持。 |
| SQL Server Connector | SQL Server | 使用变更数据捕获(CDC)或变更跟踪(CT)功能。 | 需启用数据库和表级别的 CDC。 |
| Oracle Connector | Oracle | 基于 LogMiner 或 XStream API 读取 redo log。 | 配置复杂,需 Oracle 高级权限。 |
| Db2 Connector | IBM 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 Connect | Apache Kafka 提供的可扩展、容错的数据集成框架,用于连接 Kafka 与其他系统。 | 支持 Source(数据源)和 Sink(目标)连接器。 |
| Debezium 作为 Source Connector | Debezium 实现了 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
| 组件 | 版本要求 | 说明 | 注意事项 |
|---|---|---|---|
| Java | JDK 8 或 JDK 11(推荐) | Kafka 与 Kafka Connect 基于 Java 开发,需安装 JDK。 | 建议使用 OpenJDK 或 Oracle JDK,避免使用 JRE。 |
| Apache Kafka | 2.8+(推荐 3.x) | 消息中间件,用于传输 Debezium 产生的变更事件。 | 需提前下载并解压 Kafka 发行包。 |
| ZooKeeper | 3.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.jar、mysql-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.servers | Kafka 集群地址 | localhost:9092 | 确保 Kafka 已启动。 |
| group.id | Connect 集群标识 | connect-cluster | 同一集群所有 Worker 使用相同 ID。 |
| key.converter | Key 序列化方式 | org.apache.kafka.connect.json.JsonConverter | 建议生产环境使用 Avro + Schema Registry。 |
| value.converter | Value 序列化方式 | org.apache.kafka.connect.json.JsonConverter | 同上。 |
| key.converter.schemas.enable | 是否在序列化中包含 schema | true | 便于解析结构化数据。 |
| value.converter.schemas.enable | 是否在 value 中包含 schema | true | 必须开启以支持 Debezium 事件格式。 |
| config.storage.topic | 存储连接器配置的 topic | connect-configs | 自动创建,需设置 replication.factor ≥ 3(生产)。 |
| offset.storage.topic | 存储 offset 的 topic | connect-offsets | 同上。 |
| status.storage.topic | 存储连接器状态的 topic | connect-status | 同上。 |
| rest.port | REST 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"。 |
| 检查特定连接器是否存在 | 查看返回列表中是否包含 MySqlConnector、PostgresConnector 等。 | 响应示例见下方。 | 若未出现,检查 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_format | Binlog 格式必须为 ROW 模式 | ROW | 在 my.cnf 中设置 binlog-format=ROW。 |
| binlog_row_image | 行日志内容完整性 | FULL | 确保 UPDATE 前后镜像都能捕获。 |
| server-id | 每个 MySQL 实例必须有唯一 ID | 1、2 等 | 若为 0,Debezium 将无法读取 binlog。 |
| log_bin | 启用 binlog | mysql-bin.log | 必须开启。 |
| expire_logs_days | binlog 保留天数 | 7 | 避免日志过早清理导致 Debezium 断流。 |
| 用户权限 | Debezium 连接用户所需权限 | SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT | 可使用如下 SQL 创建用户。 |
| 数据库引擎 | 表必须使用 InnoDB | InnoDB | MyISAM 不支持行级变更。 |
创建用户 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 | 整数 | 最大任务数 | 1 | MySQL 通常设为 1(单点写入)。 |
| database.hostname | 字符串 | MySQL 主机名 | localhost | 支持 IP 或域名。 |
| database.port | 整数 | MySQL 端口 | 3306 | 默认 3306。 |
| database.user | 字符串 | 连接用户名 | debezium | 需具备复制权限。 |
| database.password | 字符串 | 连接密码 | dbz | 明文存储,生产建议使用密钥管理。 |
| database.server.id | 整数 | MySQL server-id | 223344 | 必须唯一,避免与主从冲突。 |
| database.server.name | 字符串 | 逻辑服务器名,用于 Kafka topic 命名 | dbserver1 | topic 将以 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 | 字符串 | 快照模式 | initial | initial: 首次全量 + 增量;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 值 | before | after | 识别方式 | 处理建议 |
|---|---|---|---|---|---|
| INSERT | c | null | 包含数据 | op == 'c' | 将 after 数据插入目标系统。 |
| UPDATE | u | 包含旧值 | 包含新值 | op == 'u' | 使用 after 更新目标,可用 before 做条件更新。 |
| DELETE | d | 包含旧值 | null | op == 'd' | 根据 before 中的主键删除目标数据。 |
| Snapshot(初始快照) | r | null | 包含数据 | op == 'r' | 通常作为全量同步,需注意去重。 |
注意事项:
- 所有事件均通过 Kafka topic 发送,topic 名为
{database.server.name}.{database}.{table},如dbserver1.inventory.customers。 - 消费者需解析 JSON 或 Avro 格式事件,提取
op字段判断操作类型。 - 建议消费者实现幂等性,避免重复处理导致数据错乱。
第五章:PostgreSQL CDC 实践
5.1 PostgreSQL 的逻辑复制(Logical Replication)配置
| 配置项 | 说明 | 推荐值/操作 | 注意事项 |
|---|---|---|---|
| wal_level | WAL 日志级别,必须为 logical 才支持逻辑复制 | logical | 在 postgresql.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 包含 pgoutput | pgoutput 是 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 | 整数 | 最大任务数 | 1 | PostgreSQL 通常设为 1。 |
| database.hostname | 字符串 | PostgreSQL 主机名 | localhost | 支持 IP 或域名。 |
| database.port | 整数 | PostgreSQL 端口 | 5432 | 默认端口。 |
| database.user | 字符串 | 连接用户名 | debezium | 需具备复制权限。 |
| database.password | 字符串 | 连接密码 | dbz | 明文存储,生产建议使用密钥管理工具。 |
| database.dbname | 字符串 | 要连接的数据库名 | inventory | 必填。 |
| database.server.name | 字符串 | 逻辑服务器名,用于 Kafka topic 命名 | pgserver1 | topic 格式:pgserver1.public.customers。 |
| plugin.name | 字符串 | 逻辑解码插件名 | pgoutput | PostgreSQL 官方插件,Debezium 默认使用。 |
| slot.name | 字符串 | 复制槽名称 | debezium_slot | 必须唯一,不能包含 -,推荐使用字母数字下划线。 |
| snapshot.mode | 字符串 | 快照模式 | initial | initial: 初始快照 + 增量;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" 表示正常。 |
| 暂停连接器 | 暂停捕获,不提交 offset | curl -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 用户需具备 clusterMonitor 和 read 权限 | 见下方创建用户命令 | clusterMonitor 用于访问集群状态,read 用于读取数据。 |
| 认证机制 | 支持 SCRAM-SHA-1 或 SCRAM-SHA-256 | 配置 authenticationMechanism | 默认自动检测。 |
| SSL/TLS | 可选加密连接 | 配置 ssl.enabled=true | 生产环境建议启用。 |
| Storage Engine | 存储引擎 | 必须为 WiredTiger | MMAPv1 不支持变更流。 |
创建用户:
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 | 整数 | 任务数 | 1 | MongoDB Connector 通常为 1。 |
| mongodb.hosts | 字符串 | MongoDB 主机地址 | rs0/localhost:27017 | 格式:<replica_set_name>/host:port。 |
| mongodb.user | 字符串 | 用户名 | debezium | 具备 clusterMonitor 和 read 权限。 |
| 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 历史 topic | schema-changes.sales | 自动创建。 |
| snapshot.mode | 字符串 | 快照模式 | initial | initial: 先全量再增量;never: 仅增量。 |
| heartbeat.interval.ms | 整数 | 心跳间隔(毫秒) | 10000 | 保持连接活跃,避免超时。 |
| poll.interval.ms | 整数 | 轮询变更流间隔 | 1000 | 控制延迟与负载。 |
| ssl.enabled | 布尔值 | 是否启用 SSL | false 或 true | 生产建议启用。 |
6.3 变更流(Change Streams)事件解析
| 字段名 | 说明 | 示例值 | 注意事项 |
|---|---|---|---|
| operationType | 操作类型 | insert, update, delete, replace, drop, rename 等 | 对应 MongoDB 操作。 |
| fullDocument | INSERT/REPLACE 时的完整文档 | { "_id": "abc", "name": "Alice" } | UPDATE/DELETE 为 null。 |
| ns.db | 数据库名 | sales | 源数据库。 |
| ns.coll | 集合名 | customers | 源集合。 |
| documentKey | 文档主键(_id) | { "_id": "abc" } | 所有事件都包含。 |
| updateDescription | UPDATE 操作的修改详情 | { "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 | 字符串 | 数据库主机名或 IP | localhost 或 192.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 | 生产环境建议设为 required 或 verified。 |
| 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) | 包含的 schema | public,app_data | PostgreSQL 适用。 |
| schema.exclude.list | 字符串 | 排除的 schema | pg_catalog,information_schema | — |
| table.include.list | 字符串 | 包含的表/集合 | inventory.customers,inventory.orders | 格式:db.table 或 schema.table。 |
| table.exclude.list | 字符串 | 排除的表/集合 | %.audit_log,%.temp_% | 支持通配符 % 和 _。 |
| collection.include.list | 字符串(MongoDB) | 包含的集合 | sales.orders,sales.customers | MongoDB 专用。 |
| collection.exclude.list | 字符串(MongoDB) | 排除的集合 | sales.logs | — |
| topic.regex.replacement | 字符串 | 自定义 topic 命名规则 | $1-$2-$3 | 配合 SMT 使用,重写 topic 名称。 |
注意:
include.list和exclude.list不能同时用于同一层级(如不能同时设置table.include.list和table.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 中包含 schema | true | 建议开启,便于解析结构化数据。 |
| value.converter.schemas.enable | 布尔值 | 是否在 value 中包含 schema | true | 必须开启以支持 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 topic | connect-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(单消息转换)介绍与常用转换器
| 转换器名称 | 类名 | 用途 | 注意事项 |
|---|---|---|---|
| ExtractField | org.apache.kafka.connect.transforms.ExtractField | 从结构体中提取字段作为 value 或 key | 支持嵌套字段($.field.subfield)。 |
| ValueToKey | org.apache.kafka.connect.transforms.ValueToKey | 将 value 中的某个字段提取为消息 key | 常用于按主键分区。 |
| ReplaceField | org.apache.kafka.connect.transforms.ReplaceField | 重命名、排除或类型转换字段 | 支持嵌套字段操作。 |
| MaskField | org.apache.kafka.connect.transforms.MaskField | 对字段值进行掩码(如手机号、身份证) | 用于脱敏场景。 |
| Filter | org.apache.kafka.connect.transforms.Filter | 过滤掉满足条件的消息 | 可结合 predicate 实现复杂逻辑。 |
| FilterTopic | org.apache.kafka.connect.transforms.FilterTopic | 根据 topic 名称过滤消息 | 非 record 级别过滤。 |
| InsertField | org.apache.kafka.connect.transforms.InsertField | 添加静态或动态字段(如 timestamp、header) | 支持 RecordHeader、timestamp 等。 |
| HoistField | org.apache.kafka.connect.transforms.HoistField | 将整个 value 包装为一个字段 | 如将 {name: "Alice"} 变为 {"value": {name: "Alice"}}。 |
| Flatten | org.apache.kafka.connect.transforms.Flatten | 展平嵌套结构(用 . 连接字段名) | 如 {"addr.city"} 替代 {"addr": {"city": ...}}。 |
| TimestampRouter | org.apache.kafka.connect.transforms.TimestampRouter | 根据时间戳重写 topic 名称 | 如按天分表:topic-${YYYYMMdd}。 |
| RegexRouter | org.apache.kafka.connect.transforms.RegexRouter | 使用正则表达式重写 topic 名称 | 强大灵活,可用于路由。 |
SMT 链说明:多个 SMT 可串联执行,顺序由配置名决定,如
transforms=A,B,C。
8.2 使用 SMT 进行字段过滤、重命名、类型转换
| 操作 | 配置示例 | 说明 | 注意事项 |
|---|---|---|---|
| 字段过滤 | 见下方示例 | 排除 credit_card 和 ssn 字段 | 使用 blacklist 排除,whitelist 保留指定字段。 |
| 字段重命名 | 见下方示例 | 将 col1 改为 field1 | 多个用逗号分隔。 |
| 类型转换 | 不直接支持类型转换 | 需结合 Schema 或下游处理 | SMT 不改变字段物理类型,仅逻辑重命名。 |
| 提取字段为 Key | 见下方示例 | 使用 id 字段作为 Kafka 消息 key | schema=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 字段。 |
| 仅捕获 UPDATE | predicates.IsUpdate.value=u | 将 value 改为 u | — |
| 仅捕获 DELETE | predicates.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=500、poll.timeout.ms=30000、max.queue.size=8192、max.batch.size=2048 | 提高拉取频率和队列容量。 |
| 网络延迟高 | 调整心跳与超时 | heartbeat.interval.ms=10000、heartbeat.action.query=SELECT 1 | PostgreSQL/MongoDB 适用。 |
| 消费者积压(Lag) | 增加消费者或优化处理逻辑 | — | Debezium 本身是 Source,需下游消费者优化。 |
| 快照期间负载高 | 分批快照或低峰期执行 | snapshot.mode=initial、snapshot.delay.ms=60000 | 避免影响线上业务。 |
| Binlog/WAL 清理过快 | 延长日志保留时间 | MySQL: expire_logs_days=7,PostgreSQL: wal_keep_size=1GB | 防止 Debezium 断流。 |
| 连接中断自动恢复 | 启用错误重试 | errors.tolerance=all、errors.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
}
不同操作的事件结构对比:
| 操作 | before | after | op |
|---|---|---|---|
| INSERT | null | 新数据 | c |
| UPDATE | 旧数据 | 新数据 | u |
| DELETE | 旧数据 | null | d |
| 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 Registry | http://localhost:8081 |
| Avro Converter | 将数据序列化为 Avro 并注册 schema | 配置在 Kafka Connect 中 | io.confluent.connect.avro.AvroConverter |
| key.converter | Key 序列化器 | key.converter | io.confluent.connect.avro.AvroConverter |
| value.converter | Value 序列化器 | value.converter | io.confluent.connect.avro.AvroConverter |
| schema.registry.url | Schema Registry 地址 | key.converter.schema.registry.url、value.converter.schema.registry.url | http://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 / DATETIME | org.apache.kafka.connect.data.Timestamp | 见下方 | 转换为毫秒级时间戳 | 值为 1970-01-01 起的毫秒数。 |
| DATE | org.apache.kafka.connect.data.Date | 见下方 | 转换为天数(自 1970-01-01) | — |
| TIME | org.apache.kafka.connect.data.Time | 见下方 | 毫秒数(当日零点起) | — |
| JSON (MySQL/PostgreSQL) | org.apache.kafka.connect.data.Json | 见下方 | 存储为字符串,但标记为 JSON 类型 | 实际值为 JSON 字符串。 |
| UUID | org.apache.kafka.connect.data.Decimal 或 string | 通常映射为 string | PostgreSQL 中为 uuid 类型 | 可通过 SMT 转换。 |
| BYTEA / BLOB | bytes | { "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 方法 | 说明 | 示例请求 | 响应关键字段 |
|---|---|---|---|---|
/connectors | GET | 列出所有连接器 | curl http://localhost:8083/connectors | ["mysql-connector", "pg-connector"] |
/connectors/{name} | GET | 获取连接器配置 | curl http://localhost:8083/connectors/mysql-connector | 返回完整配置 JSON |
/connectors/{name}/status | GET | 获取连接器运行状态 | curl http://localhost:8083/connectors/mysql-connector/status | "state": "RUNNING", "tasks": [{"state": "RUNNING"}] |
/connectors/{name}/tasks | GET | 获取任务列表 | curl http://localhost:8083/connectors/mysql-connector/tasks | 任务 ID、状态、worker_id |
/connectors/{name}/tasks/{taskid}/status | GET | 获取单个任务状态 | curl http://localhost:8083/connectors/mysql-connector/tasks/0/status | 同上,更细粒度 |
/connectors/{name}/config | GET | 获取配置 | curl http://localhost:8083/connectors/mysql-connector/config | 仅返回 config 部分 |
/connectors/{name}/config | PUT | 更新配置(热更新) | 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-cluster | Confluent Control Center、Prometheus + Grafana |
| Offset 提交频率 | offset.flush.interval.ms 控制 | 查看 Kafka topic connect-offsets 的消息速率 | Kafka Manager |
| Debezium 内部 Metrics | JMX 暴露的指标 | 启用 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 removed | binlog/WAL 被清理,连接器断流 | Could not find first log file name from binlog index、requested 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 found | PostgreSQL slot 被手动删除 | replication slot "debezium_slot" does not exist | 重建 slot 或修改 slot.name;避免手动删除。 |
| Oplog entry not found | MongoDB oplog 太小,记录被覆盖 | CappedPositionLost | 增大 oplog 大小(mongod --oplogSize 4096);重启连接器触发快照。 |
| Out of memory | JVM 内存不足 | 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,需手动清理! |
安全删除流程:
DELETE /connectors/{name}- 检查数据库端资源:
- PostgreSQL:
SELECT * FROM pg_replication_slots;→pg_drop_replication_slot('slot_name'); - MySQL:无额外资源。
- MongoDB:无额外资源。
- PostgreSQL:
- (可选)清理 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/Sub | quarkus.google.cloud.pubsub. | project-id、credentials-path、topic | project-id=my-gcp-project、credentials-path=/path/to/creds.json、topic=debezium-events | 需启用 Pub/Sub API,服务账号需有 Publisher 权限。 |
| Amazon SNS | quarkus.amazon.sns. | access-key、secret-key、region、topic-arn | access-key=AKIA...、secret-key=xxxxxx、region=us-east-1、topic-arn=arn:aws:sns:... | SNS 为广播模式,适合通知类场景。 |
| Amazon SQS | quarkus.amazon.sqs. | access-key、secret-key、region、queue-url | queue-url=https://sqs.us-east-1.amazonaws.com/... | SQS 为队列模式,支持多个消费者。 |
| Redis | quarkus.redis. | host、port、password、streams-key | host=localhost、port=6379、password=mypassword、streams-key=dbz-stream | 使用 Redis Streams 存储事件,支持消费者组。 |
| Azure Event Hubs | quarkus.azure.eventhubs. | connection-string、event-hub-name | connection-string=Endpoint=sb://...、event-hub-name=debezium-hub | 需配置 Event Hubs 实例。 |
| Apache Pulsar | quarkus.pulsar. | service-url、topic | service-url=pulsar://localhost:6650、topic=persistent://public/default/debezium | 支持 Pulsar 的持久化主题。 |
| Infinispan | quarkus.infinispan-client. | server-list、auth-username、auth-password | server-list=127.0.0.1:11222、auth-username=admin、auth-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-server → cd debezium-server → mvn clean package -Passembly |
| 2. 解压并配置 | 进入目录,编辑 conf/application.properties | tar -xzf debezium-server-dist-*.tar.gz → cd debezium-server → vim conf/application.properties |
| 3. 启动服务 | 使用启动脚本 | bin/debezium-server run |
| 4. 后台运行 | 使用 nohup 或 systemd | nohup bin/debezium-server run & |
| 5. 查看日志 | 日志文件位置 | log/server.log |
| 6. 停止服务 | 发送 SIGTERM | pkill -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 Source | MySQL 连接器 | 见下方 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字段作为文档内容。
12.2 构建实时数仓:MySQL → Kafka → Flink → Data Warehouse
| 阶段 | 组件 | 作用 | 说明 |
|---|---|---|---|
| 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>,如OrderCreated、CustomerDeleted - 使用 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 凭据。