監控並觀察自動裝填機

自動載入器管線需要主動監控,以偵測如積壓、結構漂移、資料損壞及串流停滯等問題,避免影響下游消費者。 本頁說明如何監控關鍵指標、查詢檔案層級狀態、建立可觀察性儀表板,以及排除常見問題。

關於生產設定細節,請參見 「設定生產工作負載的自動載入器」。 關於設定最佳實務,請參閱 自動載入器最佳實務。

先決條件

本頁的多項監控工作流程依賴 cloud_files_state() 觀察每個檔案的擷取狀態——包括待辦清單查詢、延遲計算及結構漂移偵測。 cloud_files_state() 是一個以表格為值的函式,會回傳自動載入器檢查點的檔案層級擷取狀態。 並非所有欄位預設皆可用。 可用性取決於您的 Databricks 執行版本與設定:

  • Databricks 執行環境 18.2 以上版本:discovery_time、processed_time 和 commit_time 會自動提供。 在 Databricks 執行環境 16.4–18.1 中,這些欄位僅在啟用時 cloudFiles.cleanSource 可用。
  • 已啟用 cloudFiles.cleanSource 的 Databricks 執行階段 16.4 及更新版本:archive_time、archive_mode 和 move_location可使用。

啟用 cloudFiles.cleanSource 會帶來一些效能負擔。 在正式環境中啟用之前,請先在預備生產環境中針對您的工作負載進行基準測試。

此外:

  • 用 _metadata 欄位註解被輸入的資料。 至少擷取 file_path 和 file_modification_time。 請參閱 檔案資料資料列。
  • 啟用 _rescued_data 與 _corrupt_record 欄位。

主要自動裝載機指標

下表彙整了 Auto Loader 管線需要監控的最重要指標。 這些指標可從 StreamingQueryListener 進度事件中取得,Auto Loader 專屬數值會 metrics 顯示在每個來源的地圖下方。

Metric 它告訴你的
numFilesOutstanding 積壓待處理的檔案數量
numBytesOutstanding 檔案積壓的大小(位元組)
approximateQueueSize 雲端佇列深度(僅檔案通知模式)
numInputRows 每批處理的行數
inputRowsPerSecond 資料到達率
processedRowsPerSecond 處理吞吐量
durationMs 細分 每批次所花費的時間

需注意的事項

以下情況表示您的管線可能需要留意。

  • 增加中 numFilesOutstanding:積壓量正在增加。 你的管線正在落後於輸入資料。
  • processedRowsPerSecond < inputRowsPerSecond:管線處理資料的速度比到達速度慢。
  • 大型 durationMs.latestOffset:檔案發現速度很慢。 考慮改用檔案事件。
  • 大型 durationMs.addBatch:資料處理速度緩慢。 考慮擴增運算資源或最佳化轉換作業。

完整指標參考,請參閱 Auto Loader 來源指標。

查詢檔案層級狀態 cloud_files_state

cloud_files_state()以表格值計算的函式會提供 Auto Loader 發現的每個檔案的詳細資訊。 以下欄位可用。 標示為需要 Databricks Runtime 16.4 以上或 18.2 以上的欄位,僅會在 先決條件 所述的條件下填入值。

Field 類型 Description
path STRING 檔案的路徑
size BIGINT 檔案大小 (以位元組為單位)
create_time TIMESTAMP 檔案建立的時間
discovery_time TIMESTAMP 當 Auto Loader 發現該檔案時(Databricks Runtime 16.4 以上版本)
processed_time TIMESTAMP 當 Auto Loader 處理檔案時(Databricks Runtime 16.4 及更新版本)
commit_time TIMESTAMP 當檔案被提交到檢查點(Databricks Runtime 16.4 及以上版本)時
archive_time TIMESTAMP 檔案被歸檔的時間(需要 cloudFiles.cleanSource)
archive_mode STRING MOVE、 DELETE,或 NULL (需要 cloudFiles.cleanSource)
move_location STRING 當 cloudFiles.cleanSource 為 MOVE 時的目的路徑
ingestion_state STRING 目前檔案匯入狀態

檢查檔案匯入狀態

以下查詢涵蓋常見的診斷情境。

查找所有未處理的檔案(目前積壓的檔案):

SELECT * FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state != 'COMMITTED';

計算平均擷取延遲(從建立檔案到提交的時間):

SELECT avg(unix_timestamp(commit_time) - unix_timestamp(create_time)) AS avg_latency_seconds
FROM cloud_files_state('path/to/checkpoint')
WHERE commit_time IS NOT NULL AND create_time IS NOT NULL;

查找損壞或跳過的檔案:

SELECT path, ingestion_state, size, create_time
FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state LIKE 'SKIPPED%';

追蹤封存進度(需要 cloudFiles.cleanSource):

SELECT archive_mode, count(*) AS file_count
FROM cloud_files_state('path/to/checkpoint')
GROUP BY archive_mode;

尋找從發現到提交延遲偏高的檔案,以找出瓶頸:

SELECT
  path,
  size,
  unix_timestamp(commit_time) - unix_timestamp(discovery_time) AS processing_latency_seconds,
  unix_timestamp(commit_time) - unix_timestamp(create_time) AS end_to_end_latency_seconds
FROM cloud_files_state('path/to/checkpoint')
WHERE commit_time IS NOT NULL
ORDER BY end_to_end_latency_seconds DESC
LIMIT 20;

完整 SQL 參考,請參見 cloud_files_state table-valued function。

監控 Lakeflow 管線中的自動裝載機

Databricks 建議對用於正式環境的 Auto Loader 管線使用 Lakeflow 管線。 為了善用其內建的監控功能:

  • 將 Lakeflow 管線事件日誌儲存在 Delta 表格中,以便查詢可觀察性資料。 透過管線的進階設定或 API 來設定。 詳情請參閱 管線事件日誌。

  • 設計你的管線以促進可觀察性。 Lakeflow 管線中,結構良好的 Auto Loader 管線包含{table}_source檢視表(Auto Loader 來源定義)、{table}_bronze串流資料表(用於擷取原始資料,並包含_rescued_data和_corrupt_record欄位)、corrupt_records_sink用於隔離含有無法剖析資料之資料列的資料表,以及{table}供下游使用的乾淨檢視表。

  • 在青銅層串流資料表上設定預期條件,以監控結構漂移和資料毀損。 _rescued_data IS NULL 偵測突發結構變更及 _corrupt_record IS NULL 無法解析的資料。 Lakeflow 管線會隨著資料的到來評估這些預期,並產生可觀察性軌跡。 你可以將預期設定為發出警告、捨棄資料列,或使管線失敗。

為您的管線建立 event_log_raw 檢視之後,請使用下列查詢來取得 Auto Loader 專用指標。

監控每個流程的資料擷取吞吐量:

SELECT
  origin.flow_name,
  origin.update_id,
  timestamp,
  TRY_CAST(details:flow_progress.metrics.num_output_rows AS BIGINT) AS rows_written
FROM event_log_raw
WHERE event_type = 'flow_progress'
ORDER BY timestamp DESC;

監視各流程的資料積壓情況:

SELECT
  origin.flow_name,
  timestamp,
  DOUBLE(details:flow_progress.metrics.backlog_bytes) AS backlog_bytes
FROM event_log_raw
WHERE event_type = 'flow_progress'
  AND details:flow_progress.metrics.backlog_bytes IS NOT NULL
ORDER BY timestamp DESC;

彙整期望違反情形,以偵測綱要漂移和損壞資料:

SELECT
  origin.flow_name,
  explode(from_json(
    details:flow_progress.data_quality.expectations,
    'array<struct<name:string, dataset:string, passed_records:bigint, failed_records:bigint>>'
  )) AS expectation
FROM event_log_raw
WHERE event_type = 'flow_progress'
  AND details:flow_progress.data_quality.expectations IS NOT NULL;

關於一般的 Lakeflow 管線監控指引,請參閱 「監控管線 」及 「管線事件日誌」。

使用結構化串流監控自動載入器

在 Lakeflow 管線外執行 Auto Loader 時,請使用以下結構化串流監控方法。

  • 實作一個 StreamingQueryListener,以便透過從 source.metrics 讀取資料,擷取每個批次的 Auto Loader 專屬指標。
from pyspark.sql.streaming import StreamingQueryListener

class AutoLoaderMonitor(StreamingQueryListener):
    def onQueryStarted(self, event):
        pass

    def onQueryProgress(self, event):
        for source in event.progress.sources:
            if "CloudFilesSource" in source.description:
                metrics = source.metrics
                files_outstanding = metrics.get("numFilesOutstanding", "0")
                bytes_outstanding = metrics.get("numBytesOutstanding", "0")
                rows_per_sec = source.processedRowsPerSecond
                # Push metrics to your monitoring system (for example, write to a Delta table)

    def onQueryIdle(self, event):
        pass

    def onQueryTerminated(self, event):
        pass

spark.streams.addListener(AutoLoaderMonitor())

Note

監聽器中的處理邏輯可能會減緩查詢處理速度。 限制監聽器回呼中的運算,並避免在其中進行同步的外部寫入;而是以非同步方式發出輕量的遙測資料,或將指標交由獨立的工作負責持久化。

  • 使用來源進度資訊中的 numInputRows、inputRowsPerSecond 和 processedRowsPerSecond 來計算吞吐量——亦即每個批次的每秒檔案數和每秒資料列數。

  • 若要計算擷取延遲,請比較來自 create_time 的 commit_time 與 cloud_files_state(),以計算端對端延遲。 對於處理延遲,請使用 durationMs 細分資訊(例如 latestOffset、addBatch 及其他已回報的批次階段)來找出哪個階段是瓶頸。

  • 使用 df.observe() 直接在串流 DataFrame 上定義內嵌資料品質指標。 指標可在 StreamingQueryListener 下的 observedMetrics 進度事件中看到。

from pyspark.sql.functions import count, lit, col

observed_df = df.observe(
    "auto_loader_quality",
    count(lit(1)).alias("total_rows"),
    count(col("_rescued_data")).alias("rescued_rows"),
    count(col("_corrupt_record")).alias("corrupt_rows")
)
  • 請為 .queryName() 每個串流設定唯一名稱,方便在 Spark UI 串流標籤和監控儀表板中辨識自動載入串流。

完整的結構化串流監控參考,請參見 Azure Databricks 上的結構化串流監控查詢。

建立可觀察性儀表板

整合多個來源的資料,建立完整的可觀察性儀表板,以支援您的 Auto Loader 管線。 此表展示了一些建議的來源,供你用來結構可觀察性儀表板。

數據源 可觀測性資料
cloud_files_state() 檔案層級的擷取狀態:每個檔案的發現、處理、提交及歸檔時間戳記
湖流量管線事件日誌 管線運行歷史、每批次流量指標及資料品質期望結果
管線輸出表 每個已擷取資料表所寫入的資料列數與資料量

接著你可以將可觀察性資料彙整成專用資料表,作為儀表板和警示的基礎:

  • 根據 event_type = 'update_progress' 事件,彙整管線隨時間變化的執行狀態摘要(成功或失敗)。
  • 彙整檔案擷取指標(積壓量、吞吐量、每批次延遲),衍生自 cloud_files_state() 和 event_type = 'flow_progress' 事件。
  • 利用事件日誌中的 num_output_rows 列數與資料量來建立資料表統計。
  • 從詳細的錯誤日誌和每次更新的預期違規中收集除錯資訊,這些都是根據 event_type = 'flow_progress' 事件 data_quality 填充的。

這些聚合的資料表能驅動 AI/BI 儀表板及 SQL 警示。 推薦的儀表板面板包括管線執行狀態時間軸、資料擷取待辦清單趨勢、吞吐量趨勢、資料擷取延遲分布、資料品質指標、結構演化事件及檔案歸檔狀態。

監控結構演化事件

請使用下列方法,在結構描述變更發生時加以偵測。

  • 預期違規計數中的 _rescued_data 若出現非 NULL 值,表示發生了綱要漂移。 查詢事件日誌 failed_records > 0 中的 no rescued data 期望值。
  • 對已設定的 _schemas 內的 cloudFiles.schemaLocation 目錄所做的變更(或者,僅在未另外設定結構描述位置時,在 checkpoint 內的變更)表示已發生結構描述演進。 你可以透過另一個監控作業輪詢此目錄。
  • 不要僅因為在相同串流名稱下出現一個 onQueryTerminated 事件,後面接著一個 onQueryStarted 事件,就將其視為結構描述演變的充分證據。 串流重啟的原因有很多(叢集重啟、程式碼部署、暫時性儲存錯誤)。 在判定已發生綱要演化之前,先將重新啟動與獨立訊號建立關聯——例如 _schemas 目錄變更或 _rescued_data 預期違反——。
  • 用 _metadata.file_path 來辨識哪些檔案引入了結構變更。 根據 cloud_files_state() 欄位將此項與 path 連接,以將綱要變更與特定檔案和批次建立關聯。

使用此範例查詢,透過期望條件違反來偵測最近的結構漂移:

SELECT
  timestamp,
  origin.flow_name,
  exp.name AS expectation_name,
  exp.failed_records
FROM (
  SELECT
    timestamp,
    origin,
    explode(from_json(
      details:flow_progress.data_quality.expectations,
      'array<struct<name:string, dataset:string, passed_records:bigint, failed_records:bigint>>'
    )) AS exp
  FROM event_log_raw
  WHERE event_type = 'flow_progress'
    AND details:flow_progress.data_quality.expectations IS NOT NULL
)
WHERE exp.name = '<rescued-data expectation name>'
  AND exp.failed_records > 0
ORDER BY timestamp DESC;

針對常見問題設置警示

使用 Databricks SQL 警示或管道通知,在問題影響下游消費者之前偵測問題。

以下 SQL 偵測日益增長的積壓,並可作為 Databricks SQL 警示的基礎。 設定它定期執行一次(例如每 5 分鐘一次),當結果非空時會提醒。

-- Alert when backlog exceeds threshold or trends upward across recent batches
WITH recent_backlog AS (
  SELECT
    origin.flow_name,
    timestamp,
    DOUBLE(details:flow_progress.metrics.backlog_bytes) AS backlog_bytes,
    ROW_NUMBER() OVER (PARTITION BY origin.flow_name ORDER BY timestamp DESC) AS rn
  FROM event_log_raw
  WHERE event_type = 'flow_progress'
    AND details:flow_progress.metrics.backlog_bytes IS NOT NULL
)
SELECT flow_name, backlog_bytes, timestamp
FROM recent_backlog
WHERE rn = 1
  AND backlog_bytes > 1073741824  -- alert when backlog exceeds 1 GB

下表總結了建議的警示條件:

要偵測什麼 如何偵測 何時發出警報
積壓案件日益嚴重 numFilesOutstanding 呈上升趨勢 在多個批次中持續增加
停滯溪流 沒有進度事件 N分鐘內無事件(根據預期觸發間隔)
高吞嚥延遲 commit_time - create_time 超過你的服務等級標準(SLA)門檻
資料品質退化 預期失敗率 未達預期的資料列比例增加
結構演化事件 _rescued_data IS NOT NULL 期望違反計數中的任何非 NULL 值
檔案發現速度慢 durationMs.latestOffset 明顯高於基準線

排除常見問題

下表說明常見的 Auto Loader 管線問題、可能原因及建議的解決措施。

Issue 可能原因 建議的動作
積壓的增長速度比處理速度快 運算能力不足、資料偏斜或節流的速率限制 擴充運算資源,使用 Spark UI 檢查資料傾斜,並檢查 maxFilesPerTrigger 設定以控制批次大小
找不到檔案 檔案事件設定錯誤、權限問題,或串流在 7 天內無法執行 確認外部位置權限,檢查 Unity Catalog UI 中檔案事件的設定,並確保串流至少每 7 天執行一次,以避免 RocksDB 狀態失效
直播啟動時間太慢 大型檢查點狀態下載(RocksDB) 升級至 Databricks Runtime 15.3 及以上版本以實現非同步狀態載入,啟動時間可減少約 90%
重複檔案處理 過於積極的 cloudFiles.maxFileAge 設定,或檢查點損毀 使用保守的 maxFileAge(至少 90 天),驗證檢查點完整性,並避免對檢查點儲存套用生命週期政策
結構演化導致管線重新啟動 頻繁或不相容的結構變更 檢閱 schemaEvolutionMode、切換至 addNewColumnsWithTypeWidening 以進行型別提升,或使用 Variant 類型來處理高度動態的結構描述
接收端累積損壞資料 資料來源品質問題 檢查 _corrupt_record 隔離接收端的模式特徵,檢視來源資料的產生方式,並考慮新增上游驗證
discovery_time 且 commit_time 無人居住 在 Databricks 執行環境低於 18.2 的情況下運行,且沒有 cleanSource 升級至 Databricks 執行環境 18.2 或以上版本,或在 Databricks 執行環境 16.4–18.1 啟用 cloudFiles.cleanSource

如需更多故障排除,請參閱 自動裝載機常見問題。