欢迎光临
我们一直在努力

拒绝数据孤岛!基于 Flink CDC 实现 MySQL 到 ElasticSearch 的“毫秒级”数据同步

🌊 前言:为什么传统的同步方案“不香”了?

在 Flink CDC 普及之前,要实现 MySQL 到 ES 的实时同步,架构师通常会画出这样一张图:
MySQL -> Canal/Debezium -> Kafka -> Logstash/Java应用 -> Elasticsearch。

痛点非常明显:

  • 链路太长:引入了 Kafka 中间件,维护成本飙升。
  • 数据一致性难保证:一旦发生故障,如何保证数据不丢、不重?
  • 全量+增量割裂:你需要先写个脚本导全量,再开启 Canal 追增量,切换瞬间容易丢数据。
  • 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 算法(无锁算法):

  • 分片读取 (Chunk Splitting):将大表切分成无数个小 Chunk。
  • 动态监测:在读取 Chunk 的过程中,利用 Binlog 修正读取期间发生的数据变更。
  • 结果:全程无锁,对线上业务几乎零影响,且支持断点续传。

  • 💻 三、 实战: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;


    ⚡ 四、 避坑指南 (生产环境必看)

  • 关于删除 (Delete) 操作:
    • 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 任务吧!

    赞(0)
    未经允许不得转载:171主机测评 » 拒绝数据孤岛!基于 Flink CDC 实现 MySQL 到 ElasticSearch 的“毫秒级”数据同步
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址