用自動載入器引入資料
當新的資料檔案持續進入雲端儲存時,你需要一種有效率的方式來處理它們,而不需手動追蹤哪些檔案已被吞入。 Auto Loader 透過在新檔案出現時自動偵測並擷取,提供確實一次的保證,無需您手動管理狀態或檢查點,藉此解決此挑戰。
了解自動裝填機的運作方式
自動載入器是一種結構化串流來源,負責監控雲端儲存位置並逐步處理新檔案。 自動載入器不會重複掃描整個目錄,而是會維持已處理的檔案狀態,確保每個檔案只被擷取一次。
當你啟動自動載入器串流時,它可以處理目錄中現有的檔案,並持續監控新檔案的到來。 串流會將進度資訊儲存在檢查點位置,若中斷,則能從停止處繼續。
自動載入器會透過兩種模式之一偵測新檔案:
- 目錄列表模式:自動載入器會定期列出輸入目錄以發現新檔案。 此方法除了儲存存取外,不需額外設定。
- 檔案通知模式:自動載入器使用雲端通知服務,接收新檔案抵達時的事件。 此模式對大型工作負載更有效率,因為它避免了重複的目錄列表。
對於大多數生產工作負載,Databricks 建議在 Unity Catalog 的外部位置啟用檔案通知模式並開啟檔案事件功能。 透過檔案事件,自動載入器會在檔案抵達時直接接收通知,降低延遲與雲端 API 成本。
設定 Auto Loader 以載入檔案
要使用自動載入器(Auto Loader)來擷取資料,請用 cloudFiles 與 spark.readStream 格式。 以下範例是從 Azure Data Lake Storage 讀取 JSON 檔案,並將其寫入 Unity 目錄資料表:
base_path = "abfss://container@storage.dfs.core.windows.net/autoloader/orders"
schema_path = f"{base_path}/schema"
checkpoint_path = f"{base_path}/checkpoint"
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", schema_path)
.load("abfss://container@storage.dfs.core.windows.net/incoming/orders/")
.writeStream
.option("checkpointLocation", checkpoint_path)
.trigger(availableNow=True)
.toTable("sales.bronze.orders"))
該 cloudFiles.format 選項指定原始檔案格式。 自動載入器支援 JSON、CSV、Parquet、Avro、ORC、XML、文字及二進位檔案。 它 cloudFiles.schemaLocation 儲存推斷出的結構,使自動載入器能追蹤結構隨時間的變化。 雖然結構位置和檢查點位置可以相同,但保持分開可以讓你重置串流而不失去推斷的結構。
設定會處理所有目前可用的檔案,然後停止。 這種模式對於排程批次工作效果很好。 對於連續處理,你可以省略觸發器,或使用trigger(processingTime="1 minute")來間隔處理檔案。
搭配 read_files 使用 SQL 語法
Auto Loader 也支援透過 read_files table-valued 函式的 SQL 語法。 這種方法適合你偏好 SQL 或使用 SQL 筆記本時:
SELECT * FROM read_files(
'abfss://container@storage.dfs.core.windows.net/incoming/orders/',
format => 'json',
schemaHints => 'order_id INT, amount DECIMAL(10,2)'
)
在 Lakeflow pipelines 中使用 read_files 串流資料表時,自動載入功能會自動啟用。 使用 STREAM 關鍵字來啟用增量處理:
CREATE OR REFRESH STREAMING TABLE bronze_orders
AS SELECT * FROM STREAM read_files(
'abfss://container@storage.dfs.core.windows.net/incoming/orders/',
format => 'json'
)
在此情境下,管線自動管理檢查點與結構演化。
小提示
使用 Unity Catalog 時,請將檢查點和結構位置存放在受管理的儲存空間。 這確保治理一致,避免巢狀路徑帶來的權限問題。
處理模式推論與演化
自動載入器能自動偵測來源檔案的結構,免除手動定義結構的需求。 啟用結構推論時,自動載入器會取樣檔案以決定欄位名稱與資料型態。
對於 JSON 和 CSV 檔案,Auto Loader 預設會推斷所有欄位為字串。 這種保守的做法能防止類型不匹配,避免導致資料遺失。 要啟用自動型別偵測,請設定以下 cloudFiles.inferColumnTypes 選項:
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.option("cloudFiles.schemaLocation", checkpoint_path)
.load(source_path))
當你的來源資料新增欄位時,自動載入器會偵測變更並自動演化目標結構。 該 cloudFiles.schemaEvolutionMode 選項控制此行為:
| Mode | 行為 |
|---|---|
addNewColumns |
串流失敗。 新的欄位會新增至結構。 重新啟動後,串流會以更新的結構繼續。 當沒有提供結構時,這是預設的。 |
addNewColumnsWithTypeWidening |
串流失敗。 新欄位會加入到結構描述中,支援的資料型別變更(例如從 INT 變更為 LONG)會自動擴展。 未支援的類別更改被記錄在 _rescued_data中。 需要 Databricks 運行環境 16.4+。 |
rescue |
結構從不演化,串流也不會失敗。 所有新增欄位都會被擷取在該 _rescued_data 欄位中。 |
failOnNewColumns |
串流失敗。 除非所提供的結構更新或移除有問題的資料檔案,否則串流不會重新啟動。 |
none |
架構不會演化,新欄位會被忽略,資料也不會被救援,除非設定了 rescuedDataColumn 這個選項。 這是提供結構時的預設。 |
救援資料欄位(_rescued_data)捕捉任何與預期模式不符的資料,包括型別不匹配和意外欄位。 此功能可防止來源資料不符合預期時資料遺失。 自動載入器在推斷結構時會自動新增此欄位。
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.option("cloudFiles.schemaEvolutionMode", "rescue")
.schema(expected_schema)
.load(source_path))
要重新命名已救回的資料欄位,請使用 rescuedDataColumn 以下選項:
.option("rescuedDataColumn", "my_rescued_data")
如果你知道某些欄位應該有特定類型,請使用 結構提示 來覆蓋推斷的類型:
.option("cloudFiles.schemaHints", "order_date DATE, amount DECIMAL(10,2)")
監視並最佳化 Auto Loader 串流
對於生產工作負載,監控有助於追蹤資料匯入進度並找出問題。 Auto Loader 提供一個 SQL 函式來查詢發現檔案的狀態:
SELECT * FROM cloud_files_state('/path/to/checkpoint');
此查詢會回傳 Auto Loader 發現的檔案的元資料,包括其處理狀態。 你可以利用這些資訊來驗證檔案是否如預期被匯入。
Auto Loader 透過結構化串流儀表板回報進度指標,包括:
-
numFilesOutstanding:已發現但尚未處理的檔案 -
numBytesOutstanding: 待處理的總位元組數
對於高容量工作負載,請使用 cloudFiles.maxFilesPerTrigger 或 cloudFiles.maxBytesPerTrigger來控制處理速率。 這些設定可防止單一微批次過大:
.option("cloudFiles.maxFilesPerTrigger", "1000")
.option("cloudFiles.maxBytesPerTrigger", "1g")
這很重要
檢查點位置儲存關鍵狀態資訊。 避免套用可能刪除檢查點檔案的雲端儲存生命週期政策,因為這會破壞串流狀態,且需要從頭重啟。
要存取您所接收資料中的檔案層元資料,請選擇欄位 _metadata 。 此隱藏欄位包含每個來源檔案的資訊,包括 file_path、 file_name、 file_size及 file_modification_time:
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "parquet")
.load(source_path)
.select("*", "_metadata.file_path", "_metadata.file_modification_time"))
Auto Loader 與 Lakeflow 管線無縫整合,能自動處理檢查點與結構管理。 在下一單元,你將學習如何在宣告式管線定義中使用自動載入器,以進行端到端的資料擷取工作流程。