
做数据同步这些年从最早的mysqldump全量拷贝到后来折腾Canal、binlog、Flink CDC踩过的坑能写一本小册子。今天把MySQL数据出海这件事彻底拆开聊聊——不是给你堆一堆工具名词而是把方案选型背后的逻辑、每条链路的真实优缺点、以及我实际落地时踩过的坑一次性讲透。所谓“数据出海”本质就是让MySQL里沉淀的数据流动起来进入更合适它的地方。可能是离线数仓做分析可能是实时大屏做指标展示也可能是像TDengine这样的时序数据库承载海量监控数据。这篇文章适合正在做数据同步选型、或者已经在同步路上被各种诡异问题折磨的读者无论你是架构师、后端开发还是DBA都能从中找到可直接抄作业的方案。1. 数据出海场景拆解为什么要把MySQL数据搬出去1.1 数据同步的典型业务场景先别急着选工具得先搞清楚“出海”到底要去哪。我接触过的项目里最典型的是下面几类第一类是离线分析场景。业务库跑着线上交易几千万行的订单表不可能让分析师直接在上面跑复杂聚合查询哪怕加了索引也会把主库拖垮。这时候需要把MySQL数据同步到数据仓库比如Hive、Doris、ClickHouse做T1甚至T0的分析。第二类是实时消费场景。比如大屏要实时显示订单量或者下游搜索服务要实时同步商品信息数据从MySQL产生到下游可用的延迟必须控制在秒级甚至毫秒级。第三类是异构存储场景。比如设备上报的时序数据、监控指标存在MySQL里不仅存储成本高查询性能也完全不对路这时候需要同步到TDengine这类时序数据库。第四类是数据分发场景。多个业务系统都需要同一份基础数据比如用户中心的数据要同步给订单系统、客服系统、推荐系统不能每个系统都直连主库。每一种场景对应的同步方案差别非常大离线场景可以接受小时级延迟实时场景要求秒级异构存储还得考虑表结构映射和类型转换。这决定了你不可能一套方案走天下。1.2 方案选型的四个关键维度选方案之前先学会画坐标轴。我通常用四个维度来评估一套同步方案同步延迟从MySQL数据变更到目标端可见能接受多少延迟分钟级还是秒级这直接淘汰一批方案。同步粒度是全量同步还是增量同步是否需要同步表结构变更DDL还是只需要同步数据变更DML一致性要求是否允许重复数据是否需要保证事务的原子性下游有没有幂等消费能力技术栈匹配团队熟悉Java还是Python愿不愿意引入额外组件比如Canal、Kafka还是希望尽量少依赖、少运维拿这四个维度去套很多选择自然就清晰了。比如你只需要每天晚上全量同步一次订单数据给数仓那搞一套CDC实时同步成本反而高得没必要反过来如果业务方要求实时看到用户下单那你用mysqldump凌晨三点导一次显然也是交不了差的。1.3 主流方案全景对比我把这些年见过、用过的方案整理成一张对比表方便你做选择题方案同步延迟增量支持部署复杂度适用场景核心痛点mysqldump 全量导出无实时性不支持极低小表备份、一次性迁移数据量大了导出时间长会锁表MySQL 主从复制秒级以内原生支持低MySQL到MySQL的同构同步只支持MySQL源和MySQL目标异构完全不行Canal Kafka秒级支持中高异构数据源、实时数仓需要维护Canal和Kafka两套组件Flink CDC秒级支持中实时数仓、异构存储学习成本偏高状态后端调优有门槛DataX分钟级离线低离线批量同步准实时能力弱调度靠自己搭业务双写实时天然支持无小规模/新项目侵入业务代码改造风险大这里多说一句双写看着最省事实际最坑。我见过一个团队在订单系统里加了双写逻辑结果下游存储抖了一下事务提交就失败连带主库写也回滚线上订单全挂。业务双写只能用在非核心链路核心数据还是得靠CDC老老实实抓binlog。2. 核心机制原理解析binlog是如何成为同步枢纽的2.1 binlog的三种格式与CDC的必然选择聊MySQL数据同步绕不开binlog。它本质是MySQL服务端记录的所有数据变更事件的日志文件你可以把它理解成数据库的“黑匣子”谁在什么时间改了哪张表的哪一行、改之前是什么值、改之后是什么值全都记录在案。但这三种格式差别很大STATEMENT格式记录的是SQL语句原文比如UPDATE t SET namex WHERE id1。优点是日志量小但缺点是语句执行依赖上下文比如用了NOW()函数、自增列从库回放时结果可能不一样。ROW格式记录的是每一行变更前后的具体值比如“id1这一行的name从y变成了x”。优点是绝对准确缺点是日志量比STATEMENT大很多。MIXED格式MySQL自己判断默认用STATEMENT遇到不确定函数时自动切换ROW。做数据同步必须用ROW格式。原因很简单下游需要知道具体是哪一行变了、旧值是什么、新值是什么才能在其他存储里做对应的更新。STATEMENT格式虽然日志小但对异构同步来说等于没有信息——你拿到了UPDATE t SET namex WHERE id1这句SQL却根本不知道在目标端要改哪一行。还有一个重要参数是binlog_row_imageFULL这个参数决定ROW格式下binlog是否记录所有列的“前镜像”和“后镜像”。默认值是FULL即记录整行所有列有些团队为了省日志量改成MINIMAL只记录被修改的列和能定位主键的列。省了空间但会带来一个大问题下游如果拿到binlog数据后想对比“到底改了哪些列”就会出现信息不全的情况。做同步建议保持FULL。2.2 从主从复制到CDCCanal/Flink CDC伪装成从库的秘密理解了binlog也就理解了Canal和Flink CDC的工作原理。MySQL原生主从复制是这样的主库开启binlog从库启动两个线程一个IO线程连上主库伪装成从库把主库的binlog拉过来写到本地relay log另一个SQL线程读relay log并回放完成数据同步。Canal的高明之处在于它把自己的身份伪装成了一个从库。它向MySQL主库发送复制协议请求注册一个假的server-id主库就会像对待真从库一样把binlog源源不断推给它。而Canal拿到binlog后不会回放SQL而是解析成结构化的数据变更事件再推送给Kafka、RocketMQ或者直接写到下游存储。这个设计有三个好处第一完全不侵入业务代码业务方根本感知不到多了一个“订阅者”第二利用的是MySQL成熟的复制协议稳定性远比自己写的轮询脚本可靠第三解析binlog能拿到精确的变更前后值天然适合异构同步。有一个细节我特别提醒每个订阅binlog的消费者都必须有唯一的server-id。如果你起了两个Canal实例却配了同一个server-id它们会互相踢对方下线表现就是日志疯狂报错、同步一会儿停一会儿。2.3 关键参数配置与原理要让MySQL乖乖“交出”数据有几组参数必须配置对# my.cnf 核心配置 [mysqld] server-id 100 log-bin mysql-bin binlog_format ROW binlog_row_image FULL expire_logs_days 7 gtid_mode ON enforce_gtid_consistency ON逐一解释这些参数的用意server-idMySQL复制集群中每个节点的唯一标识。Canal之所以要配置一个不冲突的server-id就是为了不让主库把它误认成其他从库。log-bin开启binlog开关不开启的话后面全部白搭。binlog_format ROW上面已经讲过异构同步必需。expire_logs_daysbinlog保留时长。这个参数非常关键又最容易被忽略——如果保留时间太短比如只有1天而你的同步任务因为故障停了20个小时恢复时Canal会报“找不到binlog文件”因为它要读的位置已经被清理了。gtid_mode ON开启GTID后每个事务有了全局唯一IDCanal/Flink CDC可以根据GTID精确判断哪些事务已经处理过断点续传时不会重复消费。5.7及以上版本建议直接开启。3. 实操落地一套完整的MySQL实时同步链路搭建3.1 基于Canal Kafka 应用消费的链路设计下面这条链路是我在实践中用得最多、也最稳的一套组合适合大部分实时同步场景MySQL(binlog) → Canal → Kafka → 下游消费应用 → 目标存储选Canal而不选其他组件我的理由很实在Canal是Java写的部署就是一个jar包配置简单社区活跃度够最关键是它支持的MySQL版本跨度大5.6到8.0都能兼容。Kafka在这里承担了缓冲和削峰的角色——Canal从binlog读到的变更消息先落到Kafka下游消费程序按自己的节奏处理不会因为目标存储抖动就把Canal阻塞住。3.2 Canal端配置要点Canal的配置文件主要看两个canal.properties和instance.properties。canal.properties里重点配置 server 模式# canal.properties 核心配置 canal.serverMode kafka canal.mq.servers 192.168.1.10:9092 canal.mq.producer.group canal canal.mq.topic my-topicinstance.properties里配置的是数据源信息和订阅规则# instance.properties 核心配置 canal.instance.mysql.slaveId 1234 canal.instance.master.address 192.168.1.20:3306 canal.instance.dbUsername canal canal.instance.dbPassword canal_pass canal.instance.connectionCharset UTF-8 canal.instance.filter.regex mydb\\..*几个容易踩坑的点第一订阅规则canal.instance.filter.regex的写法。mydb\\..*表示订阅mydb库的所有表注意两个反斜杠是Java字符串转义后的结果写成单个\\.在properties文件里会被当成普通字符导致过滤规则完全无效。第二canal.instance.mysql.slaveId必须和MySQL主库、其他从库的server-id都不冲突。我习惯用一个独立网段的值比如512避免和线上实例的ID撞车。第三Canal账号的权限。MySQL这边要建一个专门账号只授予复制权限不要用rootCREATE USER canal% IDENTIFIED BY canal_pass; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; FLUSH PRIVILEGES;3.3 下游消费应用的数据处理策略Kafka里拿到Canal消息后数据格式大致长这样{ id: 12345, database: mydb, table: orders, type: UPDATE, ts: 1710000000, data: [ { id: 1001, order_no: A1001, amount: 99.90, status: PAID } ], old: [ { status: PENDING } ] }type字段是变更类型常见的有INSERT、UPDATE、DELETEdata是变更后的数据行old是变更前的旧值仅UPDATE有。处理时的三个关键点第一是幂等性。MySQL里同一行被更新了两次Canal会产生两条UPDATE消息。但下游如果先处理了第二条、再处理第一条比如Kafka分区内乱序就会导致最终值错误。解决办法是下游处理时以主键更新时间做幂等校验或者干脆因为Kafka同一个key的消息在分区内天然有序设计时保证同一主键的消息路由到同一个分区。第二是DELETE类型的处理。不要只按data字段删除因为Canal的DELETE消息里通常只有主键值。正确做法是从data中取主键字段到目标端执行删除然后忽略其他字段。第三是DDL变更的处理。默认配置下Canal会把DDL也发到Kafka但消息格式和DML不一样。如果下游没有处理DDL的能力建议在canal.properties里关闭DDL同步canal.instance.ddl.filter true否则下游消费程序解析到未知格式的JSON会直接抛异常。3.4 不引入Canal的轻量替代方案如果团队不想引入Canal和Kafka这两套组件还有一个轻量级选择直接用Flink CDC。Flink CDC内置了Debezium引擎只需要写一个Flink SQL作业就能实现MySQL到Doris/ClickHouse/TDengine的同步。一个简单的Flink SQL同步任务-- 创建MySQL CDC源表 CREATE TABLE mysql_orders ( id INT PRIMARY KEY, order_no STRING, amount DECIMAL(10,2), status STRING ) WITH ( connector mysql-cdc, hostname 192.168.1.20, port 3306, username flink, password flink_pass, database-name mydb, table-name orders ); -- 创建目标表以Doris为例 CREATE TABLE doris_orders ( id INT PRIMARY KEY, order_no STRING, amount DECIMAL(10,2), status STRING ) WITH ( connector doris, fenodes 192.168.1.30:8030, table.identifier doris_db.orders, username root, password , sink.label-prefix doris_orders ); -- 实时同步 INSERT INTO doris_orders SELECT * FROM mysql_orders;如果你只需要每天离线同步一次那直接用DataX或者mysqldump导出就够没必要把CDC全家桶都拉起来。方案的大忌是过度设计实时同步的维护成本是离线同步的好几倍能用定时任务解决的事不要天天伺候Kafka和Canal。4. 异构场景实战MySQL表结构自动转TDengine超级表子表4.1 时序数据场景下的存储选型困境聊完通用链路我想单独把TDengine这个场景拿出来讲。最近几年IoT项目井喷很多团队一开始图省事把设备上报数据直接写MySQL等到每天几千万条记录进来后才发现根本扛不住——表现在三个方面存储成本爆炸就不用说了单表几亿条记录后查询性能急剧下降尤其是按时间范围聚合这种时序查询MySQL的B树在这种场景下发挥不出优势。更麻烦的是MySQL表结构是固定死的设备多了想加个标签维度还得去改表业务代码一起跟着改。把这类数据同步到TDengine本质上不是简单的“搬运”而是要做一次表结构重构。因为TDengine的存储模型和MySQL不一样它专门为时序数据设计了“超级表 子表”的模型这是整个同步方案里最需要动脑子的一环。4.2 超级表与子表的核心概念先花两分钟讲清楚TDengine的模型否则后面的映射规则看不懂。超级表Super Table可以理解成一个抽象模板它定义了时序数据的“采集量”字段比如温度、湿度、电压和“标签”字段比如设备ID、设备类型、位置。子表Child Table是基于超级表实例化出来的具体表一个设备一张子表标签值固定采集量字段随时间不断追加。举个例子你有一个超级表devices_metric定义如下CREATE STABLE devices_metric ( ts TIMESTAMP, temperature FLOAT, humidity FLOAT, voltage INT ) TAGS ( device_id VARCHAR(32), device_type VARCHAR(16), location VARCHAR(64) );然后每个设备会生成一张子表CREATE TABLE device_001 USING devices_metric TAGS (DEV001, sensor, beijing); CREATE TABLE device_002 USING devices_metric TAGS (DEV002, sensor, shanghai);这种模型的好处是查询时可以直接在超级表上做聚合TDengine会自动跨所有子表并行查标签过滤走的是元数据索引性能远超传统数据库的“一张大表到处扫”。4.3 MySQL普通表如何映射为超级表子表现在问题来了MySQL里的业务表是普通二维表怎么把它变成TDengine的超级表和子表我在实际项目里拿到一张MySQL表会先做一次“字段分类”时间列如果表里有一个时间字段比如create_time或report_time它就是TDengine的时间戳主键。采集量列那些不断变化、需要按时间记录的数值字段温度、速度、电量等对应超级表的普通列。标签列设备ID、站点ID、类型等用于过滤和分组的字段对应超级表的TAG。一个具体的映射案例。假设MySQL里有张表长这样CREATE TABLE device_log ( id BIGINT PRIMARY KEY, device_id VARCHAR(32), device_type VARCHAR(16), location VARCHAR(64), temperature FLOAT, humidity FLOAT, report_time DATETIME );转换成TDengine的DDL就是-- 先建超级表 CREATE STABLE device_log ( report_time TIMESTAMP, temperature FLOAT, humidity FLOAT ) TAGS ( device_id VARCHAR(32), device_type VARCHAR(16), location VARCHAR(64) );关键点来了MySQL的主键id在这里是不保留的。因为TDengine按照时间戳 设备维度来组织数据MySQL主键在时序场景没有太大意义。同理每个设备ID都会创建一张子表-- 为每个不同的device_id创建子表 CREATE TABLE device_log_dev001 USING device_log TAGS ( DEV001, sensor, beijing );这个自动建表的过程可以写脚本完成。从MySQL里SELECT DISTINCT device_id, device_type, location然后循环拼接SQL。类型转换上也需要特别处理我整理了实战中用到的映射规则MySQL类型TDengine类型说明DATETIME / TIMESTAMPTIMESTAMP直接转换注意时区FLOAT / DOUBLEFLOAT / DOUBLE直接转换INT / BIGINTINT / BIGINT直接转换VARCHAR(n)VARCHAR(n)用于TAG字段时建议给定合理长度DECIMAL(p, s)FLOAT / DOUBLETDengine精度有限高精度DECIMAL会丢精度TINYINT(1)BOOL / INTMySQL布尔语义建议转INTTEXT / JSONBINARY / NCHAR大字段尽量避免时序场景基本用不到4.4 同步落地全量增量的数据写入策略表结构映射做好了数据搬运本身相对简单但有两条路可以走历史数据全量迁移用DataX或者直接用SQL查询后批量写入。注意TDengine写入时要用参数绑定模式逐条INSERT的吞吐量很低。Java示例// 批量写入示例 String sql INSERT INTO device_log_dev001 VALUES (?, ?, ?); PreparedStatement pstmt conn.prepareStatement(sql); for (DataPoint point : batchList) { pstmt.setTimestamp(1, point.getTs()); pstmt.setFloat(2, point.getTemperature()); pstmt.setFloat(3, point.getHumidity()); pstmt.addBatch(); } pstmt.executeBatch();增量数据实时同步可以把Flink CDC的sink指向TDengine或者用前面的Canal链路下游消费程序收到消息后构造对应子表的INSERT语句。这里有个容易踩坑的细节TDengine的数据模型要求同一张子表的时间戳必须单调递增如果消费程序是乱序处理Kafka消息的可能会往子表里写入时间戳倒退的数据TDengine会拒绝写入。解决办法是在消费端做一次按时间戳排序的缓存窗口或者干脆保证Kafka分区内消息按主键有序。5. 常见问题与排查技巧实录5.1 数据同步时效性排查时区不同差8小时这类问题在跨时区部署、或者MySQL和下游存储时区不一致时非常常见。MySQL里的DATETIME不带时区信息但Canal解析binlog时会按照MySQL实例所在时区处理然后转成时间戳传给下游下游如果按目标库的时区再转一次就会差出8小时甚至更多。我的排查思路是先看Canal端的canal.instance.connectionCharset和MySQL的time_zone是否一致再看下游消费程序解析JSON里的时间字段时用的是哪个时区。建议统一约定全链路所有组件一律使用UTC8的Asia/Shanghai时区消息传递统一用时间戳毫秒值或者ISO8601字符串不要在消息里传无时区的裸时间。5.2 表结构变更导致同步突然中断这是实时同步最恶心的问题之一没有任何报错提示Canal日志安静得跟什么都没发生一样但下游数据就是不更新了。大概率是MySQL执行了DDL比如ALTER TABLE orders ADD COLUMN new_field VARCHAR(32)。Canal默认会把DDL也同步到下游但下游消费程序还在按旧的消息结构处理解析失败后直接跳过或者抛出异常。处理方式有两种一是确保下游有DDL同步能力消费程序能识别DDL消息并按需变更目标表结构二是更简单的处理在Canal配置里显式关闭DDL同步canal.instance.ddl.filter true然后在MySQL DDL变更时手动维护下游表结构。5.3 字段类型不兼容导致数据漂移MySQL和下游存储的字段类型完全对齐是同步方案里最容易出问题的地方。拿DECIMAL来说MySQL的DECIMAL(18,4)在Canal解析后会变成BigDecimal对象如果你转成JSON序列化时用了toString下游用浮点类型接收精度就可能丢位。另一个典型是TINYINT(1)MySQL逻辑上就是布尔值但binlog里的数值是0和1部分下游驱动会把它解析成true/false导致下游数据从1变成了true。我的建议是在同步链路的边界层做一次显式类型转换不要依赖框架自动处理。比如Kafka消息的JSON序列化里把decimal统一转成字符串下游解析时再用BigDecimal还原tinyint(1)统一约定转成INT避免布尔歧义。5.4 同步延迟监控与告警最后说说延迟问题。Canal/Kafka链路跑久了难免出现下游消费变慢导致消息堆积。延迟从秒级变成分钟级如果没人发现等业务方来投诉就晚了。我自己的做法是在消费程序里记录每个消息的处理时间戳和消息里携带的MySQL binlog位置时间做差值每30秒上报一次最大延迟超过阈值就告警。常见的原因包括下游存储写入变慢、目标表索引缺失导致update语句全表扫描、Kafka消费者线程数不够。另外给一个独家心得大事务是延迟的隐形杀手。MySQL里一个更新100万行的事务binlog会产生100万条变更记录Canal会一次性推到Kafka。下游消费时即使一条只处理5毫秒也要5000秒才能处理完延迟瞬间拉满。遇到这种情况要么在业务侧控制单事务更新行数要么给下游准备临时扩容的应急预案。做个同步方案先想清楚场景再选型运维优先级大于炫技优先级。我就是那个把Canal、Kafka、TDengine、Flink CDC都轮流折腾过一遍的人——有些坑真的不必亲自去踩。但凡你要动MySQL数据同步希望这篇能帮你省掉至少一周的排障时间。