{T}

场景案例:如何用 Flink SQL CDC 实现实时数据同步?

0. 引言

业务增长后,一份数据往往需要写入多种数据库:MySQL(主库)、Redis(缓存)、Elasticsearch(查询分析)。如何低成本地完成多库同步?本文从"改业务代码"与"中间件方案"的痛点出发,介绍 Flink SQL CDC 的原理与实现——只需几行 SQL 即可完成全量 + 增量实时同步。

1. 业务场景:多库同步的三种方案

1.1 方案一:业务代码直写多库(不推荐)

在业务代码中同时写入 MySQL、Redis、Elasticsearch。缺点:

  1. 代码严重耦合:每次修改都要重新测试、上线;
  2. 性能受损:业务系统多次写不同数据库;
  3. 全量同步靠人工:增量前需人工做一次全量。

1.2 方案二:消息中间件解耦

图表渲染中…

优点:业务与写入逻辑解耦、业务只需写一次 Kafka。缺点:仍需埋点改代码、需运维 Kafka、多库多表时开发量激增、全量同步仍需人工。

CDC(Change Data Capture,变化数据捕获) 是捕获数据变化的技术——MySQL 主从同步跟随 binlog 就是典型 CDC 场景。Flink CDC 用 Flink 实现 CDC:先全量同步,再跟随 binlog 实时增量同步,且全量/增量封装一体、支持 SQL,无需引入 Kafka、无需改业务代码

2. 实现原理:快照 + binlog 增量

以 MySQL → 目标库为例,Flink CDC 分两步:全量同步(快照)→ 增量同步(binlog)

2.1 全量同步:快照流程

图表渲染中…

关键点:可重复读事务 + 全程只读 → 无幻读 → 扫描数据与 binlog 位置时间点对齐,保证"一致快照"(不多不少、与源库完全相同)。

2.2 增量同步:binlog 跟随 + Checkpoint

全量完成后,跟随源库 binlog 将每次变更实时同步到目标库。Flink 周期性执行 Checkpoint 记录已同步的 binlog 位置,作业故障重启后从最近 Checkpoint 恢复。配合目标写入幂等,整个同步过程可达 exactly once 可靠性。

3. 具体实现

3.1 方式一:基于 DataStream

java
SourceFunction<String> sourceFunction = MySQLSource.<String>builder()
        .hostname("127.0.0.1").port(3306)
        .databaseList("db001")
        .username("root").password("123456")   // 测试用,生产禁用 root 与弱口令
        .deserializer(new StringDebeziumDeserializationSchema())
        .build();
// Elasticsearch 目标库 Sink(省略构造细节)
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.addSource(sourceFunction)          // 源:MySQL(封装全量+增量)
   .addSink(esSinkBuilder.build())     // 目标:Elasticsearch
   .setParallelism(1);                 // 保持消息顺序
env.execute("FlinkCdcDemo");

MySQLSource 连接器封装了全部复杂操作。缺点:写入 ES 的是整串字符串(未解析字段,查询效率低)、仍需写代码。

sql
-- 1. 配置源:MySQL
CREATE TABLE sourceTable (
  id INT, name STRING, counts INT, description STRING
) WITH (
  'connector' = 'mysql-cdc',
  'hostname' = '192.168.1.7', 'port' = '3306',
  'username' = 'root', 'password' = '123456',
  'database-name' = 'db001', 'table-name' = 'table001'
);

-- 2. 配置目标:Elasticsearch
CREATE TABLE sinkTable (
  id INT, name STRING, counts INT
) WITH (
  'connector' = 'elasticsearch-7',
  'hosts' = 'http://192.168.1.7:9200', 'index' = 'table001'
);

-- 3. 启动同步作业(可指定字段、支持 GROUP BY/Window 等复杂逻辑)
INSERT INTO sinkTable SELECT id, name, counts FROM sourceTable;

三个 SQL 完成全量 + 增量实时同步,且目标库字段与 SELECT 指定完全一致。Flink SQL 底层也会转化为 DataStream 执行。

4. 小结

  • 三种方案对比:业务直写(耦合高)→ Kafka 中间件(解耦但复杂)→ Flink CDC(无埋点、无 Kafka、全量+增量一体)
  • 原理:全局读锁 + 可重复读事务的一致快照 → binlog 位置 → 全表扫描 → Checkpoint 增量续传 → 幂等写入 = exactly once;
  • 两种实现:DataStream(灵活,适合 SQL 暂不支持的功能)与 Table & SQL(简单,适合常规同步);
  • 生产建议:禁 root 账号、弱口令;DataStream 与 SQL 两种方式都需掌握。

下一章讲解:竟然还有分布式的 JVM?