自動載入器管線需要主動監控,以偵測如積壓、結構漂移、資料損壞及串流停滯等問題,避免影響下游消費者。 本頁說明如何監控關鍵指標、查詢檔案層級狀態、建立可觀察性儀表板,以及排除常見問題。
關於生產設定細節,請參見 「設定生產工作負載的自動載入器」。 關於設定最佳實務,請參閱 自動載入器最佳實務。
先決條件
本頁的多項監控工作流程依賴 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 |
如需更多故障排除,請參閱 自動裝載機常見問題。