設計擷取邏輯與資料來源設定
設計攝取邏輯需要你做出關鍵決策,決定資料如何從來源系統傳送到 Azure Databricks。 這些決策影響資料的新鮮度、處理效率及系統可靠性。 在選擇工具或撰寫程式碼前,你需要了解資料來源的特性,並定義符合業務需求的擷取策略。
在本單元中,您將探討 Azure Databricks 中資料擷取邏輯與資料來源配置的關鍵設計考量。 重點在於概念邏輯與設計模式——具體工具選擇與實作細節會在後續單元中介紹。
了解萃取類型
擷取類型 決定了你如何從來源系統擷取資料,並影響資料管線的 時效 性與 效率 。 Azure Databricks 支援三種主要的擷取模式,每種模式都適合不同的情境。
完全抽取
完整擷取 會在每次擷取執行時,從來源系統讀取整個資料集。
這種方法在以下情況下效果良好:
- 來源系統不支援變更追蹤
- 資料量足夠小,可以有效率地重新處理
- 你需要定期從頭重建目的資料表
- 資料關係複雜,增量邏輯容易出錯
隨著資料量增加,完全擷取變得不切實際。 當只有幾千條變更時,重新處理數百萬筆紀錄會浪費計算資源並延遲下游消費者的處理。
增量萃取
增量擷取只會 處理自上次擷取以來新增或變更的紀錄。 引擎會追蹤已處理的資料,並只查詢更新。
此圖案要求:
- 來源資料(時間戳、序號或版本欄位)中的可靠變更指示器
- 狀態管理用來追蹤最後處理的位置
- 處理 遲到 或 亂序 紀錄的邏輯
增量擷取顯著提升大型資料集的效率。 每次執行可能只需要處理幾千次變更,而不是讀取數百萬個資料列。
串流擷取
串流擷取 透過持續處理提供 近乎即時 的資料擷取。 串流工作不再按排程運行,而是保持活躍狀態,並在記錄到來時處理。
考慮在以下情況下串流:
- 業務流程需要以秒數或分鐘計算的資料新鮮度
- 來源系統會發佈事件到訊息匯流排或變更摘要。
- 你需要對事件發生時做出反應
串流擷取需要始終連線的基礎設施,但能提供最低的延遲以提升資料可用性。
變更資料擷取設計
變更資料擷取 (CDC)是一種專門的擷取模式,能擷取個別列層級的變更,包括插入、更新與刪除。 CDC 與簡單的增量擷取不同,因為它會保留每次變更的操作類型。
在為 CDC 設計時,請考慮您的來源系統如何產生變更紀錄:
- 資料庫交易日誌:像 SQL Server 這類資料庫可利用其基於交易日誌的 CDC 功能,暴露列層級變更。 像 Azure Databricks 這類下游系統可以讀取這些變更,並將其作為擷取流程的一部分來處理。
- 變更資料串流:Delta 表格和部分資料庫如 Azure Cosmos DB 內建變更追蹤功能,能以元資料呈現列變更。
- 定期快照集:有些系統不支援連續 CDC,需要您比較快照集來確定變更。
CDC 能實現簡單的增量擷取無法支援的應用場景。 您可以準確地複製刪除、維護稽核線索以及處理修改關鍵資料行的更新。
備註
CDC要求排序資訊才能正確處理錯序事件。 你的來源資料需要一個欄位(例如時間戳或序號),來建立正確的變更順序。
識別原始檔案類型特徵
當資料以檔案形式而非來自交易系統時,檔案格式會影響你的擷取設計。 每種格式都有其獨特的特性,會影響處理過程。
基於文字的格式
CSV 和 JSON 檔案是常見的原始格式,因為它們容易被人類閱讀且廣泛支援。 然而,它們在效率處理上也帶來挑戰:
- 結構變異性:檔案可能有欄位順序不一致、缺失欄位或資料型別變化
- 無壓縮效率:文字格式通常比二進位替代方案佔用更多儲存空間
- 剖析額外負荷:在讀取過程中,每個值都需要進行文字到類型的轉換。
當來源系統只能產生這些格式,或資料量足夠小、效率問題不大時,請使用文字格式。
列式格式
Parquet 與 ORC 檔案以欄式格式儲存資料,使分析查詢更有效率。 這些格式提供:
- 壓縮:柱式儲存透過將相似的數值分組,能達到更好的壓縮比
- 欄位修剪:查詢只讀取所需的欄位,減少輸入輸出
- 嵌入式結構:檔案包含描述欄位名稱與類型的元資料
當你能控制上游系統如何匯出資料時,請設計你的擷取方式,讓它偏好欄位式原始格式。
半結構化格式
XML 和巢狀 JSON 檔案所代表的階層資料,並不適合用平面表格結構表示。 在攝取半結構化資料時:
- 規劃結構描述擷取或定義。
- 決定是要將結構壓平還是保留巢穴。
- 同時考慮同一來源內不同的文件結構。
資料來源連線考量
資料來源設定定義了 Azure Databricks 如何連接並驗證來源系統。 你的設計必須涵蓋連線安全、憑證管理與存取模式。
雲端物件儲存來源
雲端儲存體 (Azure Data Lake Storage、Amazon S3 或 Google Cloud Storage) 是最常見的擷取登陸區域。 在設定儲存連接時:
- 使用 Unity Catalog 的外部位置來管理存取儲存路徑。
- 配置適當的認證方法(管理身份、服務主體或儲存憑證)。
- 設計支援高效增量處理的資料夾結構。
資料夾結構大幅影響 Azure Databricks 辨識與處理新檔案的效率。 請考慮以下模式:
基於日期的分割:依日期階層組織檔案,以便在增量讀取時實現分割剪枝:
/data/sales/year=2025/month=11/day=28/
/data/sales/year=2025/month=11/day=29/
此結構允許查詢只掃描相關的日期分割區,而非列出所有檔案。
基於來源的組織:當從多個來源匯入資料時,應依來源系統分隔資料,以簡化處理邏輯:
/landing/erp/orders/
/landing/crm/customers/
/landing/web/clickstream/
處理狀態資料夾:使用獨立資料夾來追蹤檔案處理狀態:
/incoming/ # New files arrive here
/processing/ # Files currently being processed
/archive/ # Successfully processed files
/failed/ # Files that failed processing
設計摺疊結構時,請考慮:
- 檔案發現開銷:扁平結構中數千個檔案會拖慢列出作業。 分割結構減少了 Azure Databricks 必須列舉的檔案數量。
- 自動載入器相容性:Azure Databricks 自動載入器 能有效追蹤雲端儲存中的新檔案。 讓資料夾結構與自動載入器發現及檢查點檔案的方式保持一致。
-
分割區欄位擷取:使用Hive風格的分割(
key=value),讓Azure Databricks能自動將分割區值當作欄位擷取,而無需解析檔案內容。
資料庫來源
連接交易式資料庫需要與雲端儲存不同的考量:
- Azure Databricks 與資料庫伺服器之間的網路連線。
- 認證憑證被安全儲存(使用 秘密 而非硬編碼值)。
- 連線池 與 查詢下推 功能。
- 擷取過程中對來源系統效能的影響。
串流來源
訊息匯流與事件串流(如 Azure Event Hubs 或其他常見訊息平台)需要:
- 經紀商連線細節與認證。
- 取用者群組平行處理設定。
- 偏移管理 以追蹤消費狀況。
- 用於訊息格式驗證的結構登錄整合。
SaaS 應用來源
像 Salesforce、ServiceNow 或 SAP 這類企業應用程式,透過具有獨特特性的 API 揭露資料:
- 限制擷取吞吐量的 API 速率限制。
- 大型結果集的分頁策略。
- 認證流程(OAuth、API 金鑰或憑證式)。
- 可用的變更追蹤或 webhook 功能。
應用需求導向的設計框架
有效的擷取設計始於分析需求,再選擇擷取類型或檔案格式。 請使用以下可重複的框架來指導你的設計決策。
步驟 1:收集需求
在做出任何技術決策前,請先回答以下關鍵問題:
| 需求類別 | 待解問題 |
|---|---|
| 資料量 | 原始碼中存在多少資料? 兩次運行之間有多少變動? 成長率是多少? |
| 延遲需求 | 對下游消費者來說,數據必須有多新鮮? 有需要遵守的服務等級協議(SLA)嗎? |
| 原始碼功能 | 來源是否支援變更追蹤、CDC,還是僅支援完整擷取? 它能產生哪些檔案格式? |
| 資料品質 | 來源資料有多可靠? 需要什麼樣的肯定或淨化? |
| 網路限制 | 來源與 Azure Databricks 之間的頻寬是多少? 防火牆或 VPN 有需求嗎? |
| 成本影響 | 不同方法的運算、儲存和資料傳輸成本是多少? |
步驟二:將需求對應到提取類型
利用您的需求分析來選擇合適的萃取模式:
| 如果您的需求是... | 然後選擇...... |
|---|---|
| 資料量很小(< 10萬列),無法追蹤變更,簡單重建可接受 | 完全抽取 |
| 資料量龐大,且存在可靠的變更指標(時間戳/順序),且有足夠的小時或每日新鮮度 | 增量萃取 |
| 需要即時或近即時的新鮮度,來源會發布事件,並具備串流基礎設施 | 串流擷取 |
| 需要擷取刪除與操作類型,需要稽核追蹤,來源支援交易日誌擷取 | CDC 提取 |
步驟 3:選擇原始檔案格式
當你控制來源格式或需要選擇著陸區格式時,請將格式與你的需求相匹配:
| 如果你的情境涉及...... | 那麼建議選擇...... |
|---|---|
| 大型分析資料集、欄位查詢,需要壓縮 | Parquet |
| 與多種系統的互通性,需要人類可讀性 | CSV 或 JSON |
| 階層式或巢狀資料結構 | JSON 或 XML |
| 結構描述演進的高頻串流 | Avro 與結構註冊表 |
步驟四:驗證你的設計
在實施前,請確認你的設計能解決以下問題:
擷取類型符合延遲要求。 否則,下游消費者會收到過時的資料,可能導致錯誤的商業決策或違反服務等級協議(SLA)。
此方法可擴展以應對預期的數據成長。 若不行,現今運作良好的管線可能會失效,或隨著資料量增加而變得過於緩慢。
來源系統的能力支援你選擇的模式。 若不然,管線會在執行時失效,或需要昂貴的變通方法增加複雜度。
成本估算是在預算限制內。 如果沒有,這個方案在技術上可能合理,但在生產用途上的財務狀況可能不可持續。
管線計畫進行資料品質檢查。 若未達標,無效或損壞的資料會傳入下游系統,導致錯誤或錯誤分析。
應用框架:實務範例
電子商務平台需要接收訂單資料。 訂單會儲存在 PostgreSQL 資料庫中,且業務需要每小時報告。
步驟一 - 收集需求:
- 銷量:每日新增訂單 50,000 筆,總筆數 200 萬筆,年成長率 15%。
- 延遲:報告時可接受一小時新鮮度。
- 來源:PostgreSQL 支援基於時間戳記的查詢,但沒有原生的 CDC。
- 品質:來源資料於入錄時驗證;只需極少的清潔。
- 網路:透過 Azure Private Link 直接連接。
- 成本:盡量減少運算使用;儲存成本是次要的。
步驟 2 - 對應到擷取類型:
大容量資料,需要可靠的時間戳記資料行和每小時延遲要求 → 增量擷取。
為什麼不選其他的?
- 完整擷取:每小時重新處理 200 萬筆資料浪費運算資源,而每小時僅有約 2,000 筆訂單發生變更。
- 串流:一小時的延遲無法證明維持持續運作的串流基礎設施成本是合理的。
- CDC:PostgreSQL 不會暴露原生的交易日誌 CDC,且使用情境不需要追蹤刪除或操作類型。
步驟 3 - 選擇檔案格式:
資料會匯入 Azure Data Lake Storage 進行處理→ Parquet 則用於高效的欄式儲存。 用於增量發現的日期分區佈局→ /landing/orders/year=2025/month=11/day=28/。
為什麼不選其他的?
- CSV/JSON:檔案大小較大,沒有壓縮效益,且分析工作負載查詢效能較慢。
- Avro:更適合應用於具有結構演變的串流;批次擷取不需要串流導向的功能。
- XML:增加了解析的複雜度,但對結構化表格順序資料沒有幫助。
步驟四 - 驗證設計:
- ✓ 按小時增量,符合一小時延遲要求。
- ✓ 處理過程僅能隨著成長而有效率地擴大規模。
- ✓ PostgreSQL
order_updated_at欄位支援增量查詢。 - ✓ 處理數千列與數百萬列相比,能降低運算成本。
最終決定:使用 order_updated_at 時間戳記欄位進行增量擷取,按小時排程,並將資料以 Parquet 檔案形式放入 Azure Data Lake Storage。
小提示
記錄你的設計決策及其背後的理由。 當需求變更或出現問題時,這些文件能幫助你了解為何選擇目前的方法,以及做出了哪些取捨。 考慮為每條主要管線建立設計決策記錄。
基於這個以需求為導向的基礎,您已準備好評估實施您資料擷取設計的特定工具與載入方法。