简介:本资源是面向大数据开发工程师的Flink CDC 3.0实战指南,聚焦MySQL到Doris的实时数据同步场景,解决流式ETL中变更捕获、低延迟同步与端到端一致性等核心问题。内容由尚硅谷研究院编写,覆盖CDC原理对比(基于查询 vs Binlog)、flink-cdc-connectors组件机制、DataStream与Flink SQL双路径实操,以及MySQL→Doris全链路部署——含Binlog开启、检查点配置、多库表路由、环境搭建与任务验证等关键细节。资源为1个145KB的docx文档,结构清晰,含完整代码示例、命令行操作步骤及配置说明,便于快速复现和工程落地。目前已有1788人学习下载,适合具备Flink基础、正开展实时数仓建设或数据集成项目的中高级研发人员系统掌握Flink CDC生产级应用能力。
1. Flink CDC 3.0 不是“又一个同步工具”:它让 MySQL 到 Doris 的实时链路第一次真正摆脱双写、避免反查、绕过 Binlog 解析黑匣子,落地工业级秒级延迟
你有没有遇到过这样的现场:业务刚上线,运营要查“过去5分钟新注册用户地域分布”,BI 系统刷半天——数据还在 Kafka 里排队;运维半夜被告警叫醒,发现 Doris 表里某张订单明细表比 MySQL 少了 23 条记录,但日志里既没报错也没重试痕迹;更常见的是,开发改了个 MySQL 字段类型,Flink 作业直接挂掉,重启后全量重同步,Doris 里出现重复 ID 或空值……这些不是偶发故障,而是旧有同步范式(如 Canal + 自定义消费 + 手写 SQL 写入)在真实生产环境中的必然代价。Flink CDC 3.0 的核心价值,恰恰在于把“变更捕获”和“流式计算引擎”深度耦合:它不再把 Binlog 当作原始日志流来解析转发,而是用 Flink 原生状态管理做 checkpointed 解析、用 Watermark 控制乱序、用 Exactly-Once 语义保障端到端一致性——这意味着,当 MySQL 一条UPDATE order_status=2 WHERE id=1001落库,Doris 对应行在 800ms 内完成原子更新,且中间不依赖任何中间存储、不触发额外查询、不产生脏读幻读。它适合正在构建实时数仓、需要强一致 OLAP 查询、且 MySQL 版本 ≥5.7 / 8.0(含 GTID)、Doris 部署 ≥2.0.0 的中大型团队。如果你还在用 rsync 模拟实时、用定时任务拉全量、或靠业务双写保一致性——这篇就是你该停下手头工作立刻看的落地笔记。
2. 从零启动:Flink CDC 3.0 + MySQL + Doris 三组件最小可行链路搭建
Flink CDC 3.0 不是独立服务,它是 Flink 生态内建的 Source Connector。这意味着你不需要单独部署 Canal Server、不用维护 ZooKeeper 协调、也不用写 Java Consumer。但正因如此,它的正确启动极度依赖三者的版本兼容性、网络连通性与权限配置。本章带你用最简路径跑通端到端:MySQL 开启 Binlog → Flink 作业提交 → Doris 表自动创建并写入。所有操作均在 Linux x86_64 环境下验证(CentOS 7.9 / Ubuntu 22.04),不依赖 Docker 或 Kubernetes,纯裸机可复现。
2.1 MySQL 侧:必须显式启用的 5 项配置(缺一不可)
Flink CDC 3.0 依赖 MySQL 的 Binlog event 进行增量捕获,但默认安装的 MySQL 往往只开启基础 Binlog,缺少事务标识、行格式支持和 GTID,这会导致 CDC 作业无法识别 DML 类型、无法处理大事务、甚至直接抛Unsupported binlog format。以下配置必须写入/etc/my.cnf的[mysqld]段并重启 MySQL:
# 必须开启 Binlog(基础) log-bin=mysql-bin # 必须设为 ROW 格式(Flink CDC 仅支持 ROW) binlog-format=ROW # 必须开启 GTID(CDC 3.0 强依赖 GTID 定位位点,否则无法断点续传) gtid-mode=ON enforce-gtid-consistency=ON # 必须设置 server-id(集群唯一,单机可设为 1) server-id=1 # 可选但强烈建议:增大 Binlog 缓存,避免大事务截断 binlog-row-image=FULL提示:修改后执行
sudo systemctl restart mysqld,然后登录 MySQL 验证:SHOW VARIABLES LIKE 'log_bin'; -- 应返回 ON SHOW VARIABLES LIKE 'binlog_format'; -- 应返回 ROW SHOW VARIABLES LIKE 'gtid_mode'; -- 应返回 ON SHOW MASTER STATUS; -- 确认有 File 和 Position 输出
2.2 创建专用同步账号并授权(最小权限原则)
Flink CDC 作业以数据库用户身份连接 MySQL,该用户必须拥有SELECT,RELOAD,SHOW DATABASES,REPLICATION SLAVE,REPLICATION CLIENT权限。严禁使用 root 或应用账号——这是线上事故高发点。执行以下 SQL(替换cdc_user和your_password):
CREATE USER 'cdc_user'@'%' IDENTIFIED BY 'your_password'; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'cdc_user'@'%'; FLUSH PRIVILEGES;参数说明:
'cdc_user'@'%'允许从任意 IP 连接(生产环境请限定为 Flink 集群 IP 段,如'cdc_user'@'192.168.10.0/24')REPLICATION SLAVE是 CDC 读取 Binlog 的底层权限,缺失则报Access denied; you need (at least one of) the SUPER or SLAVE privilege(s)FLUSH PRIVILEGES不可省略,否则权限不生效
2.3 Doris 侧:启用 Broker Load 并创建目标数据库
Flink CDC 3.0 官方推荐通过 Doris 的Stream Load接口写入(HTTP POST),但实测中,当单条消息体积 >1MB 或并发写入量 >500 QPS 时,Stream Load 易触发413 Request Entity Too Large或503 Service Unavailable。因此,我们采用更稳定的Broker Load方式:Flink 将变更数据写入 HDFS/S3/OSS,Doris 通过 Broker 读取并导入。本例使用本地文件系统模拟(生产环境请替换为 HDFS):
-- 1. 创建数据库(CDC 作业会自动建表,但库需手动创建) CREATE DATABASE IF NOT EXISTS cdc_demo; -- 2. 创建 Broker(指向本地文件系统,Doris 2.0+ 默认内置 LocalFileSystemBroker) CREATE EXTERNAL CATALOG broker_catalog PROPERTIES ( "type"="hms", "hive.metastore.type"="local", "hive.metastore.uris"="thrift://localhost:9083" ); -- ⚠️ 注意:若未启用 Hive Metastore,请跳过此步,直接用 Stream Load(见 2.4) -- 3. 创建目标表(Flink CDC 可自动建表,但建议提前建好以控制分区、副本数) CREATE TABLE IF NOT EXISTS cdc_demo.orders ( id BIGINT NOT NULL COMMENT "订单ID", user_id BIGINT COMMENT "用户ID", amount DECIMAL(10,2) COMMENT "金额", status TINYINT COMMENT "状态:0-待支付,1-已支付,2-已完成", create_time DATETIME COMMENT "创建时间", update_time DATETIME COMMENT "更新时间" ) DUPLICATE KEY(id) DISTRIBUTED BY HASH(id) BUCKETS 10 PROPERTIES( "replication_allocation" = "tag.location.default: 3" );关键点说明:
DUPLICATE KEY(id)是 Doris 的明细模型,适合 CDC 场景(支持 UPSERT)DISTRIBUTED BY HASH(id)确保相同订单 ID 路由到同一 BE 节点,避免分布式 UPDATE 冲突replication_allocation设置副本数为 3,保障高可用(单节点测试可改为"tag.location.default: 1")
2.4 Flink 集群准备:JAR 包、配置与启动命令
Flink CDC 3.0 是 Flink 1.17+ 的原生功能,无需额外插件 JAR。但 Doris Sink 需要官方提供的flink-doris-connector。按以下顺序操作:
下载 Flink 1.17.2 二进制包(官网
flink-1.17.2-bin-scala_2.12.tgz),解压至/opt/flink下载 Doris Connector:从 Apache Doris 官网 GitHub Release 页面获取
flink-doris-connector-1.17_2.12-2.0.0.jar(注意 Scala 和 Flink 版本匹配),放入/opt/flink/lib/配置
flink-conf.yaml(关键项):# 必须启用 checkpoint,CDC 依赖其保存 Binlog 位点 state.backend: filesystem state.checkpoints.dir: file:///opt/flink/checkpoints execution.checkpointing.interval: 10s execution.checkpointing.mode: EXACTLY_ONCE # 提高并发容忍度(根据 CPU 核数调整) parallelism.default: 4 # 关键:禁用 operator chaining,避免 CDC Source 和 Sink 绑定导致背压传导 pipeline.operator-chaining: false启动 Flink Standalone 集群(单机测试):
cd /opt/flink ./bin/start-cluster.sh # 验证:curl http://localhost:8081/v1/jobs | jq '.jobs[] | select(.status=="RUNNING")'
3. 构建 CDC Pipeline:SQL Client 一键提交作业(无 Java 编码)
Flink CDC 3.0 最大进化是支持纯 SQL 定义整个 ETL 流程。你不再需要写 Java/Scala 代码、不再需要编译打包、不再需要管理 UDF。所有逻辑——从 MySQL 表发现、字段映射、变更事件过滤,到 Doris 写入策略——全部用标准 SQL 描述。本节提供可直接粘贴运行的完整 SQL 脚本,并逐行解释其作用域与风险点。
3.1 创建 MySQL CDC Source Table(自动发现表结构)
-- 在 Flink SQL Client 中执行(启动方式:./bin/sql-client.sh embedded) CREATE TABLE mysql_orders ( id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status TINYINT, create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (id) NOT NULL ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = '192.168.1.100', -- 替换为你的 MySQL IP 'port' = '3306', 'username' = 'cdc_user', 'password' = 'your_password', 'database-name' = 'shop_db', 'table-name' = 'orders', 'scan.startup.mode' = 'initial', -- initial=全量+增量;latest-offset=仅增量 'server-time-zone' = 'Asia/Shanghai', 'connect.timeout' = '30s', 'heartbeat.interval' = '30s' );参数深挖:
'scan.startup.mode' = 'initial':首次运行时先读全量快照(SELECT * FROM orders),再追 Binlog。生产环境首次上线必须用此模式,否则 Doris 数据为空。'server-time-zone':必须与 MySQLtime_zone一致(SELECT @@time_zone;查看),否则 DATETIME 字段时区错乱。'heartbeat.interval':CDC 会定期向 MySQL 发送心跳 SQL(SELECT 1),避免连接超时断开。设为30s是经验值,低于wait_timeout(MySQL 默认 28800s)即可。
3.2 创建 Doris Sink Table(声明式写入目标)
CREATE TABLE doris_orders ( id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status TINYINT, create_time STRING, -- Doris 不支持 TIMESTAMP 直接写入,需转 STRING update_time STRING ) WITH ( 'connector' = 'doris', 'fenodes' = '192.168.1.101:8030', -- Doris FE 地址(多个用逗号分隔) 'table.identifier' = 'cdc_demo.orders', 'username' = 'root', 'password' = '', 'sink.batch.size' = '1000', 'sink.max-retries' = '3', 'sink.enable-delete' = 'true', -- 启用 DELETE 事件同步(对应 MySQL DELETE) 'sink.enable-upsert' = 'true', -- 启用 UPSERT(对应 MySQL UPDATE) 'sink.properties.format' = 'json', 'sink.properties.strip_outer_array' = 'true' );关键参数说明:
'sink.enable-delete' = 'true':若 MySQL 表有 DELETE 操作,必须开启,否则删除事件被忽略。'sink.enable-upsert' = 'true':Doris 的DUPLICATE KEY模型依赖此参数实现 UPDATE 语义(主键存在则覆盖)。'sink.properties.format' = 'json':Doris Stream Load 接口要求 JSON 格式,Flink Connector 自动序列化。'create_time STRING':Flink 的TIMESTAMP(3)类型无法直接映射 DorisDATETIME,必须转STRING,Doris 会自动解析(前提是格式为yyyy-MM-dd HH:mm:ss.SSS)。
3.3 执行 INSERT INTO:启动实时同步管道
-- 将 MySQL 表变更实时写入 Doris 表 INSERT INTO doris_orders SELECT id, user_id, amount, status, DATE_FORMAT(create_time, 'yyyy-MM-dd HH:mm:ss.SSS') AS create_time, DATE_FORMAT(update_time, 'yyyy-MM-dd HH:mm:ss.SSS') AS update_time FROM mysql_orders;执行后验证:
- 查看 Flink Web UI(http://localhost:8081)→ Jobs → 该作业状态应为
RUNNING- 在 MySQL 执行
INSERT INTO orders VALUES (1001, 101, 99.99, 0, NOW(), NOW());- 立即查 Doris:
SELECT * FROM cdc_demo.orders WHERE id=1001;—— 数据应在 1~3 秒内出现- 再执行
UPDATE orders SET status=1, update_time=NOW() WHERE id=1001;,Doris 表中status应更新,update_time变为新值
4. 避坑指南:Flink CDC 3.0 + Doris 同步中 5 个血泪级翻车现场
Flink CDC 3.0 文档写得漂亮,但真实环境里,80% 的失败不是因为不会写 SQL,而是卡在环境细节、权限边界或隐式转换上。以下是我在线上踩过的坑,每一条都附带现象、根因和可立即执行的修复命令。
4.1 现象:Flink 作业启动后立即失败,日志报java.lang.ClassNotFoundException: com.ververica.cdc.connectors.mysql.table.MySqlTableSourceFactory
原因:Flink CDC 3.0 的 MySQL Connector 已合并进 Flink 主干,但部分发行版(如某些云厂商定制版)仍需手动添加 JAR。Flink 1.17.2 官方包默认不包含 CDC connector,必须显式引入。
解决:
# 下载 flink-sql-connector-mysql-cdc-3.0.0.jar(注意版本匹配 Flink 1.17) wget https://repo1.maven.org/maven2/com/ververica/flink-sql-connector-mysql-cdc/3.0.0/flink-sql-connector-mysql-cdc-3.0.0.jar cp flink-sql-connector-mysql-cdc-3.0.0.jar /opt/flink/lib/ # 重启 Flink 集群 ./bin/stop-cluster.sh && ./bin/start-cluster.sh4.2 现象:作业 RUNNING,但 Doris 表无任何数据,Flink Web UI 中 Source 的numRecordsInPerSecond为 0
原因:MySQL 用户缺少REPLICATION CLIENT权限,或server-id未设置。CDC 无法获取 Binlog 文件列表,只能空转。
解决:
-- 在 MySQL 中重新授权并确认 server-id GRANT REPLICATION CLIENT ON *.* TO 'cdc_user'@'%'; FLUSH PRIVILEGES; -- 检查 server-id SHOW VARIABLES LIKE 'server_id'; -- 若为 0,修改 my.cnf 并重启 mysqld4.3 现象:Doris 表中出现大量NULL值,尤其create_time和update_time字段全为 NULL
原因:MySQL 表中create_time定义为DATETIME DEFAULT CURRENT_TIMESTAMP,但 Flink CDC 读取时未触发默认值填充,且TIMESTAMP类型在 Flink 中解析为NULL。
解决:
-- 修改 MySQL 表,显式设置 DEFAULT 值并允许 NULL(CDC 能正确读取) ALTER TABLE orders MODIFY COLUMN create_time DATETIME DEFAULT CURRENT_TIMESTAMP, MODIFY COLUMN update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP; -- 或在 Flink SQL 中用 COALESCE 处理(推荐) SELECT id, user_id, amount, status, COALESCE(DATE_FORMAT(create_time, 'yyyy-MM-dd HH:mm:ss.SSS'), '1970-01-01 00:00:00.000') AS create_time, COALESCE(DATE_FORMAT(update_time, 'yyyy-MM-dd HH:mm:ss.SSS'), '1970-01-01 00:00:00.000') AS update_time FROM mysql_orders;4.4 现象:MySQL 执行DELETE FROM orders WHERE id=1001;后,Doris 表中该行未消失,反而status变为NULL
原因:Doris Sink 的'sink.enable-delete' = 'true'未生效,或 Doris 表未启用DELETE语义(DUPLICATE KEY模型需配合enable-delete才能物理删除)。
解决:
-- 确认 Doris 表模型为 DUPLICATE KEY(非 AGGREGATE 或 UNIQUE) SHOW CREATE TABLE cdc_demo.orders; -- 若非 DUPLICATE KEY,重建表 DROP TABLE cdc_demo.orders; CREATE TABLE cdc_demo.orders (...) DUPLICATE KEY(id) ... ; -- 重启 Flink 作业,确保 Sink 参数含 'sink.enable-delete' = 'true'4.5 现象:Flink 作业运行数小时后突然 Failover,日志报Caused by: com.mysql.cj.jdbc.exceptions.CommunicationsException: Communications link failure
原因:MySQLwait_timeout(默认 28800 秒 = 8 小时)超时,连接被服务端主动关闭,而 CDC 未及时重连。
解决:
-- 在 MySQL 中延长超时(推荐 24 小时) SET GLOBAL wait_timeout = 86400; SET GLOBAL interactive_timeout = 86400; -- 同时在 Flink CDC WITH 参数中增加重试配置 'connection.pool.size' = '10', 'connect.timeout' = '60s', 'scan.snapshot.fetch.size' = '1024'5. 生产级加固:从“能跑通”到“扛得住”的 4 个必调参数与 1 个监控闭环
跑通 demo 只是起点,真正的挑战在 7×24 小时稳定运行。我负责的某电商实时订单链路,峰值 QPS 12000,日增数据 8TB,曾因一个参数未调,导致连续 3 天凌晨 2 点准时断连。以下是我沉淀出的生产环境黄金配置组合,以及如何用 Prometheus + Grafana 构建可观测性闭环。
5.1 Flink CDC Source 侧:3 个影响吞吐与稳定性的核心参数
| 参数名 | 推荐值 | 作用说明 | 调优逻辑 |
|---|---|---|---|
scan.snapshot.fetch.size | 2048 | 全量阶段每次 JDBC FETCH 的行数 | 默认 1024,太小导致网络往返多;太大易 OOM。2048 是 16GB Flink TaskManager 的安全值 |
debezium.database.history.provided | true | 是否启用 Debezium 内置 history store | 设为true可避免创建schema_history表,减少 MySQL 权限依赖 |
server-id | 5400-5499(范围) | MySQL 复制的 server-id,必须唯一且为范围 | 单机部署设5400即可;集群部署每个 Flink TaskManager 分配不同 ID,避免 Binlog 冲突 |
实操命令(修改
CREATE TABLE mysql_orders的 WITH 子句):'scan.snapshot.fetch.size' = '2048', 'debezium.database.history.provided' = 'true', 'server-id' = '5400'
5.2 Doris Sink 侧:1 个决定写入成功率的关键参数
'sink.buffer-flush.max-bytes'是 Doris Connector 的内存缓冲阈值,默认3mb。当单条变更消息较大(如含长文本字段),或网络抖动导致批量发送失败时,缓冲区满会触发强制 flush,但若此时 Doris FE 不可用,整个 batch 会丢弃且无重试——这就是“数据静默丢失”的根源。
生产建议值:10485760(10MB),并配合'sink.max-retries' = '10':
'sink.buffer-flush.max-bytes' = '10485760', 'sink.max-retries' = '10', 'sink.buffer-flush.interval-ms' = '30000' -- 30秒强制 flush,防内存积压5.3 构建监控闭环:用 Prometheus 抓取 Flink 指标并告警
Flink 自带 Metrics Reporter,只需配置即可暴露 JMX/REST 指标。我们重点关注三个黄金指标:
| 指标名 | Prometheus Query | 告警阈值 | 说明 |
|---|---|---|---|
flink_taskmanager_job_task_operator_numRecordsInPerSecond | rate(flink_taskmanager_job_task_operator_numRecordsInPerSecond{job_name=~".*cdc.*"}[5m]) < 10 | <10 条/秒持续 5 分钟 | 源端无数据流入,检查 MySQL 连接或 Binlog |
flink_taskmanager_job_task_operator_numRecordsOutPerSecond | rate(flink_taskmanager_job_task_operator_numRecordsOutPerSecond{job_name=~".*cdc.*"}[5m]) < 5 | <5 条/秒持续 5 分钟 | Sink 写入阻塞,检查 Doris FE/Broker 状态 |
flink_taskmanager_job_task_operator_stateSize | flink_taskmanager_job_task_operator_stateSize{job_name=~".*cdc.*"} > 1073741824 | >1GB | State 过大,可能 checkpoint 失败,需调大state.checkpoints.dir磁盘空间 |
配置步骤:
- 修改
/opt/flink/conf/flink-conf.yaml,添加:metrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter metrics.reporter.prom.port: 9249- 启动 Prometheus,配置
scrape_configs抓取http://flink-taskmanager:9249/metrics- 在 Grafana 导入 Flink Dashboard(ID: 12658),添加上述告警规则
5.4 一个我坚持了 3 年的习惯:每日凌晨 3 点自动校验数据一致性
再稳的链路也需要“后悔药”。我写了一个 Python 脚本,每天固定时间执行:
#!/usr/bin/env python3 # cdc_consistency_check.py import pymysql from pyspark.sql import SparkSession # 1. 从 MySQL 抽样 1000 条最新订单(按 update_time 倒序) mysql_conn = pymysql.connect(host='192.168.1.100', user='reader', password='pwd', db='shop_db') with mysql_conn.cursor() as cur: cur.execute("SELECT id, status, update_time FROM orders ORDER BY update_time DESC LIMIT 1000") mysql_data = cur.fetchall() # 2. 从 Doris 查询同一批 ID spark = SparkSession.builder \ .appName("Doris Check") \ .config("spark.doris.fenodes", "192.168.1.101:8030") \ .getOrCreate() doris_df = spark.read.format("doris") \ .option("table.identifier", "cdc_demo.orders") \ .option("filter", f"id in ({','.join([str(x[0]) for x in mysql_data])})") \ .load() doris_data = [(row.id, row.status, row.update_time) for row in doris_df.collect()] # 3. 对比并发送企业微信告警 diff = [m for m in mysql_data if m not in doris_data] if diff: send_wechat_alert(f"CDC 不一致!发现 {len(diff)} 条差异:{diff[:3]}")为什么有效:
- 抽样而非全表对比,耗时 <30 秒
- 选
update_time最新的数据,覆盖高频变更场景- 差异定位到具体 ID 和字段,运维可秒级介入
- 企业微信告警带链接直达 Flink Web UI 和 Doris Query Log
这套机制上线后,数据不一致问题平均响应时间从 4.2 小时缩短到 7 分钟。它不替代监控,而是给自动化兜底加一道人工可验证的保险栓。
希望帮到你。
本文还有配套的精品资源,点击获取