🌊 前言:为什么传统的同步方案“不香”了?
在 Flink CDC 普及之前,要实现 MySQL 到 ES 的实时同步,架构师通常会画出这样一张图:
MySQL -> Canal/Debezium -> Kafka -> Logstash/Java应用 -> Elasticsearch。
痛点非常明显:
Flink CDC 的出现,直接把中间商(Kafka、Canal)给“干掉”了。 它实现了从 MySQL Binlog 直接读取数据并写入 ES,全量历史数据 + 增量实时数据自动无缝切换。
⚔️ 一、 架构对比:极简主义的胜利
传统架构 vs Flink CDC 架构 (Mermaid):
#mermaid-svg-FDvXW3ZpPWDH4mGK{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-FDvXW3ZpPWDH4mGK .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-FDvXW3ZpPWDH4mGK .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-FDvXW3ZpPWDH4mGK .error-icon{fill:#552222;}#mermaid-svg-FDvXW3ZpPWDH4mGK .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-FDvXW3ZpPWDH4mGK .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-FDvXW3ZpPWDH4mGK .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-FDvXW3ZpPWDH4mGK .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-FDvXW3ZpPWDH4mGK .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-FDvXW3ZpPWDH4mGK .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-FDvXW3ZpPWDH4mGK .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-FDvXW3ZpPWDH4mGK .marker{fill:#333333;stroke:#333333;}#mermaid-svg-FDvXW3ZpPWDH4mGK .marker.cross{stroke:#333333;}#mermaid-svg-FDvXW3ZpPWDH4mGK svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-FDvXW3ZpPWDH4mGK p{margin:0;}#mermaid-svg-FDvXW3ZpPWDH4mGK .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-FDvXW3ZpPWDH4mGK .cluster-label text{fill:#333;}#mermaid-svg-FDvXW3ZpPWDH4mGK .cluster-label span{color:#333;}#mermaid-svg-FDvXW3ZpPWDH4mGK .cluster-label span p{background-color:transparent;}#mermaid-svg-FDvXW3ZpPWDH4mGK .label text,#mermaid-svg-FDvXW3ZpPWDH4mGK span{fill:#333;color:#333;}#mermaid-svg-FDvXW3ZpPWDH4mGK .node rect,#mermaid-svg-FDvXW3ZpPWDH4mGK .node circle,#mermaid-svg-FDvXW3ZpPWDH4mGK .node ellipse,#mermaid-svg-FDvXW3ZpPWDH4mGK .node polygon,#mermaid-svg-FDvXW3ZpPWDH4mGK .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-FDvXW3ZpPWDH4mGK .rough-node .label text,#mermaid-svg-FDvXW3ZpPWDH4mGK .node .label text,#mermaid-svg-FDvXW3ZpPWDH4mGK .image-shape .label,#mermaid-svg-FDvXW3ZpPWDH4mGK .icon-shape .label{text-anchor:middle;}#mermaid-svg-FDvXW3ZpPWDH4mGK .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-FDvXW3ZpPWDH4mGK .rough-node .label,#mermaid-svg-FDvXW3ZpPWDH4mGK .node .label,#mermaid-svg-FDvXW3ZpPWDH4mGK .image-shape .label,#mermaid-svg-FDvXW3ZpPWDH4mGK .icon-shape .label{text-align:center;}#mermaid-svg-FDvXW3ZpPWDH4mGK .node.clickable{cursor:pointer;}#mermaid-svg-FDvXW3ZpPWDH4mGK .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-FDvXW3ZpPWDH4mGK .arrowheadPath{fill:#333333;}#mermaid-svg-FDvXW3ZpPWDH4mGK .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-FDvXW3ZpPWDH4mGK .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-FDvXW3ZpPWDH4mGK .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-FDvXW3ZpPWDH4mGK .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-FDvXW3ZpPWDH4mGK .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-FDvXW3ZpPWDH4mGK .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-FDvXW3ZpPWDH4mGK .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-FDvXW3ZpPWDH4mGK .cluster text{fill:#333;}#mermaid-svg-FDvXW3ZpPWDH4mGK .cluster span{color:#333;}#mermaid-svg-FDvXW3ZpPWDH4mGK div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-FDvXW3ZpPWDH4mGK .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-FDvXW3ZpPWDH4mGK rect.text{fill:none;stroke-width:0;}#mermaid-svg-FDvXW3ZpPWDH4mGK .icon-shape,#mermaid-svg-FDvXW3ZpPWDH4mGK .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-FDvXW3ZpPWDH4mGK .icon-shape p,#mermaid-svg-FDvXW3ZpPWDH4mGK .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-FDvXW3ZpPWDH4mGK .icon-shape rect,#mermaid-svg-FDvXW3ZpPWDH4mGK .image-shape rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-FDvXW3ZpPWDH4mGK .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-FDvXW3ZpPWDH4mGK .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-FDvXW3ZpPWDH4mGK :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
Flink CDC 链路 (轻、快、直连)
Flink CDC Source
Upsert
MySQL
Flink Cluster (计算引擎)
Elasticsearch
传统链路 (重、慢、繁琐)
Binlog
MySQL
Canal/Maxwell
Kafka (消息队列)
Java Consumer / Logstash
Elasticsearch
🛠️ 二、 核心原理:它是如何做到“无锁”的?
早期的 CDC 工具在读取全量数据时,需要对表加全局锁 (Global Lock),这会阻塞线上业务的写入,DBA 根本不敢让你在白天跑。
Flink CDC 2.0+ 引入了 Netfix DBLog 算法(无锁算法):
💻 三、 实战:5 分钟搭建同步任务 (Flink SQL)
我们使用最简单的 Flink SQL 来演示。
1. 环境准备
- Flink 1.14+ (建议 1.17+)
- MySQL 5.7/8.0 (开启 Binlog, binlog_format=ROW)
- Elasticsearch 7.x/8.x
2. 依赖 JAR 包
将以下 jar 包放入 Flink 的 lib 目录:
- flink-sql-connector-mysql-cdc-2.x.jar
- flink-connector-elasticsearch-7_2.12.jar
3. 编写 Flink SQL
启动 sql-client.sh,执行以下语句:
Step 1: 定义 MySQL 来源表 (Source)
— 创建 MySQL CDC 表
CREATE TABLE mysql_users (
id INT,
name STRING,
age INT,
update_time TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = '192.168.1.100',
'port' = '3306',
'username' = 'root',
'password' = '123456',
'database-name' = 'mydb',
'table-name' = 'users',
— 极其重要:服务器时区设置
'server-time-zone' = 'Asia/Shanghai'
);
Step 2: 定义 Elasticsearch 目标表 (Sink)
— 创建 ES 结果表
CREATE TABLE es_users (
id INT,
name STRING,
age INT,
update_time TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'elasticsearch-7',
'hosts' = 'http://192.168.1.200:9200',
'index' = 'users_idx',
— 开启 Upsert 模式,支持更新和删除
'sink.bulk-flush.max-actions' = '1'
);
Step 3: 启动同步任务
— 一句 SQL 搞定同步
INSERT INTO es_users
SELECT id, name, age, update_time FROM mysql_users;
⚡ 四、 避坑指南 (生产环境必看)
- Elasticsearch Sink 必须定义 PRIMARY KEY。
- 这样当 MySQL 中删除一条数据时,Flink 会自动向 ES 发送一条带有 _id 的 Delete 请求,保证数据一致性。
- Flink CDC 2.x 支持 CDAS (Create Database As Source) 语法,可以一行代码同步整个数据库几百张表,无需一张张写 SQL。
- 在全量同步阶段,如果表特别大(亿级),需要适当调大 Flink TaskManager 的内存,防止 OOM。
🎯 总结
使用 Flink CDC,我们不再需要维护复杂的 Kafka 集群,不再需要编写繁琐的 Java 代码。
仅仅几十行 SQL,就构建了一个低延迟、高可靠、不仅能 Insert 还能 Update/Delete 的实时同步链路。
这就是技术演进带来的红利。拒绝数据孤岛,从拥抱 CDC 开始!
Next Step:
快去检查一下你的 MySQL Binlog 开没开,然后去测试环境跑通你的第一个 CDC 任务吧!




