
MySQL 到 BigQuery 的數據同步最容易被低估的問題就是時間窗口。無論是定時導出還是按updated_at增量拉取本質上都屬于 periodic syncs。它們的共同點是數據庫里的變化并不會等待調度任務開始也不會按周期整齊地落入邊界。一次刪除、一條字段被改回舊值、一張表在夜間被大批量 UPDATE 后又改回來這些事件都可能發生在兩批同步任務的間隙最終 BigQuery 里的數據既不是源表的真實狀態也不是任何歷史時刻的真實狀態。CDCChange Data Capture通過讀取 MySQL binlog把每一條數據變更作為事件流送到 BigQuery正好從機制上補上了這個缺口。這篇文章圍繞周期同步會漏什么、binlog 為什么能避免漏、落地時要注意什么展開適合正在設計數據管道、給數倉接增量數據或者被批量任務數據不一致問題困擾的開發者與數據工程師。1. 周期同步在 MySQL 到 BigQuery 場景下到底漏了什么1.1 常見的三種周期同步寫法先看最常用的三種同步方式它們并不只是實現細節不同能捕獲的數據變化粒度也完全不同。第一種是全量導出覆蓋。直接把 MySQL 表導出成文件或通過 SQL 拉取寫入 BigQuery 臨時表再覆蓋目標表。這種方式能保證目標表最終狀態一致但同步窗口很長且 BigQuery 做覆蓋時下游可能讀到一半數據。數據量一旦上億這個方案基本不可持續。第二種是按自增 ID 增量拉取。記錄max(id)每次只拉大于該 ID 的行像這樣SELECT * FROM orders WHERE id :last_max_id ORDER BY id;這個方案只能捕獲新增數據。業務表一旦發生 UPDATE主鍵 ID 不變增量 SQL 永遠拉不到這一行。DELETE 更不會出現在結果里。第三種是按更新時間戳增量拉取SELECT * FROM orders WHERE updated_at :last_sync_ts;這是目前最常見的周期同步方案前提是業務表有updated_at字段并且所有寫入路徑都正確更新這個字段。實際項目里這個前提經常被破壞某些批量導入腳本沒有更新updated_at某些框架寫入時沒有映射該字段于是出現數據明明變了增量 SQL 卻查不到的問題。1.2 周期同步一定會錯過的幾類變更物理刪除是最典型的一類。DELETE 之后這條記錄不再存在于表中任何基于當前表狀態的 SELECT 都無法發現它曾經存在過。全量對拍能發現問題但只能事后補救而且對拍本身在大表上成本極高。沒有更新時間字段或者更新時沒有寫入時間戳也是一類。訂單表如果通過第三方系統直接改庫或者 DBA 手工執行 UPDATE 時沒有維護updated_at那么增量邊界從一開始就是錯的。同周期內狀態回跳同樣會被掩蓋。假設訂單在 00:00:10 從pending改為paid00:00:20 又改回pending。周期任務在 01:00 運行拉到的最終狀態還是pending。從業務角度看中間那次paid狀態也曾經是真實數據但周期同步完全感知不到。高頻更新更不用說。一張促銷表每秒更新幾千行周期任務每隔 5 分鐘拉一次單行在周期內被反復更新后最終拉到的只是最后一次值中間所有取值全部丟失。1.3 為什么不是多跑幾次就能解決周期同步的失敗模式是邏輯性漏數據不是漏跑任務。把調度頻率從小時改成分鐘只是縮小時間窗口并沒有改變讀取當前表狀態的本質。一張表在周期內發生了 100 次更新周期同步只能看到最后一行binlog 能看到 100 個事件并且每個事件都保留前鏡像和后鏡像。這就是原理層面的差異。周期同步試圖通過查詢結果反推變化而 binlog 是 MySQL 自己記錄的寫操作流水賬。流水賬不會因為業務表沒有updated_at就缺頁也不會因為 DELETE 后記錄消失就抹去歷史。1.4 三種方案的能力對比維度全量快照增量字段輪詢binlog CDC刪除事件全量對拍后才發現通常無法發現每條 DELETE 都有對應事件更新歷史只有最后狀態只有最后一次變更每次 UPDATE 都有前鏡像和后鏡像對業務表要求無必須有updated_at等字段無binlog 與業務表結構獨立實時性取決于調度周期取決于調度周期秒級到分鐘級可配置對源庫壓力大全表掃描代價高中等取決于索引較小讀取日志而不是反復掃描表從這張表能看出周期同步不是慢而是漏。CDC 的價值不是讓同步更快而是讓變化過程本身可見。2. binlog 為什么能捕捉每一次變化CDC 的原理2.1 binlog 是什么binlog 是 MySQL 的二進制日志記錄所有改變數據庫內容的操作包括 INSERT、UPDATE、DELETE以及部分 DDL。MySQL 主從復制、崩潰恢復、數據恢復都依賴它。可以理解為 MySQL 把每一次寫操作按順序寫到一本流水賬上。binlog 并不是默認可用的。MySQL 5.7 中log_bin默認關閉8.0 默認開啟但不同發行版和云廠商的默認值可能不同落地前必須先確認。如果 binlog 沒有開啟后續所有 CDC 方案都無從談起。2.2 ROW 格式給 CDC 提供了什么binlog 有三種格式STATEMENT、ROW、MIXED。STATEMENT 格式記錄的是 SQL 語句本身例如UPDATE orders SET statuspaid WHERE id1001;。這種格式日志量小但無法可靠還原每一行在語句執行前后的具體值。MIXED 格式是兩者的混合MySQL 會根據語句類型自動選擇但對于 CDC 場景依然不夠穩定。CDC 要求使用 ROW 格式。ROW 格式下binlog 直接記錄行的變化包括字段級的前鏡像和后鏡像。具體來說INSERT 事件包含插入后的完整行數據。UPDATE 事件包含變更前的整行數據和變更后的整行數據。DELETE 事件包含刪除前的整行數據。這意味著 CDC 消費者不僅能知道某張表發生了變化還能拿到 哪一行的哪個字段從什么值變成什么值。2.3 CDC 連接器如何消費 binlogDebezium、Flink CDC 這類工具在原理上會偽裝成 MySQL 從庫。它們通過 MySQL 的復制協議從主庫拉取 binlog并把 binlog 里的二進制事件解析成結構化的 JSON 變更事件。連接器需要記錄自己的消費位點。傳統方式是記錄 binlog 文件名加偏移量例如mysql-bin.000023的position 45123。更可靠的方式是使用 GTID即全局事務標識符。GTID 能唯一標識每個事務即使 binlog 文件被清理只要 MySQL 實例保留了完整的事務歷史連接器也能定位到正確的起點。CDC 連接器通常具備先快照再增量的能力。首次啟動時它會先讀取一次源表全量數據記錄當時的 binlog 位點之后繼續從該位點消費增量從而保證從啟動那一刻起不遺漏后續變更。2.4 從 binlog 到 BigQuery 的完整鏈路一個常見的生產架構是MySQL master - binlog - CDC Connector (Debezium / Flink CDC) - Kafka Topic - 流處理或寫入程序 - BigQuery Storage Write API / Load Job - BigQuery Table也可以簡化為MySQL master - Flink CDC - BigQuery 目標表無論采用哪種架構核心都是從日志讀取變化而不是定時查詢表。這也決定了后面的環境準備、配置、驗證和排錯方式。3. 前期準備MySQL、BigQuery 和權限一項都不能省3.1 版本與前置條件在配置 CDC 之前先確認環境是否滿足基本條件組件要求說明MySQL5.7 或 8.0開啟 binlog5.7 建議顯式開啟8.0 確認默認配置BigQuery數據集、目標表、服務賬號建議單獨建服務賬號避免共用管理員賬號CDC 工具Debezium 或 Flink CDC版本需要與 MySQL 和 Kafka 版本匹配網絡源庫與數倉側連通私網優先公網場景需要做好傳輸加密如果源 MySQL 是云數據庫還需要查看云廠商是否允許開啟 binlog 保留策略、是否開放復制賬號權限。有些托管數據庫默認不開放REPLICATION SLAVE這是接入 CDC 前最容易發現的阻塞點。注意開啟 binlog 并切換為 ROW 格式后binlog 日志量通常會變大磁盤占用和復制延遲都會上升。生產環境切換前需要評估磁盤余量。3.2 修改 MySQL 配置下面是一份最小可用的 MySQL CDC 配置示例[mysqld] server_id 1001 log_bin /var/log/mysql/mysql-bin.log binlog_format ROW binlog_row_image FULL expire_logs_days 7 # MySQL 8.0 可用以下參數控制 binlog 保留時長 # binlog_expire_logs_seconds 604800每個參數的作用server_idMySQL 實例在復制拓撲中的唯一標識。CDC 客戶端也會占用一個 server-id不能與主從庫中其他節點重復。log_bin開啟 binlog并指定日志文件路徑。binlog_formatROW讓 binlog 記錄行級變更。CDC 必須使用 ROW 格式。binlog_row_imageFULL讓 UPDATE 事件包含整行前鏡像和后鏡像。如果設置為 MINIMALbinlog 只包含被修改的字段和主鍵CDC 拿不到完整舊行和新行。expire_logs_days控制 binlog 文件保留天數。保留太短CDC 位點落后時可能追不上保留太長磁盤占用過大。常見建議是 3 到 7 天具體要結合源庫寫入量和磁盤容量調整。修改配置后需要重啟 MySQL。重啟前確認max_allowed_packet等參數不會限制大事務的 binlog 傳輸。3.3 創建 MySQL CDC 賬號建議為 CDC 單獨創建一個賬號避免使用 rootCREATE USER cdc_user% IDENTIFIED BY strong_password; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO cdc_user%; FLUSH PRIVILEGES;三個權限的含義SELECT用于 CDC 工具首次啟動時的全量快照以及讀取表結構信息。REPLICATION SLAVE允許該賬號通過復制協議讀取 binlog這是 CDC 的核心權限。REPLICATION CLIENT允許執行SHOW MASTER STATUS、SHOW BINARY LOG STATUS等命令用于確認位點信息。不要把ALL PRIVILEGES都授出去。CDC 賬號只需要讀取能力不需要寫源庫。3.4 BigQuery 側準備BigQuery 側需要準備數據集、目標表和服務賬號。在 Google Cloud Console 中先創建數據集例如analytics。目標表建議在接入 CDC 之前就定義好字段類型盡量與 MySQL 類型對應。如果后續依賴 BigQuery 自動加列容易遇到 schema 不一致導致寫入失敗。服務賬號需要授予 BigQuery Data Editor 或更細粒度的角色。把服務賬號的 JSON 密鑰下載到寫入服務所在機器并通過環境變量GOOGLE_APPLICATION_CREDENTIALS指向密鑰文件。BigQuery 是列式存儲目標表 schema 在寫入前就要對齊。CDC 事件字段如果比目標表多需要做過濾如果少目標表多出的列會使用默認值或 NULL。4. 最小落地鏈路Debezium 捕獲 binlog程序寫入 BigQuery4.1 兩種常用的技術選型常見方案有兩種方案鏈路適合場景Debezium KafkaMySQL - Debezium - Kafka - 寫入程序 - BigQuery已有 Kafka 基礎設施需要多消費方Flink CDCMySQL - Flink CDC - BigQuery Sink團隊熟悉 Flink希望用 SQL 處理流下面以 Debezium Kafka Python 消費者為例把鏈路拆開看。這樣更容易理解每個環節的職責。Flink CDC 只是把 Debezium 和流處理合并到一個框架里原理一致。4.2 Debezium connector 的配置Debezium 通過 Kafka Connect 運行一個典型配置如下{ name: mysql-orders-connector, config: { connector.class: io.debezium.connector.mysql.MySqlConnector, database.hostname: 10.0.0.10, database.port: 3306, database.user: cdc_user, database.password: xxxx, database.server.id: 5400, database.include.list: ecommerce, table.include.list: ecommerce.orders, database.history.kafka.bootstrap.servers: kafka:9092, database.history.kafka.topic: schema-changes.ecommerce, topic.prefix: mysql, include.schema.changes: true } }關鍵參數database.server.idDebezium 會占用一個 server-id。它必須與 MySQL 現有主從庫、其他 CDC 實例的 server-id 不沖突否則連接會被 MySQL 拒絕。database.include.list/table.include.list限定監聽的庫表。只同步需要的表能顯著減少 binlog 解析壓力。database.history.kafka.topicDebezium 用這個 topic 記錄表結構歷史。binlog 里的舊事件在解析時可能依賴歷史 schema因此這個 topic 不能隨意刪除。topic.prefix生成 Kafka topic 名稱的前綴。最終 topic 名稱一般是{topic.prefix}.{database}.{table}。4.3 變更事件長什么樣Debezium 輸出的變更事件是一段 JSON核心結構如下{ before: { id: 1001, status: pending }, after: { id: 1001, status: paid }, source: { db: ecommerce, table: orders, server_id: 1001, ts_ms: 1719900000123 }, op: u }op字段表示操作類型op 值含義事件內容cINSERT只有afteruUPDATE有before和afterdDELETE只有beforer快照讀取類似 INSERTafter為快照行注意DELETE 事件沒有after。寫入 BigQuery 時如果目標表要反映刪除必須自己定義刪除策略比如寫入一條帶刪除標記的記錄或者通過主鍵 MERGE 刪除目標行。4.4 寫入 BigQuery 的示例程序下面是一個最小 Python 消費者示例從 Kafka 讀取 MySQL 變更事件批量寫入 BigQueryimport json from google.cloud import bigquery from kafka import KafkaConsumer PROJECT my-project DATASET analytics TABLE orders client bigquery.Client(projectPROJECT) table_ref client.get_table(f{PROJECT}.{DATASET}.{TABLE}) def process_event(msg): payload json.loads(msg.value()) op payload.get(op) if op in (c, r): return payload[after] if op u: return payload[after] if op d: before payload[before] before[_is_deleted] True return before return None consumer KafkaConsumer( mysql.ecommerce.orders, bootstrap_serverskafka:9092, group_idbigquery-sync, auto_offset_resetlatest, enable_auto_commitFalse, ) rows [] batch_size 500 for message in consumer: row process_event(message) if row is not None: rows.append(row) if len(rows) batch_size: errors client.insert_rows_json(table_ref, rows) if not errors: consumer.commit() rows [] else: print(errors)這個示例說明的是思路不是完整生產代碼。insert_rows_json適合小規模驗證生產環境更推薦使用 BigQuery Storage Write API并配合監控、重試和死信隊列。enable_auto_commitFalse是為了避免消息未成功寫入就提交位點減少丟失風險但代價是重復消費因此目標表必須容忍重復。4.5 如果團隊已經用 Flink可以考慮 Flink CDCFlink CDC 可以把上面的鏈路壓縮成一個 SQL 和一套連接器。用 Flink SQL 創建 MySQL CDC 源表CREATE TABLE mysql_orders ( id INT, user_id INT, amount DECIMAL(10, 2), status STRING, updated_at TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 10.0.0.10, port 3306, username cdc_user, password xxxx, database-name ecommerce, table-name orders, server-id 5400-5406, scan.incremental.snapshot.enabled true );scan.incremental.snapshot.enabled在較新版本默認開啟。它讓 Flink CDC 以分片方式并行快照大表不需要像舊版本那樣先對全表加鎖再讀取對大表更友好。源表創建后可以再創建 BigQuery Sink 表通過INSERT INTO完成同步。具體 Sink 類名和參數取決于連接器版本落地前要以當前使用的 Flink 和連接器文檔為準。5. 怎么驗證 binlog 同步沒有漏數據5.1 先確認 binlog 真的開了進入 MySQL 命令行執行SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format; SHOW VARIABLES LIKE binlog_row_image;預期結果中log_bin為ONbinlog_format為ROWbinlog_row_image為FULL。還可以執行SHOW BINARY LOG STATUS;如果輸出包含當前 binlog 文件名和 position說明 binlog 文件正在正常寫入。5.2 驗證 Kafka 收到了哪些變更先用 Kafka 自帶的控制臺消費命令觀察 MySQL 變更是否進入 topickafka-console-consumer.sh \ --bootstrap-server kafka:9092 \ --topic mysql.ecommerce.orders \ --from-beginning然后在 MySQL 中分別執行一次 UPDATE 和一次 DELETEUPDATE orders SET status paid WHERE id 1001; DELETE FROM orders WHERE id 1002;正常情況下消費端會看到op為u和d的兩條事件。這一步直接驗證了周期同步最難做到的能力刪除和更新都能被捕獲。5.3 驗證 BigQuery 目標表觀察 BigQuery 目標表是否有新數據寫入。可以通過控制臺查詢也可以執行SELECT COUNT(*) FROM my-project.analytics.orders; SELECT MAX(updated_at) FROM my-project.analytics.orders;必須注意一個容易誤判的地方BigQuery 目標表不會因為收到了 DELETE 事件就自動刪除對應行。如果寫入程序只是把after或before以追加方式寫入刪除事件只會變成一行帶標記的數據。要真實反映刪除目標表需要按主鍵做 MERGE或者通過分區覆蓋實現。驗證時先明確自己的目標表語義是追加明細還是鏡像源表。5.4 延遲監控指標從 binlog 到 BigQuery 的同步不是一次性的必須持續監控。常見指標包括指標含義告警建議Kafka consumer lag消費程序落后的消息數持續增長則告警Debezium 位點與當前 binlog 的文件間隔連接器是否在追趕超過 binlog 保留期則高風險端到端延遲事件寫入 MySQL 到進入 BigQuery 的時間差根據業務要求設置閾值BigQuery 寫入錯誤率schema 不匹配等寫入失敗立即告警把位點落后和consumer lag 持續增長作為關鍵告警能提前發現大事務、網絡抖動或消費程序故障。6. 數據到達 BigQuery 后模式映射、DDL 和冪等才是真正的坑6.1 MySQL 與 BigQuery 類型映射字段類型映射是 CDC 鏈路里最容易踩坑的部分。下面是常見映射關系MySQL 類型BigQuery 類型注意事項INT / INTEGERINT64無符號 INT 可能超過 INT64 有符號范圍BIGINTINT64超過 2^63-1 的數據要改用 NUMERIC 或 STRINGDECIMAL(p, s)NUMERIC / BIGNUMERIC金額字段不要用 FLOAT精度會丟失DATETIMEDATETIME無時區語義按原值寫入TIMESTAMPTIMESTAMP建議統一按 UTC 存儲VARCHAR / TEXTSTRING長度和編碼要注意JSONJSONBigQuery 需要字段模式為 JSON 或先轉成 STRINGTINYINTINT64 / BOOL看業務語義確定最容易出問題的是 DECIMAL。MySQL 中的DECIMAL(10, 2)如果映射成 BigQuery 的 FLOAT640.1 這樣的值可能出現精度誤差。正確做法是映射為 NUMERIC。TIMESTAMP 也容易出問題。MySQL 的TIMESTAMP有會話時區概念CDC 事件里的ts_ms可能是 UTC 時間而業務字段本身可能是本地時間。建議在寫入端統一規范避免目標表同一列混入不同時區的數據。6.2 DDL 變更會打斷 CDC當 MySQL 表結構變化時CDC 鏈路會面臨兩個層面的問題。第一Debezium 需要依賴database.history.kafka.topic中的 schema 歷史來解析 binlog 里的舊事件。如果這個 topic 被刪除或清理連接器可能無法反序列化舊的 binlog 事件。第二BigQuery 目標表的 schema 不會自動跟隨 MySQL DDL 變化。MySQL 加了一列CDC 事件里出現了新字段但 BigQuery 目標表沒有這一列寫入就會報錯。處理建議是把 DDL 納入變更流程先審查 MySQL DDL 對同步鏈路的影響。先在 BigQuery 目標表補充或調整 schema。再在 MySQL 執行 ALTER TABLE。同步完成后核對事件是否正常。對于大表的 ALTER TABLE還可能導致源庫鎖表和復制延遲。生產環境做主從切換時要評估 DDL 對 binlog 位點的影響。注意不要依賴 BigQuery 自動加列來處理所有 DDL 變更。自動加列在不同版本和連接器里行為不一致且不能處理列重命名、刪除、類型變更等復雜操作。6.3 至少一次語義下重復是正常的binlog CDC 鏈路通常提供 at-least-once 語義。網絡閃斷、消費程序重啟、位點提交失敗都可能導致同一事件被重復消費。因此目標表必須能接受重復。常見做法按主鍵去重寫入前先判斷目標表是否已有該主鍵。使用 BigQuery MERGE按主鍵更新目標行。在記錄中增加事件版本字段如event_ts_ms或 GTID寫入時取較新的事件。下面是 BigQuery MERGE 的簡化思路MERGE my-project.analytics.orders AS t USING changes AS s ON t.id s.id WHEN MATCHED THEN UPDATE SET status s.status, amount s.amount WHEN NOT MATCHED THEN INSERT (id, user_id, amount, status, updated_at) VALUES (s.id, s.user_id, s.amount, s.status, s.updated_at);MERGE 在處理刪除事件時還可以加一個WHEN MATCHED AND s._is_deleted TRUE THEN DELETE分支。但 MERGE 的成本比流式追加高適合對一致性要求高、更新頻率可控的場景。如果表更新量極大需要考慮分區覆蓋、冷熱分離等方案。6.4 亂序事件怎么處理同一個主鍵的多條變更在 Kafka 中如果分布到不同分區消費程序收到的順序可能和源庫事務提交順序不一致。比如先提交了statuspaid后提交了statuscancelled亂序可能導致目標表最終停在paid。處理方式Kafka Topic 按主鍵 hash 分區保證同一主鍵路由到同一分區。寫入端使用 binlog 里的ts_ms或 GTID 做排序只接受更新的事件。如果業務允許短暫延遲可以在寫入端做窗口緩沖按主鍵排序后批量提交。如果源表存在刪主鍵后重新插入同一主鍵的場景還需要區分刪除后插入和舊 UPDATE 后到否則可能出現舊數據覆蓋新數據的現象。這種情況下GTID 或事務 ID 是更可靠的順序依據。7. 常見問題排查從現象倒推 binlog 鏈路故障7.1 現象連接器啟動時報權限不足或無法讀取 binlog可能原因MySQL 賬號缺少REPLICATION SLAVE權限。連接器配置的 server-id 與現有從庫沖突。binlog 未開啟或者binlog_format不是 ROW。排查命令SHOW VARIABLES LIKE binlog_format; SHOW GRANTS FOR cdc_user%; SHOW PROCESSLIST;處理方式核對 MySQL 配置和賬號權限修改后重啟連接器。server-id 沖突通常會在 MySQL 錯誤日志里看到A slave with the same server_uuid/server_id as this slave has connected to the master之類的信息。7.2 現象任務運行一段時間后Kafka 里有歷史事件但新事件遲遲不來可能原因Kafka Connect 或連接器進程掛掉后位點沒有正確恢復。MySQL 實例重啟導致 binlog 文件名變化連接器找不到舊位點對應的文件。table.include.list配置了大小寫敏感的表名實際表名大小寫不一致。排查方式kafka-consumer-groups.sh --bootstrap-server kafka:9092 --describe --group bigquery-sync重點看CURRENT-OFFSET、LOG-END-OFFSET和LAG。如果 consumer lag 為 0 但新數據沒進來檢查連接器日志里 binlog offset 是否還在推進。必要時做一次 重新快照 增量 的初始化。7.3 現象BigQuery 寫入報錯字段不存在或類型不匹配可能原因MySQL DDL 新增了列BigQuery schema 沒有同步。DECIMAL 字段映射成了 FLOAT64導致精度丟失或寫入失敗。MySQL JSON 字段映射到了 BigQuery STRING但事件里是 JSON 對象。排查方式SELECT column_name, data_type FROM my-project.analytics.INFORMATION_SCHEMA.COLUMNS WHERE table_name orders;處理方式定位是哪一列不匹配先同步 schema再重放失敗事件。不要直接丟棄報錯事件否則會在對賬時發現數據缺口。7.4 現象同步延遲持續增長可能原因源庫執行了大事務例如一次 UPDATE 超過十萬行binlog 事件量巨大。消費程序單線程寫入 BigQuery寫入速度跟不上源庫變更速度。網絡帶寬不足或者 BigQuery 寫入配額受限。處理方式在源庫側避免一次性更新超大范圍拆成小事務。寫入端改用批量并行寫并啟用 Storage Write API。增加監控觀察 binlog 保留時間是否充足。如果消費端位點落后太遠而 binlog 文件已經過期可能需要重新快照。