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