数据同步与 CDC

→ 返回 Spring Cloud

微服务拆分后,常需把 MySQL 变更 同步到缓存、ES、数仓或其它服务库。CDC(变更数据捕获) 解析 binlog / WAL,比定时全量扫表延迟低、对源库压力小。概念见 CDC;Canal 见 Canal;Debezium 见 Debezium。


典型链路

MySQL(业务库)
  │ binlog ROW
  ▼
Canal Server / Debezium
  ▼
Kafka(推荐缓冲)
  ▼
Spring Boot 消费者 → 更新 Redis / ES / 从库

与 Kafka 集成、消息队列 配合。


Canal → Kafka → Spring 消费

MySQL 前提

[mysqld]
log-bin=mysql-bin
binlog-format=ROW
binlog_row_image=FULL
server-id=1
sync_binlog=1

ROW + FULL 便于消费者获得行的前后值;sync_binlog=1 缩小数据库已提交但 binlog 尚未持久化的崩溃窗口。binlog 保留时长必须覆盖 Canal 最长停机时间,否则位点对应日志被清理后只能重新做全量 + 增量同步。

Canal 投递 MQ

Canal 支持将 Entry 投递到 Kafka / RocketMQ(见 Canal 服务端配置)。

分区键与事件信封

同一业务实体的事件必须进入同一 partition:

Kafka key = database + '.' + table + ':' + primaryKey

Kafka 只保证单 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 ConnectDebezium
需要 Avro + Schema RegistryDebezium
Java Client 直连、不经 KafkaCanal TCP 模式

其它同步方式

方式工具 / 实现文档
日志 CDCCanal、Debezium、MaxwellCDC
批处理同步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。


相关链接

Spring Cloud

中间件

架构 / 数据