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