数据同步与 CDC
微服务拆分后,常需把 MySQL 变更 同步到缓存、ES、数仓或其它服务库。CDC(变更数据捕获) 解析 binlog / WAL,比定时全量扫表延迟低、对源库压力小。概念见 CDC;Canal 见 Canal;Debezium 见 Debezium。
典型链路
MySQL(业务库)
│ binlog ROW
▼
Canal Server / Debezium
▼
Kafka(推荐缓冲)
▼
Spring Boot 消费者 → 更新 Redis / ES / 从库
Canal → Kafka → Spring 消费
MySQL 前提
[mysqld]
log-bin=mysql-bin
binlog-format=ROW
binlog_row_image=FULL
server-id=1
sync_binlog=1ROW + FULL 便于消费者获得行的前后值;sync_binlog=1 缩小数据库已提交但 binlog 尚未持久化的崩溃窗口。binlog 保留时长必须覆盖 Canal 最长停机时间,否则位点对应日志被清理后只能重新做全量 + 增量同步。
Canal 投递 MQ
Canal 支持将 Entry 投递到 Kafka / RocketMQ(见 Canal 服务端配置)。
分区键与事件信封
同一业务实体的事件必须进入同一 partition:
Kafka key = database + '.' + table + ':' + primaryKeyKafka 只保证单 partition 内的顺序,不保证跨 partition 全局有序。Canal MQ 模式应按主键配置 partitionHash;只按表分区虽然有序,但会把整张大表压到单分区。
建议统一事件字段:
{
"eventId": "mysql-bin.000123:456789:2",
"database": "shop",
"table": "t_order",
"pk": "10001",
"operation": "UPDATE",
"version": 8,
"commitTs": 1770000000000,
"before": {},
"after": {}
}eventId:binlog 文件/位置/行序号或 GTID 派生值,用于去重与追踪;version:业务表单调版本,用于拒绝旧事件覆盖新状态;before/after:删除、索引更新和审计所需字段;- 时间戳只用于观测,不能替代 version,因为跨机器时钟可能漂移。
Spring Kafka 缓存失效示例
@KafkaListener(
topics = "canal.order_db",
groupId = "cache-sync",
containerFactory = "manualAckContainerFactory"
)
public void onBinlog(String payload, Acknowledgment ack) {
CanalMessage msg = parse(payload);
if (!"t_order".equals(msg.getTable())) {
ack.acknowledge();
return;
}
// DEL 幂等。数据库仍是事实源,下一次 cache miss 从 DB 重建。
redisTemplate.delete("order:" + msg.getPk());
// 只有 Redis 操作成功后才提交 offset。
// 若此处之前崩溃,消息会重复投递,但重复 DEL 没有副作用。
ack.acknowledge();
}| 要点 | 说明 |
|---|---|
| ACK 时机 | 先处理目标系统,成功后再提交 offset;反过来会漏消息 |
| 幂等 | 缓存优先 DEL;有业务副作用时用 eventId 唯一键去重 |
| 顺序 | 同一主键进入同一 Kafka partition |
| 乱序防护 | 更新缓存/ES 时比较业务 version,只接受更大版本 |
| 重试 | 原地重试会阻塞 partition;重试 Topic/DLQ 可能乱序,必须保留版本检查 |
| 延迟 | 同时监控 Canal binlog 位点差与 Kafka consumer lag |
直接更新缓存时的版本保护
删除缓存最安全;如果热点数据必须由消费者主动预热,则需要把“比较版本 + 写入”做成 Redis 原子操作:
-- KEYS[1] = cache key; ARGV[1] = new version; ARGV[2] = payload
local old = redis.call('HGET', KEYS[1], '_version')
if (not old) or (tonumber(ARGV[1]) > tonumber(old)) then
redis.call('HSET', KEYS[1], '_version', ARGV[1], 'data', ARGV[2])
redis.call('EXPIRE', KEYS[1], 3600)
return 1
end
return 0这能避免旧事件因重试或 DLQ 回放晚到,反向覆盖新缓存。version 应来自数据库行版本或可靠的每实体序号,不要仅用消费者本机时间。
有业务副作用时的 Inbox 去重
发送短信、记账、发券等不能靠重复执行。消费者应在自己的数据库本地事务中插入 processed_event(event_id UNIQUE) 并执行业务更新:唯一键冲突表示已经处理。Redis SET NX 可做短期快速去重,但 TTL 到期、故障切换和缓存丢失后不能充当永久业务凭证。
一致性与故障窗口
T1 MySQL COMMIT
T2 Canal 读取 binlog 并推进位点
T3 Kafka 持久化事件
T4 Consumer 更新 Redis/ES
T5 Consumer 提交 offset| 故障位置 | 现象 | 恢复方式 |
|---|---|---|
| T1 后、T2 前 | 下游暂时没看到 | Canal 从持久化位点续读 |
| T2 后、T3 前 | Canal/MQ 重试可能重复 | eventId + 幂等消费 |
| T4 后、T5 前 | 目标已更新但消息重放 | 幂等 DEL、UPSERT 或 Inbox |
| binlog 已清理但 Canal 未追上 | 无法从原位点续传 | 全量快照 + 记录新起点 + 增量追平 |
| MySQL 主库切换 | 文件位点/源地址变化 | GTID、Canal HA、验证新主复制历史连续 |
端到端默认按 at-least-once + 幂等设计。Kafka 生产或消费事务无法自动把外部 MySQL 与 Redis 纳入同一个原子事务。
业务表 CDC 与 Transactional Outbox
直接订阅业务表
适合构建缓存、ES、数仓等派生读模型:数据库行发生任何变化,下游按最终状态收敛。缺点是下游容易耦合物理表结构,也不一定能从一行 UPDATE 推导完整业务事件。
Transactional Outbox
同一个 MySQL 本地事务:
UPDATE t_order ...
INSERT outbox(event_id, aggregate_id, version, event_type, payload)
Canal/Debezium 订阅 outbox → Kafka → 下游服务Outbox 把业务状态和“应发布事件”原子提交,避免应用代码直接双写 DB + MQ。它更适合订单已创建、支付已完成这类稳定业务语义;缓存失效等纯读模型同步可直接监听业务表。
重放、对账与全量重建
- 保留足够的 binlog 和 Kafka 日志,明确可恢复时间目标;
- 记录每条事件的 source、GTID/position、partition/offset,便于端到端追踪;
- 对账任务按主键范围或
updated_at比较源表与目标版本; - 全量重建时先记录 CDC 起点,再导出一致性快照,最后回放起点后的增量;
- consumer lag 超过 SLA 时,关键读绕过缓存/读模型回源,避免继续返回无限期旧数据;
- DDL 和字段演进需要兼容期,消费者先支持新旧结构,再执行数据库变更。
Debezium(Kafka Connect)
多数据库、与 Kafka 生态一体,适合 MySQL / PostgreSQL / MongoDB 等异构源:
Debezium Connector → Kafka Connect → Topic {prefix}.{db}.{table}
Spring 侧用 @KafkaListener 消费 Envelope(before / after / op / source),分区键仍建议用主键保证顺序。Connector 注册、快照模式、Schema History、PostgreSQL 逻辑复制槽等详见 Debezium。
与 Canal 选型:
| 场景 | 推荐 |
|---|---|
| 仅 MySQL,已有 Canal 运维经验 | Canal |
| 多数据库 / PG / Mongo,统一 Kafka Connect | Debezium |
| 需要 Avro + Schema Registry | Debezium |
| Java Client 直连、不经 Kafka | Canal TCP 模式 |
其它同步方式
| 方式 | 工具 / 实现 | 文档 |
|---|---|---|
| 日志 CDC | Canal、Debezium、Maxwell | CDC |
| 批处理同步 | DataX、Sqoop | 离线全量 + 增量 |
| 应用双写 | 不推荐 | 一致性难保证 |
| 本地消息表 | Spring 定时发 MQ | 分布式与数据层 |
与微服务边界
| 做法 | 说明 |
|---|---|
| 推荐 | 单服务写自己的库,CDC 同步读模型(CQRS) |
| 避免 | CDC 回调直接跨服务 Feign 写库(耦合、难幂等) |
| 搜索索引 | MySQL → Canal → Kafka → 写 ES |
| 缓存失效 | DELETE/UPDATE 事件删 Redis key |
Spring Cloud 中的位置
CDC 消费者通常是 独立 Sync 服务(Spring Boot),注册到 Nacos,与业务服务解耦;不要求 Stream,但可与 消息驱动 统一 Binder。