據(jù)同步方案選型與實戰(zhàn):從binlog到TDengine)
做數(shù)據(jù)同步這些年從最早的mysqldump全量拷貝到后來折騰Canal、binlog、Flink CDC踩過的坑能寫一本小冊子。今天把MySQL數(shù)據(jù)出海這件事徹底拆開聊聊——不是給你堆一堆工具名詞而是把方案選型背后的邏輯、每條鏈路的真實優(yōu)缺點、以及我實際落地時踩過的坑一次性講透。所謂“數(shù)據(jù)出海”本質就是讓MySQL里沉淀的數(shù)據(jù)流動起來進入更合適它的地方??赡苁请x線數(shù)倉做分析可能是實時大屏做指標展示也可能是像TDengine這樣的時序數(shù)據(jù)庫承載海量監(jiān)控數(shù)據(jù)。這篇文章適合正在做數(shù)據(jù)同步選型、或者已經在同步路上被各種詭異問題折磨的讀者無論你是架構師、后端開發(fā)還是DBA都能從中找到可直接抄作業(yè)的方案。1. 數(shù)據(jù)出海場景拆解為什么要把MySQL數(shù)據(jù)搬出去1.1 數(shù)據(jù)同步的典型業(yè)務場景先別急著選工具得先搞清楚“出?!钡降滓ツ摹N医佑|過的項目里最典型的是下面幾類第一類是離線分析場景。業(yè)務庫跑著線上交易幾千萬行的訂單表不可能讓分析師直接在上面跑復雜聚合查詢哪怕加了索引也會把主庫拖垮。這時候需要把MySQL數(shù)據(jù)同步到數(shù)據(jù)倉庫比如Hive、Doris、ClickHouse做T1甚至T0的分析。第二類是實時消費場景。比如大屏要實時顯示訂單量或者下游搜索服務要實時同步商品信息數(shù)據(jù)從MySQL產生到下游可用的延遲必須控制在秒級甚至毫秒級。第三類是異構存儲場景。比如設備上報的時序數(shù)據(jù)、監(jiān)控指標存在MySQL里不僅存儲成本高查詢性能也完全不對路這時候需要同步到TDengine這類時序數(shù)據(jù)庫。第四類是數(shù)據(jù)分發(fā)場景。多個業(yè)務系統(tǒng)都需要同一份基礎數(shù)據(jù)比如用戶中心的數(shù)據(jù)要同步給訂單系統(tǒng)、客服系統(tǒng)、推薦系統(tǒng)不能每個系統(tǒng)都直連主庫。每一種場景對應的同步方案差別非常大離線場景可以接受小時級延遲實時場景要求秒級異構存儲還得考慮表結構映射和類型轉換。這決定了你不可能一套方案走天下。1.2 方案選型的四個關鍵維度選方案之前先學會畫坐標軸。我通常用四個維度來評估一套同步方案同步延遲從MySQL數(shù)據(jù)變更到目標端可見能接受多少延遲分鐘級還是秒級這直接淘汰一批方案。同步粒度是全量同步還是增量同步是否需要同步表結構變更DDL還是只需要同步數(shù)據(jù)變更DML一致性要求是否允許重復數(shù)據(jù)是否需要保證事務的原子性下游有沒有冪等消費能力技術棧匹配團隊熟悉Java還是Python愿不愿意引入額外組件比如Canal、Kafka還是希望盡量少依賴、少運維拿這四個維度去套很多選擇自然就清晰了。比如你只需要每天晚上全量同步一次訂單數(shù)據(jù)給數(shù)倉那搞一套CDC實時同步成本反而高得沒必要反過來如果業(yè)務方要求實時看到用戶下單那你用mysqldump凌晨三點導一次顯然也是交不了差的。1.3 主流方案全景對比我把這些年見過、用過的方案整理成一張對比表方便你做選擇題方案同步延遲增量支持部署復雜度適用場景核心痛點mysqldump 全量導出無實時性不支持極低小表備份、一次性遷移數(shù)據(jù)量大了導出時間長會鎖表MySQL 主從復制秒級以內原生支持低MySQL到MySQL的同構同步只支持MySQL源和MySQL目標異構完全不行Canal Kafka秒級支持中高異構數(shù)據(jù)源、實時數(shù)倉需要維護Canal和Kafka兩套組件Flink CDC秒級支持中實時數(shù)倉、異構存儲學習成本偏高狀態(tài)后端調優(yōu)有門檻DataX分鐘級離線低離線批量同步準實時能力弱調度靠自己搭業(yè)務雙寫實時天然支持無小規(guī)模/新項目侵入業(yè)務代碼改造風險大這里多說一句雙寫看著最省事實際最坑。我見過一個團隊在訂單系統(tǒng)里加了雙寫邏輯結果下游存儲抖了一下事務提交就失敗連帶主庫寫也回滾線上訂單全掛。業(yè)務雙寫只能用在非核心鏈路核心數(shù)據(jù)還是得靠CDC老老實實抓binlog。2. 核心機制原理解析binlog是如何成為同步樞紐的2.1 binlog的三種格式與CDC的必然選擇聊MySQL數(shù)據(jù)同步繞不開binlog。它本質是MySQL服務端記錄的所有數(shù)據(jù)變更事件的日志文件你可以把它理解成數(shù)據(jù)庫的“黑匣子”誰在什么時間改了哪張表的哪一行、改之前是什么值、改之后是什么值全都記錄在案。但這三種格式差別很大STATEMENT格式記錄的是SQL語句原文比如UPDATE t SET namex WHERE id1。優(yōu)點是日志量小但缺點是語句執(zhí)行依賴上下文比如用了NOW()函數(shù)、自增列從庫回放時結果可能不一樣。ROW格式記錄的是每一行變更前后的具體值比如“id1這一行的name從y變成了x”。優(yōu)點是絕對準確缺點是日志量比STATEMENT大很多。MIXED格式MySQL自己判斷默認用STATEMENT遇到不確定函數(shù)時自動切換ROW。做數(shù)據(jù)同步必須用ROW格式。原因很簡單下游需要知道具體是哪一行變了、舊值是什么、新值是什么才能在其他存儲里做對應的更新。STATEMENT格式雖然日志小但對異構同步來說等于沒有信息——你拿到了UPDATE t SET namex WHERE id1這句SQL卻根本不知道在目標端要改哪一行。還有一個重要參數(shù)是binlog_row_imageFULL這個參數(shù)決定ROW格式下binlog是否記錄所有列的“前鏡像”和“后鏡像”。默認值是FULL即記錄整行所有列有些團隊為了省日志量改成MINIMAL只記錄被修改的列和能定位主鍵的列。省了空間但會帶來一個大問題下游如果拿到binlog數(shù)據(jù)后想對比“到底改了哪些列”就會出現(xiàn)信息不全的情況。做同步建議保持FULL。2.2 從主從復制到CDCCanal/Flink CDC偽裝成從庫的秘密理解了binlog也就理解了Canal和Flink CDC的工作原理。MySQL原生主從復制是這樣的主庫開啟binlog從庫啟動兩個線程一個IO線程連上主庫偽裝成從庫把主庫的binlog拉過來寫到本地relay log另一個SQL線程讀relay log并回放完成數(shù)據(jù)同步。Canal的高明之處在于它把自己的身份偽裝成了一個從庫。它向MySQL主庫發(fā)送復制協(xié)議請求注冊一個假的server-id主庫就會像對待真從庫一樣把binlog源源不斷推給它。而Canal拿到binlog后不會回放SQL而是解析成結構化的數(shù)據(jù)變更事件再推送給Kafka、RocketMQ或者直接寫到下游存儲。這個設計有三個好處第一完全不侵入業(yè)務代碼業(yè)務方根本感知不到多了一個“訂閱者”第二利用的是MySQL成熟的復制協(xié)議穩(wěn)定性遠比自己寫的輪詢腳本可靠第三解析binlog能拿到精確的變更前后值天然適合異構同步。有一個細節(jié)我特別提醒每個訂閱binlog的消費者都必須有唯一的server-id。如果你起了兩個Canal實例卻配了同一個server-id它們會互相踢對方下線表現(xiàn)就是日志瘋狂報錯、同步一會兒停一會兒。2.3 關鍵參數(shù)配置與原理要讓MySQL乖乖“交出”數(shù)據(jù)有幾組參數(shù)必須配置對# 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逐一解釋這些參數(shù)的用意server-idMySQL復制集群中每個節(jié)點的唯一標識。Canal之所以要配置一個不沖突的server-id就是為了不讓主庫把它誤認成其他從庫。log-bin開啟binlog開關不開啟的話后面全部白搭。binlog_format ROW上面已經講過異構同步必需。expire_logs_daysbinlog保留時長。這個參數(shù)非常關鍵又最容易被忽略——如果保留時間太短比如只有1天而你的同步任務因為故障停了20個小時恢復時Canal會報“找不到binlog文件”因為它要讀的位置已經被清理了。gtid_mode ON開啟GTID后每個事務有了全局唯一IDCanal/Flink CDC可以根據(jù)GTID精確判斷哪些事務已經處理過斷點續(xù)傳時不會重復消費。5.7及以上版本建議直接開啟。3. 實操落地一套完整的MySQL實時同步鏈路搭建3.1 基于Canal Kafka 應用消費的鏈路設計下面這條鏈路是我在實踐中用得最多、也最穩(wěn)的一套組合適合大部分實時同步場景MySQL(binlog) → Canal → Kafka → 下游消費應用 → 目標存儲選Canal而不選其他組件我的理由很實在Canal是Java寫的部署就是一個jar包配置簡單社區(qū)活躍度夠最關鍵是它支持的MySQL版本跨度大5.6到8.0都能兼容。Kafka在這里承擔了緩沖和削峰的角色——Canal從binlog讀到的變更消息先落到Kafka下游消費程序按自己的節(jié)奏處理不會因為目標存儲抖動就把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里配置的是數(shù)據(jù)源信息和訂閱規(guī)則# 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\\..*幾個容易踩坑的點第一訂閱規(guī)則canal.instance.filter.regex的寫法。mydb\\..*表示訂閱mydb庫的所有表注意兩個反斜杠是Java字符串轉義后的結果寫成單個\\.在properties文件里會被當成普通字符導致過濾規(guī)則完全無效。第二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 下游消費應用的數(shù)據(jù)處理策略Kafka里拿到Canal消息后數(shù)據(jù)格式大致長這樣{ 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是變更后的數(shù)據(jù)行old是變更前的舊值僅UPDATE有。處理時的三個關鍵點第一是冪等性。MySQL里同一行被更新了兩次Canal會產生兩條UPDATE消息。但下游如果先處理了第二條、再處理第一條比如Kafka分區(qū)內亂序就會導致最終值錯誤。解決辦法是下游處理時以主鍵更新時間做冪等校驗或者干脆因為Kafka同一個key的消息在分區(qū)內天然有序設計時保證同一主鍵的消息路由到同一個分區(qū)。第二是DELETE類型的處理。不要只按data字段刪除因為Canal的DELETE消息里通常只有主鍵值。正確做法是從data中取主鍵字段到目標端執(zhí)行刪除然后忽略其他字段。第三是DDL變更的處理。默認配置下Canal會把DDL也發(fā)到Kafka但消息格式和DML不一樣。如果下游沒有處理DDL的能力建議在canal.properties里關閉DDL同步canal.instance.ddl.filter true否則下游消費程序解析到未知格式的JSON會直接拋異常。3.4 不引入Canal的輕量替代方案如果團隊不想引入Canal和Kafka這兩套組件還有一個輕量級選擇直接用Flink CDC。Flink CDC內置了Debezium引擎只需要寫一個Flink SQL作業(yè)就能實現(xiàn)MySQL到Doris/ClickHouse/TDengine的同步。一個簡單的Flink SQL同步任務-- 創(chuàng)建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 ); -- 創(chuàng)建目標表以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. 異構場景實戰(zhàn)MySQL表結構自動轉TDengine超級表子表4.1 時序數(shù)據(jù)場景下的存儲選型困境聊完通用鏈路我想單獨把TDengine這個場景拿出來講。最近幾年IoT項目井噴很多團隊一開始圖省事把設備上報數(shù)據(jù)直接寫MySQL等到每天幾千萬條記錄進來后才發(fā)現(xiàn)根本扛不住——表現(xiàn)在三個方面存儲成本爆炸就不用說了單表幾億條記錄后查詢性能急劇下降尤其是按時間范圍聚合這種時序查詢MySQL的B樹在這種場景下發(fā)揮不出優(yōu)勢。更麻煩的是MySQL表結構是固定死的設備多了想加個標簽維度還得去改表業(yè)務代碼一起跟著改。把這類數(shù)據(jù)同步到TDengine本質上不是簡單的“搬運”而是要做一次表結構重構。因為TDengine的存儲模型和MySQL不一樣它專門為時序數(shù)據(jù)設計了“超級表 子表”的模型這是整個同步方案里最需要動腦子的一環(huán)。4.2 超級表與子表的核心概念先花兩分鐘講清楚TDengine的模型否則后面的映射規(guī)則看不懂。超級表Super Table可以理解成一個抽象模板它定義了時序數(shù)據(jù)的“采集量”字段比如溫度、濕度、電壓和“標簽”字段比如設備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會自動跨所有子表并行查標簽過濾走的是元數(shù)據(jù)索引性能遠超傳統(tǒng)數(shù)據(jù)庫的“一張大表到處掃”。4.3 MySQL普通表如何映射為超級表子表現(xiàn)在問題來了MySQL里的業(yè)務表是普通二維表怎么把它變成TDengine的超級表和子表我在實際項目里拿到一張MySQL表會先做一次“字段分類”時間列如果表里有一個時間字段比如create_time或report_time它就是TDengine的時間戳主鍵。采集量列那些不斷變化、需要按時間記錄的數(shù)值字段溫度、速度、電量等對應超級表的普通列。標簽列設備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按照時間戳 設備維度來組織數(shù)據(jù)MySQL主鍵在時序場景沒有太大意義。同理每個設備ID都會創(chuàng)建一張子表-- 為每個不同的device_id創(chuàng)建子表 CREATE TABLE device_log_dev001 USING device_log TAGS ( DEV001, sensor, beijing );這個自動建表的過程可以寫腳本完成。從MySQL里SELECT DISTINCT device_id, device_type, location然后循環(huán)拼接SQL。類型轉換上也需要特別處理我整理了實戰(zhàn)中用到的映射規(guī)則MySQL類型TDengine類型說明DATETIME / TIMESTAMPTIMESTAMP直接轉換注意時區(qū)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 同步落地全量增量的數(shù)據(jù)寫入策略表結構映射做好了數(shù)據(jù)搬運本身相對簡單但有兩條路可以走歷史數(shù)據(jù)全量遷移用DataX或者直接用SQL查詢后批量寫入。注意TDengine寫入時要用參數(shù)綁定模式逐條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();增量數(shù)據(jù)實時同步可以把Flink CDC的sink指向TDengine或者用前面的Canal鏈路下游消費程序收到消息后構造對應子表的INSERT語句。這里有個容易踩坑的細節(jié)TDengine的數(shù)據(jù)模型要求同一張子表的時間戳必須單調遞增如果消費程序是亂序處理Kafka消息的可能會往子表里寫入時間戳倒退的數(shù)據(jù)TDengine會拒絕寫入。解決辦法是在消費端做一次按時間戳排序的緩存窗口或者干脆保證Kafka分區(qū)內消息按主鍵有序。5. 常見問題與排查技巧實錄5.1 數(shù)據(jù)同步時效性排查時區(qū)不同差8小時這類問題在跨時區(qū)部署、或者MySQL和下游存儲時區(qū)不一致時非常常見。MySQL里的DATETIME不帶時區(qū)信息但Canal解析binlog時會按照MySQL實例所在時區(qū)處理然后轉成時間戳傳給下游下游如果按目標庫的時區(qū)再轉一次就會差出8小時甚至更多。我的排查思路是先看Canal端的canal.instance.connectionCharset和MySQL的time_zone是否一致再看下游消費程序解析JSON里的時間字段時用的是哪個時區(qū)。建議統(tǒng)一約定全鏈路所有組件一律使用UTC8的Asia/Shanghai時區(qū)消息傳遞統(tǒng)一用時間戳毫秒值或者ISO8601字符串不要在消息里傳無時區(qū)的裸時間。5.2 表結構變更導致同步突然中斷這是實時同步最惡心的問題之一沒有任何報錯提示Canal日志安靜得跟什么都沒發(fā)生一樣但下游數(shù)據(jù)就是不更新了。大概率是MySQL執(zhí)行了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 字段類型不兼容導致數(shù)據(jù)漂移MySQL和下游存儲的字段類型完全對齊是同步方案里最容易出問題的地方。拿DECIMAL來說MySQL的DECIMAL(18,4)在Canal解析后會變成BigDecimal對象如果你轉成JSON序列化時用了toString下游用浮點類型接收精度就可能丟位。另一個典型是TINYINT(1)MySQL邏輯上就是布爾值但binlog里的數(shù)值是0和1部分下游驅動會把它解析成true/false導致下游數(shù)據(jù)從1變成了true。我的建議是在同步鏈路的邊界層做一次顯式類型轉換不要依賴框架自動處理。比如Kafka消息的JSON序列化里把decimal統(tǒng)一轉成字符串下游解析時再用BigDecimal還原tinyint(1)統(tǒng)一約定轉成INT避免布爾歧義。5.4 同步延遲監(jiān)控與告警最后說說延遲問題。Canal/Kafka鏈路跑久了難免出現(xiàn)下游消費變慢導致消息堆積。延遲從秒級變成分鐘級如果沒人發(fā)現(xiàn)等業(yè)務方來投訴就晚了。我自己的做法是在消費程序里記錄每個消息的處理時間戳和消息里攜帶的MySQL binlog位置時間做差值每30秒上報一次最大延遲超過閾值就告警。常見的原因包括下游存儲寫入變慢、目標表索引缺失導致update語句全表掃描、Kafka消費者線程數(shù)不夠。另外給一個獨家心得大事務是延遲的隱形殺手。MySQL里一個更新100萬行的事務binlog會產生100萬條變更記錄Canal會一次性推到Kafka。下游消費時即使一條只處理5毫秒也要5000秒才能處理完延遲瞬間拉滿。遇到這種情況要么在業(yè)務側控制單事務更新行數(shù)要么給下游準備臨時擴容的應急預案。做個同步方案先想清楚場景再選型運維優(yōu)先級大于炫技優(yōu)先級。我就是那個把Canal、Kafka、TDengine、Flink CDC都輪流折騰過一遍的人——有些坑真的不必親自去踩。但凡你要動MySQL數(shù)據(jù)同步希望這篇能幫你省掉至少一周的排障時間。