Flink CDC + ClickHouse:MySQL实时同步至数仓的完整实践

发布时间:2026/9/26 17:27:40
Flink CDC + ClickHouse:MySQL实时同步至数仓的完整实践
上个月帮一个做电商业务的团队搭实时订单看板他们在MySQL里跑了两年的线上库单表已经超过千万行几个常用的聚合SQL要三秒以上运营那边要的“实时GMV”只能靠半小时一次的定时任务勉强将就。聊了一圈之后方案很自然地落到“实时链路”这三个字上把MySQL的变更数据近乎实时地搬进ClickHouse让分析查询和线上事务彻底分家。选型时的候选方案不少——定时批量同步、业务代码双写、基于日志的CDC最终我采用的是目前数据工程圈子里最主流的做法用Flink CDC这类SQL化工具以binlog为数据源把MySQL里的增删改事件流式投递到ClickHouse。这篇文章就是这次从选型、联调到上线踩坑的完整记录适合正在搭实时数仓、或准备把MySQL分析业务迁到ClickHouse的团队参考。1. 业务问题与链路设计目标为什么非要把MySQL数据搬进ClickHouse1.1 场景与痛点线上交易系统一般都以MySQL为存储核心订单、支付、用户这些关键业务表全在里面。业务侧的需求是点查和短事务这正好是MySQL的强项但分析侧的需求是另一码事——大范围扫描、多表关联、复杂聚合这类查询在B树索引和行存储面前非常吃力几个大表join一次CPU和IO就吃满了。更麻烦的是分析查询会在高峰期跟线上事务抢资源。白天业务量上来之后稍微跑一个稍大的报表线上库的慢查询数量和连接数就开始飙运气不好还会拖慢支付确认、订单创建这种核心接口。这不是加个索引能解决的因为报表SQL是动态多变的索引根本追不上。把数据搬到ClickHouse就是很自然的选择。ClickHouse是列式存储数据压缩比高单表扫描速度快得离谱尤其适合GROUP BY、聚合、漏斗分析这类场景。但企业里让ClickHouse直接当业务主库也不现实它的点查能力弱、事务支持有限、真正的UPDATE和DELETE做起来很别扭。所以实践中几乎都采用“MySQL当源ClickHouse当分析侧副本”的架构把两边擅长的东西分开。1.2 链路设计的三条硬约束和那团队对需求的时候我给他们列了三条约束后续所有技术选型都是围绕这三条来的延迟从MySQL提交到ClickHouse可见目标秒级到分钟级至少要比定时批处理高一个量级。他们要看的实时GMV、实时转化率延迟超过5分钟就没有参考价值了。准确性不能丢数据也不能重复至少保证最终一致。对账单、库存这类数据错一条都可能引发业务投诉。低成本不改造业务代码尽量少引入额外组件。双写方案被否掉的主要原因就是侵入性太强而且业务团队根本不愿意在每个写接口里再塞一段ClickHouse的调用。还有一条很容易被忽略链路必须把全量和增量统一起来。首次上线时要把MySQL里的历史数据灌进ClickHouse之后要持续同步增量的增删改。全量和增量之间的衔接如果断掉就会出现老数据缺失或重复的隐患。很多方案表面上能跑实际就是在这一环露馅的。2. 选型推演为什么是CDC以及为什么用Flink CDC2.1 三种同步路线的现实对比我在动手前把市面上常见的同步方案过了一遍整理成对比就很直观了方案延迟业务侵入性数据一致性实现复杂度定时批量同步DataX/Sqoop分钟到小时级低差非实时低业务代码双写实时高要改所有写路径一般容易出偏差中CDC日志捕获binlog解析秒级低不改业务代码好有前后镜像中定时批量同步的做法最成熟但本质上是“每天拍一次全楼照片”调度时间为准数据必然不是实时的。双写方案听起来简单实际等于在每个业务接口里多塞一条通路MySQL写完ClickHouse没写、或者ClickHouse写失败了两边数据就对不上排查起来极其痛苦。CDC的思路完全不同。它不碰业务代码而是直接读取MySQL的binlog模拟一个从库把增量日志消费下来。每次INSERT、UPDATE、DELETE都变成了带前后镜像的流式事件比如更新前叫什么、更新后变成什么全都拿得到。用一个生活化的类比定时批量是每天拍照双写是每个房间改动都要自己多抄一份送到隔壁而CDC是直接装了管道系统水龙头一开流水自动到目标端。2.2 CDC实现的关键前提MySQL binlog的正确姿势CDC能工作的前提是MySQL的binlog配置正确。binlog有三种格式这里值得展开说一下STATEMENT记录原始SQL语句。优点是占用空间小但同一个SQL在不同节点执行的结果可能不一致比如用了NOW()、UUID()这类函数。ROW记录每一行变更前后的完整镜像。CDC拿到的是真实的数据变化不是SQL能精确识别哪一行变了。MIXED混合模式MySQL自己判断。CDC必须用ROW格式并且binlog_row_image要设为FULL这样才会记录所有列的前后值。另一个关键点是server_id必须唯一。每个binlog dump连接在MySQL看来都是一个从库如果跟现有主从复制的server_id重复MySQL会直接踢掉其中一个连接。MySQL侧最小配置如下[mysqld] server-id1 log-binmysql-bin binlog_formatROW binlog_row_imageFULL gtid_modeON enforce_gtid_consistencyON binlog_expire_logs_seconds604800binlog_expire_logs_seconds604800意思是binlog保留7天这个值要根据你的同步任务挂了多久能恢复来定太短容易丢位点太长会占磁盘。还需要一个同步专用账号权限给到最小CREATE USER cdc% IDENTIFIED BY cdc_password; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO cdc%; FLUSH PRIVILEGES;这里REPLICATION SLAVE和REPLICATION CLIENT是核心很多权限不全的案例表面现象是Flink任务连上了但一直不输出数据实际上就是这两个权限没给到位。2.3 为什么最终选了Flink CDC而不是Canal自研Canal是经典的binlog同步中间件Java系团队很熟悉但它本身不带计算框架数据推给下游之后怎么处理要自己写。更关键的是Canal拿全量历史数据的能力很弱通常是靠单独导一次MySQL快照再启动增量这正好踩中了我前面说的“全量和增量衔接难”的问题。Flink CDC是另一个路子。它内嵌了Debezium核心能直接对接Flink SQL和DataStream API最舒服的是全量阶段支持按主键切分chunk并行扫描也不需要全程持有全表锁。Flink CDC 3.0之后甚至提供了YAML Pipeline模式连代码都不用写纯配置就能搭一条同步链路。对于咱们这类场景Flink CDC还带一个杀手锏状态管理和Checkpoint。任务重启后能从最近一次Checkpoint恢复继续从对应binlog位点消费不用从头再全量扫一遍。这一点在生产环境实在太重要了后面我会细讲。3. 联调实录用Flink CDC 3.x Pipeline直通ClickHouse3.1 环境准备版本清单与前置配置我这次用的版本组合如下不一定是最新的但足够稳MySQL 8.0.x开启binlog并验证配置ClickHouse 23.8建好目标库表Flink 1.17 Flink CDC 3.0 standalone部署包JDK 11MySQL开启binlog之后先验证一下配置真的生效了别改完my.cnf没重启就以为完事了SHOW VARIABLES LIKE binlog_format; SHOW VARIABLES LIKE binlog_row_image;ClickHouse侧建库建表我习惯把同步表设计成ReplacingMergeTree引擎CREATE DATABASE IF NOT EXISTS app_db ON CLUSTER default; CREATE TABLE app_db.orders ( id UInt64, order_no String, user_id UInt64, amount Decimal64(10, 2), create_time DateTime64(3), _version UInt64, _sign Int8 ) ENGINE ReplacingMergeTree(_version) PRIMARY KEY (id) ORDER BY (id);这里解释一下_version和_sign的用途。MySQL里对同一条记录做多次UPDATE在ClickHouse这边其实会生成多行不可变数据ReplacingMergeTree后台合并时按主键去重靠_version字段判断哪一行才是最新的——值越大越新。_sign用来标记删除MySQL的DELETE事件落过来时我们把它翻译成_sign -1正常变更填1查询时过滤一下就能拿到逻辑上的最新状态。3.2 YAML Pipeline配置逐行拆解Flink CDC 3.0的YAML Pipeline模式确实省心整个链路就一个配置文件source: type: mysql hostname: 127.0.0.1 port: 3306 username: cdc password: cdc_password tables: app_db.orders server-id: 5400-5404 startup-mode: initial sink: type: clickhouse hostname: 127.0.0.1 port: 8123 username: default password: database: app_db table: orders batch-size: 5000 flush-interval: 1000 retry-times: 3 pipeline: name: mysql-to-clickhouse parallelism: 2几个字段值得单独说tables支持库表白名单多个表用逗号分隔也支持正则。这里先只跑app_db.orders验证通过后再加其他表。server-id我给了一个范围5400-5404不要只写单个数字。因为全量阶段并行读多个chunk时每个reader都需要一个独立的server-id去拉binlog范围够用才能并行。startup-mode: initial这个模式会在启动时先做全量快照再自动切到增量正好覆盖前面说的“全量和增量衔接”问题。如果只想从当前时刻开始同步可以选latest-offset。batch-size和flush-interval控制攒批写入ClickHouse的阈值。ClickHouse特别不喜欢单条写入攒到5000条或1秒后再批量flush写入效率能差出一个数量级。parallelismYAML模式下的全局并行度。单表同步2够用多表同步可以调高但要注意上游MySQL的server-id范围也要跟着扩大。有的Flink CDC版本对ClickHouse sink的配置字段命名可能略有差异以你实际使用的版本Release文档为准这个没关系结构是一样的。启动命令很简单./flink-cdc.sh mysql-to-clickhouse.yaml3.3 首次全量增量验证的完整过程启动之后我先在MySQL里执行一批INSERT、UPDATE、DELETE再跑到ClickHouse查结果INSERT INTO orders (id, order_no, user_id, amount, create_time) VALUES (1, ORD001, 101, 99.90, NOW()), (2, ORD002, 102, 199.00, NOW()); UPDATE orders SET amount 129.00 WHERE id 1; DELETE FROM orders WHERE id 2;然后查询ClickHouseSELECT * FROM app_db.orders FINAL;注意这里加了FINAL关键字。ReplacingMergeTree的合并是异步的刚写入完不一定会立即合并不加FINAL可能在联调时看到同一主键的多版本数据造成“同步丢了数据”的错觉。FINAL会让查询时强制合并结果准确但开销大生产环境大表不建议频繁用。验证结果表明id1的订单金额已经变成129.00id2的记录带上了删除标记。全量同步的部分我在启动前预置了一些历史数据任务启动后也完整出现在了ClickHouse里说明initial模式的全量增量衔接没有断。4. 用Flink SQL把同一条链路重写一遍这才是日常用法4.1 为什么还要强调SQLYAML Pipeline确实快但它解决的是“同步”问题如果要说清楚标题里的“基于SQL”得重点看Flink SQL这一层。SQL是数据工程师最熟悉的抽象用一句CREATE TABLE把MySQL表暴露成Flink的虚拟表再用一句INSERT INTO把数据写进ClickHouse这种体验比写DataStream API的代码直观太多。实际生产里我们往往不只是把数据搬运过去还要做过滤、字段改名、多表关联、数据清洗。这些活在YAML Pipeline里实现起来很别扭但在Flink SQL里就是几行SQL的事。这也是Flink CDC比Canal自研生态强的地方——CDC能力被SQL API封装掉了业务逻辑用SQL表达谁来都能看懂。4.2 Flink SQL建表与源端配置在Flink SQL Client里第一步是把jar包加载进去包括Flink CDC连接器和ClickHouse连接器ADD JAR /path/to/flink-sql-connector-mysql-cdc.jar; ADD JAR /path/to/flink-sql-connector-clickhouse.jar;然后建MySQL源表CREATE TABLE mysql_orders ( id BIGINT, order_no STRING, user_id BIGINT, amount DECIMAL(10, 2), create_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 127.0.0.1, port 3306, username cdc, password cdc_password, database-name app_db, table-name orders, server-id 5400-5404, scan.startup.mode initial );再建ClickHouse目标表。要说明的是ClickHouse目前不在Flink官方内置连接器清单里社区和官方维护的连接器在WITH字段上可能不统一下面是我用的那个版本的常见写法CREATE TABLE ch_orders ( id BIGINT, order_no STRING, user_id BIGINT, amount DECIMAL(10, 2), create_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector clickhouse, url clickhouse://127.0.0.1:8123, username default, password , database-name app_db, table-name orders, sink.batch-size 5000, sink.flush-interval 1000, sink.max-retries 3 );最后执行INSERT INTO ch_orders SELECT id, order_no, user_id, amount, create_time FROM mysql_orders;任务就起来了。这句INSERT背后Flink CDC会持续消费binlog把MySQL的增删改事件流式写入ClickHouse。4.3 SQL方式的进阶能力实时订单宽表纯粹做表复制YAML Pipeline就够了。SQL方式真正的价值体现在处理逻辑上。比如那家电商团队要做实时订单宽表需要把orders和users关联起来YAML Pipeline做不到Flink SQL只需要一句CREATE VIEW order_wide AS SELECT o.id AS order_id, o.order_no, u.user_name, o.amount, o.create_time FROM mysql_orders o LEFT JOIN mysql_users FOR SYSTEM_TIME AS OF o.proceed_time u ON o.user_id u.id;这里用到了Flink SQL的Temporal Join时态表关联也就是“关联的是用户表在订单发生时刻的快照”而不是用户表当前的最新状态。做实时宽表时这个细节特别重要如果不加FOR SYSTEM_TIME AS OF关联上的是“当前用户表”历史订单的用户名会被后面修改的数据污染。这背后的含义是使用SQL并打开CDC能力之后链路不只是搬运工具而是一个实时流计算管道可以在数据从MySQL到ClickHouse的过程中完成实时的JOIN、过滤、聚合这是我觉得“基于SQLCDC”最核心的价值。4.4 这种方式的局限有得必有失。SQL方式最大的短板是ClickHouse sink不提供事务性写入Flink无法做到严格意义上的端到端exactly-once实际上通常是at-least-once也就是可能重复但不会丢。要保住最终一致就得靠ClickHouse目标表的幂等去重。另外社区连接器在某个版本下可能有类型映射的坑用之前一定要先压测、再核对数据。链路简单时我推荐用YAML Pipeline一旦涉及加工逻辑还是老老实实回到Flink SQL。5. 链路能跑不等于数据对核心机制与一致性保证5.1 无锁全量增量增量快照到底做了什么很多第一次接触Flink CDC的人会好奇它做全量同步的时候MySQL不需要锁表吗答案是基本不需要。Flink CDC用的是增量快照机制先把目标表按主键切分成多个chunk多个task并行扫描这些chunk每个chunk读取时会记录当时的binlog位点全量片段读完后再消费这段时间内的binlog事件做补偿。以orders表为例假设有100万行按主键切成10个chunk10个reader并行扫每个chunk扫完之后交给一个补偿线程去拉对应区间的binlog增量。这个机制让全量阶段极快同时不阻塞线上业务写入。相比之下Canal要读全量历史数据就麻烦得多一般得靠别的工具先导一次快照再拼接增量中间很容易出现缝隙。5.2 server-id为什么要给一整个范围前面YAML配置里我写了server-id: 5400-5404这里解释得更细一些。Flink CDC的每个reader本质上是一个模拟从库要用一个server-id向MySQL发送binlog dump请求。当并行读多个chunk时多个reader如果共用同一个server-idMySQL那边会认为是从库连接冲突直接踢掉一个连接。分配一个范围让每个reader各用各的互不干扰。生产环境里还要注意和MySQL现有主从复制的server-id区分开。我可以负责任地说一半以上的“Flink CDC连上后同步卡住”案例都是server-id跟线上已有从库撞了。5.3 Checkpoint与重启恢复在Flink SQL和YAML Pipeline模式下CDC源会把binlog消费位点记到Flink的状态里状态会周期性写入Checkpoint。任务异常重启后Flink从最近一次Checkpoint恢复把binlog位点回退到那个时刻重新消费没处理完的事件从而做到“不丢”。实操中我给两个经验值Checkpoint间隔不要设太长默认1分钟在关键链路上有点粗我习惯设30秒恢复时少补一些数据状态后端生产环境用RocksDB别把状态全放堆内存否则大表全量阶段容易OOM。5.4 ClickHouse端的幂等ReplacingMergeTree的正确姿势ClickHouse的存储模型决定了同一行数据可以有多版本。每次INSERT都会生成新的不可变数据段part不会原地更新旧数据。ReplacingMergeTree在后台合并part时按主键保留版本最新的那行从而实现最终去重。两个实操细节一是版本字段的来源。最省事的是取binlog事件里的ts_ms毫秒时间戳但如果同一毫秒内同一行发生了多次变更就会因为版本相同而随机保留造成结果不确定。生产环境更稳妥的做法是给每条记录生成一个单调递增的序号比如用binlog的filenameposition拼一个seq作为_version。二是删除事件的处理。MySQL的DELETE到了Flink CDC会变成一个DELETE操作如果我们在Sink端直接把这一行忽略ClickHouse里的旧版本会一直留在表里无法体现“已删除”状态。这也是我在3.1建表时加_sign字段的原因DELETE时写一条_sign -1的记录查询时用sum(_sign)或last row by version来过滤已删除行。如果只用ReplacingMergeTree不管删除标记数据看起来“多了”日志跟源库对不上排查时会非常痛苦。6. 生产环境踩坑清单六个真实经验6.1 MySQL侧权限、server-id、binlog保留时长第一个坑是server-id冲突。我一个项目里同时起了两个Flink任务分别同步不同库结果两个任务轮流报“Got error reading packet from server”。排查下来发现两个任务的server-id段配置成了同一个MySQL把其中一个连接踢了。解决很简单给每个任务分配独立的server-id范围比如任务A用5400-5404任务B用5500-5504。第二个坑是binlog保留时间不够。有一次Flink任务挂了超过一天恢复了之后从Checkpoint重启发现binlog位点已经过期直接在日志里抛异常。这里没有捷径只能调大binlog_expire_logs_seconds或者把Checkpoint的生命周期调好保证任务恢复时间小于binlog的保留窗口。第三个坑是同步账号权限不足。前面给的权限清单里REPLICATION SLAVE和REPLICATION CLIENT缺一不可。曾经遇到一个环境只给了SELECT权限任务启动后日志看似正常实际上一个事件都读不出来。用SHOW GRANTS FOR cdc%;确认一下最快。6.2 ClickHouse侧part数量、合并风暴与批量写入ClickHouse的经典问题是写入过于碎片化。如果你把batch-size设成几十条每秒钟会产生大量小part后台merge线程根本来不及合并一段时间之后就会报too many parts错误。我这边优化后的参数是batch-size调到5000~10000flush-interval控制在1~2秒每个part至少有个几百KB。效果立竿见影part数量降了一个数量级。还有合并风暴的问题。ReplacingMergeTree的合并是异步的如果某一天数据量突然暴涨merge跟不上查询就会变慢。这时候不要慌先在system.merges里观察合并状态如果积压严重可以临时调大background_pool_size一般能缓过来。千万不要在生产大表上频繁执行OPTIMIZE TABLE ... FINAL那会让ClickHouse临时把所有part强制合并IO打满影响线上查询。6.3 类型与时区DECIMAL、DATETIME、String类型映射是同步链路最琐碎又最容易出问题的地方。MySQL的DECIMAL(10,2)对应ClickHouse要用Decimal64(10,2)或Decimal128精度一旦超出就会报数据溢出。MySQL的DATETIME是没有时区概念的本地时间TIMESTAMP则有时区转换Flink CDC读取时会根据table.local-time-zone参数统一处理如果这个参数跟ClickHouse服务器时区不一致同步过去的create_time会出现整小时偏移。我的建议是建表之前先构造几条已知时间戳的数据跑一遍链路做对照别等数据量大了才发现时间对不上。String字段在ClickHouse用String还是Nullable(String)也要商量好。不涉及空值的统一用String可空字段用Nullable(String)但Nullable在列存里有额外标记开销能避免就避免。实际项目中我见过因为两边类型宽泛导致写入报错的排查起来特别费劲。6.4 数据核对与监控告警链路搭完只是开始如何保证数据从MySQL到ClickHouse之后是对的才真正考验工程能力。我常用的三板斧总量对比每天定时跑一次SELECT COUNT(*), SUM(amount), MAX(create_time)跟MySQL源表对比几条大指标对不上就说明链路有问题。增量对账抽最近一小时的数据做一次MD5或哈希聚合对比能把细微的类型转换错误揪出来。延迟监控在ClickHouse侧执行SELECT now() - max(create_time) FROM app_db.orders;把结果接入告警一旦延迟超过30秒就报警。Flink侧也要盯重点看Checkpoint是否持续成功、taskmanager的忙闲状态、是否有restart。一个常见的隐蔽问题是内部反复重启但外部看着任务还活着日志一直刷反序列化异常这种多半要人工介入处理。6.5 DDL变更别指望全自动Flink CDC在部分较新版本里能捕获MySQL的ALTER TABLE操作但这不意味着ClickHouse会自动同步字段变更。两边的类型系统、存储结构差异太大生产上我从来不指望DDL自动透传。我的流程是MySQL要加字段时先在ClickHouse手动ALTER TABLE ... ADD COLUMN再同步更新Flink SQL里的表定义或者重启YAML Pipeline。新字段默认值在两边要配一致否则历史数据没值、新数据有值查询一跑就露馅。最后分享两个小体会这套链路跑稳定之后我的体会是实时链路真正的难点不在“搭起来”而在“持续稳定”。丢几条数据的根因通常不是工具本身而是某个字段类型、某个时区、某次DDL没对齐。做实时链路的本质是管理数据在时间线上的版本想清楚这一点排障的时候会少走很多弯路。另外一个习惯强烈建议保持把YAML配置和Flink SQL脚本都放到Git里管理。链路的全量重建、环境迁移、配置变更全部能追溯比拷来拷去的配置文件干净太多。我后来接手过几个团队的实时同步工程有没有版本管理排查问题的效率能差出好几倍。