亚马逊AWS官方博客
大规模数据库迁移中的 CDC 吞吐优化实践
摘要:当 AWS Database Migration Service (DMS) 单任务无法追平高写入量源库的 CDC 延迟时,如何在不增加源库压力的前提下提升同步吞吐?本文介绍基于 Amazon MSK Connect + Debezium 的读写解耦方案:单连接读取 binlog,多 Task 并行写入 Aurora MySQL,实现 N 倍吞吐提升。
一、概述
在企业将自建 My SQL 数据库迁移到 AWS 的过程中,持续数据复制 (CDC) 是保证业务不中断的关键技术手段。然而,当源库写入量较大时,AWS DMS (Database Migration Service) 的单任务复制模式可能无法追平源端的数据变更,导致迁移割接窗口无法收敛。
本文介绍如何利用 Amazon MSK Connect 配合开源的 Debezium MySQL Connector,构建一条从自建 MySQL 到 Amazon Aurora MySQL 的高吞吐 CDC 同步管道。该方案的核心优势在于:通过多 Task 并行消费和写入,线性提升同步吞吐量,解决 DMS 单任务模式下数据延迟无法追平的问题。
适用场景说明:本方案主要针对自建 MySQL(IDC/其他云)迁移到 AWS 的场景。如果源端已经是 Amazon RDS for MySQL,迁移到 Aurora 有更简单的方案(如 Aurora Read Replica 提升、RDS 快照恢复、或原生 binlog 复制)。
二、痛点分析:DMS 多任务并行的代价
DMS 确实可以通过创建多个任务(每个任务负责不同的表)来提升并行度。但这种方式有一个根本性问题:
每多拆一个 DMS 任务,就会多一条到源库的 binlog 连接。
- 3 个 DMS 任务 = 3 条 binlog 读取连接,源库需要并行服务 3 次 binlog dump
- 10 个 DMS 任务 = 10 条 binlog 连接,对源库 I/O 和 CPU 造成显著压力
- 在源库本身已经是高负载的情况下,额外的 binlog 读取开销可能导致源库性能进一步恶化
这在生产迁移中是不可接受的——迁移的前提是不能影响线上业务。
三、方案选型:Debezium + MSK Connect 如何解决?
Debezium 方案的核心优势在于读写解耦:
- 源端只需一条 binlog 连接:Debezium Source Connector 通过单个数据库连接读取全部 binlog 事件,然后按表自动分发到不同的 Kafka Topic
- 目标端多 Task 并行写入:JDBC Sink Connector 可配置多个 Task,从不同的 Topic/Partition 并行消费并写入目标库
换言之:1 次 binlog 读取 → N 倍并行写入,在不增加源库压力的前提下线性提升同步吞吐量。
| 对比维度 | DMS 多任务拆分 | Debezium + MSK Connect |
| 源库 binlog 连接数 | N 条(每个任务一条) | 仅 1 条 |
| 对源库性能影响 | 线性增长,任务越多压力越大 | 恒定,与并行度无关 |
| 写入并行度 | N 个任务并行 | N 个 Task 并行 |
| 追平能力 | 受限于源库承受能力 | 仅受限于目标库写入能力 |
| 按表路由到 Kafka Topic | 不支持 | 原生支持 |
| 运维模式 | 全托管 | 全托管(MSK Connect) |
| 配置复杂度 | 低 | 中等(一次性配置) |
四、架构设计
整体架构分为三层:
[图. CDC 高吞吐同步架构 — 1 次 Binlog 读取 → N 倍并行写入] |
4.1 数据捕获层(Source Connector)
Debezium MySQL Connector 通过读取自建 MySQL 源库的 binlog(二进制日志)捕获数据变更事件。它以单连接方式读取 binlog 流,对源库的性能影响极小。每张表的变更事件自动写入独立的 Kafka Topic,命名规则为 {prefix}.{database}.{table}。
4.2 消息传输层(Amazon MSK)
Amazon MSK 作为消息中间件,承担变更事件的可靠传输和缓冲。每个 Topic 可配置多个 partition,为下游消费提供并行能力。
4.3 数据写入层(Sink Connector)
JDBC Sink Connector 从 Kafka Topic 中消费变更事件,通过 Upsert 模式将数据写入目标 Aurora MySQL。通过配置多个 Task 实现并行写入,充分利用 Aurora 的写入吞吐能力。
五、同步能力
| 操作类型 | 支持情况 |
| INSERT | ✅ 支持 |
| UPDATE | ✅ 支持 |
| DELETE | ✅ 支持 |
| DDL 变更 | 记录到独立 Topic(不自动执行) |
端到端延迟:在典型负载下,从源库写入到目标库可查,延迟约 2-10 秒。
六、前置条件
- Amazon MSK 集群(建议 Provisioned 模式,kafka.m5.large,3 个 Broker)
- MSK 集群配置已开启 auto.create.topics.enable=true
- 源端自建 MySQL 已开启 binlog(binlog_format=ROW)
- 目标端 Aurora MySQL 已创建对应的数据库和表结构
- IAM 角色(信任 kafkaconnect.amazonaws.com)
- S3 存储桶(用于存放 Connector 插件包)
- 源库与 MSK 集群之间的网络连通(VPN/Direct Connect/VPC Peering)
七、部署步骤
7.1 步骤一:准备 Connector 插件
下载并打包 Debezium MySQL Connector 和 JDBC Sink Connector(含 MySQL 驱动和 Debezium SMT 依赖),上传到 S3 存储桶。在 MSK Connect 控制台创建 Custom Plugin。
7.2 步骤二:配置 Debezium Source Connector
核心配置参数:
| 参数 | 作用 |
| topic.prefix=cdc | 生成的 Topic 名为 cdc.<库名>.<表名> |
| snapshot.mode=schema_only | 仅捕获表结构,从当前 binlog 位置开始增量同步 |
| schemas.enable=true | 消息中携带 schema 信息,JDBC Sink 需要此信息解析字段类型 |
| time.precision.mode=connect | 避免 Debezium 默认的微秒精度导致下游兼容问题 |
| tasks.max=1 | 源端仅需 1 个 Task,单连接读取全部 binlog |
7.3 步骤三:配置 JDBC Sink Connector
⚠️ 重要:必须在 Source Connector 运行并创建 Topic 后,再创建 Sink Connector。
| 参数 | 作用 |
| topics.regex=cdc\.mydb\..* | 自动匹配所有按表分的 Topic |
| insert.mode=upsert | 根据主键执行 INSERT 或 UPDATE |
| delete.enabled=true | 同步 DELETE 操作 |
| transforms.unwrap | 展平 Debezium 消息信封,提取实际数据 |
| transforms.routeTopic | 将 cdc.mydb.<表名> 映射为目标表名 |
| transforms.convertTimestamp | ISO 8601 转为 MySQL 兼容的 yyyy-MM-dd HH:mm:ss |
| tasks.max=N | 多 Task 并行写入,根据表数量和数据量调整 |
八、并行度调优
| 场景 | 推荐配置 |
| 2-10 张表 | 单个 JDBC Sink,tasks.max=2~4 |
| 10-50 张表 | 按业务域拆分 2-3 个 Sink Connector,每个 tasks.max=4~8 |
| 50+ 张表 | 多个 Sink Connector + Topic partition 数 ≥ task 数 |
九、常见问题与解决方案
问题 1:Debezium 时间格式与 MySQL 不兼容
现象:Debezium 输出的 ZonedTimestamp 为 ISO 8601 格式,MySQL 的 TIMESTAMP/DATETIME 列无法接受。
解决方案:使用 TimestampConverter SMT 将时间字段转换为 yyyy-MM-dd HH:mm:ss 格式。
问题 2:修改 schemas.enable 后消息格式不一致
现象:Topic 中同时存在新旧格式的消息,导致 Sink Connector 解析失败。
解决方案:删除相关 Kafka Topic 并重新创建,确保消息格式统一。
问题 3:JDBC Sink 报 ClassNotFoundException
现象:配置 ExtractNewRecordState SMT 后,Connector 启动时报类找不到。
解决方案:JDBC Sink 插件包中需要额外包含 debezium-core 和 debezium-api 两个 jar 文件。
问题 4:MSK Connect 更新插件不生效
现象:更新 S3 上的 zip 文件后,MSK Connect 仍使用旧版本。
解决方案:MSK Connect 的 Custom Plugin 是不可变的。需要删除旧 Plugin 并重新创建。
问题 5:MSK 集群不自动创建 Topic
现象:Debezium 启动后报错,无法创建按表分的 Topic。
解决方案:在 MSK 集群配置中确认 auto.create.topics.enable=true 已开启。
十、总结
本文介绍了一种基于 Amazon MSK Connect 和 Debezium 的全托管高吞吐 CDC 同步方案,专为自建 MySQL 迁移到 AWS 时解决数据追平延迟问题而设计。
核心优势
- 并行写入解决追平问题:多 Task 并行消费和写入,线性提升目标库写入吞吐量
- 源库零额外压力:仅 1 条 binlog 连接,与写入并行度无关
- 全托管:基于 MSK Connect,无需管理服务器或 Kafka Connect 集群
- 按表路由:每张表自动拥有独立的 Kafka Topic,支持细粒度的并行消费
- 低延迟:端到端 2-10 秒,满足大多数实时同步场景
何时选择此方案
- 自建 MySQL 迁移到 AWS,DMS 的 CDC 复制延迟持续增长无法追平
- 源库写入 TPS 高,需要更高的同步吞吐能力且不能增加源库压力
- 需要按表将变更数据路由到多个下游系统
何时选择其他方案
- RDS for MySQL → Aurora:使用 Aurora Read Replica 提升或 RDS 快照恢复
- 源库写入量不大、延迟可接受:DMS 即可满足需求
- 仅需一次性全量迁移:使用 mysqldump/XtraBackup + DMS 增量追平
➡️ 下一步行动:
相关产品:
- Amazon Connect — AI 客户体验解决方案
- Amazon MSK — 完全托管式 Apache Kafka 服务
- Amazon DMS — 数据库迁移服务
- Amazon Aurora — 适用于 PostgreSQL、MySQL 和 DSQL 的无服务器关系数据库服务
- Amazon RDS — 完全托管的关系数据库服务
相关文章:
十一、参考资料
- Debezium Connector for MySQL
- Amazon MSK Connect 开发者指南
- Confluent JDBC Sink Connector
- Debezium ExtractNewRecordState
- Amazon MSK 最佳实践
*前述特定亚马逊云科技生成式人工智能相关的服务目前在亚马逊云科技海外区域可用。亚马逊云科技中国区域相关云服务由西云数据和光环新网运营,具体信息以中国区域官网为准。
本篇作者
AWS 架构师中心:云端创新的引领者探索 AWS 架构师中心,获取经实战验证的最佳实践与架构指南,助您高效构建安全、可靠的云上应用 |
![]() |

