在 Azure Databricks 上,將生產環境的 Structured Streaming 工作負載作為排定的 Lakeflow Jobs 執行。 請參閱 Lakeflow 職位。
Databricks 建議你一定要設定以下事項:
- 從會傳回結果的筆記本中移除不必要的程式碼,例如
display和count。 - 不要用萬能運算來執行結構化串流工作負載。 務必使用作業運算,將串流排程為 Lakeflow Jobs。
- 使用
Continuous模式 排程 Lakeflow 作業。 這裡指的是Azure Databricks工作排程功能,而非結構化串流觸發間隔。 - 不要啟用結構化串流作業的自動縮放計算。
某些工作負載受益於下列各項:
Databricks 推出了 Lakeflow 管線,以降低結構化串流工作負載生產基礎設施管理的複雜性。 Databricks 建議使用 Lakeflow 管線來開發新的結構化串流管線。 參見 Spark 宣告式管線。
注意
在縮小結構化串流工作負載的叢集大小時,計算自動調整有其限制。 Databricks 建議在 Lakeflow 上使用 Spark 宣告式管線,並加強自動擴展功能來處理串流工作負載。 請參見「使用自動調整最佳化 Lakeflow 管道叢集使用率」。
:::note 無伺服器運算
在無伺服器運算中,僅支援 Trigger.AvailableNow() 和 Trigger.Once()。 Databricks 建議 Trigger.AvailableNow()。
在無伺服器運算中,以連續模式使用觸發模式與連續流水線模式。
請參閱 串流限制。
:::
降低營運串流的延遲
營運串流工作負載能近乎即時地接收、轉換並執行資料。 常見的例子包括詐欺偵測、異常偵測、個人化,以及即時監控與警示,這些延遲處理會直接影響業務成果。 這些工作負載的低延遲通常意味著數十到數百毫秒,儘管許多團隊會在幾秒級間設定服務水準協議(SLA),以考量較高百分位數的變異性。
為了達到最低端對端延遲,建議使用即時模式,該模式在尾端延遲低於一秒,常見情況下約為 300 毫秒。 詳見 即時模式概念。
當即時模式不適合你的工作負載時,以下最佳實務可降低微批次結構化串流的延遲:
- 輸出模式:在查詢運算子和接收端支援更新模式時,請使用更新模式。 更新模式會在每次觸發後輸出更新後的資料列,並在浮水印失效前持續更新這些資料列,因此請確保你的下游接收端具備冪等性,以處理更新後的結果。 對於更新模式不支援的工作負載,例如串流連結,或是可以丟棄遲到資料時,請使用附加模式。 不要用完整模式來降低延遲。 請參閱 選取結構化串流的輸出模式。
-
觸發器:使用
processingTime帶有0區間的觸發器,當前一個微批次結束且有新資料可用時,立即開始下一個微批次。 這能帶來最低的微批次延遲,但也會增加雲端儲存 API 的成本。 不要使用AvailableNow、 ,Once或Continuous用於操作負載。 請參閱《設定結構化串流觸發間隔》。 - 浮水印:將浮水印設定足夠長,包含晚到的資料,確保你的工作量不會下降。 浮水印控制查詢在丟棄亂序事件時間資料並清除狀態之前,會接受這類資料多久,因此過短的浮水印會悄悄丟棄有效的延遲到達記錄。 在此限制下,較短的浮水印降低延遲並保留較少狀態,較長的浮水印則能容忍較多延遲資料,但代價是延遲與狀態。 延遲 SLA 的小倍數(例如 2 倍)可作為調校的合理起點。 請參閱套用浮水印來控制資料處理閾值。
-
來源與接收端:可從低延遲來源讀取資料,例如訊息匯流排(Apache Kafka、Amazon Kinesis、Apache Pulsar 或 Google Cloud Pub/Sub),或讀取來自 Delta Lake 和 Apache Iceberg 資料表的變更資料饋送。 寫入至低延遲、高吞吐量的接收器,如訊息匯流排、操作資料庫或
foreach接收器。 將寫入端操作設計為冪等,使下游消費者能夠處理重複資料和延遲到達的資料。 - 狀態與檢查點:對於具狀態查詢,請使用 RocksDB 狀態存放區,因為變更日誌檢查點和非同步狀態檢查點都需要它。 啟用變更日誌檢查點功能,僅儲存增量狀態變更。 當狀態檢查點成為批次時間的瓶頸時,請啟用非同步狀態檢查點,讓檢查點寫入與下一個微批次重疊,並檢視其失敗復原及叢集調整大小的注意事項。 在耐用的雲端儲存中,為每個查詢設置自己的檢查點目錄。 請參閱 Azure Databricks 上的 Configure RocksDB state store、有狀態查詢的非同步狀態檢查點,以及結構化串流檢查點。
-
偏移管理:為減少連續串流偏移檢查點的延遲,啟用非同步進度追蹤,更新偏移與提交日誌而不阻擋資料處理。 它與
Once或AvailableNow觸發器不相容。 請參見非同步進度追蹤。 - 儲存體跳轉:盡可能將運算維持在單一串流管線內。 將邏輯分散到多個工作或管線會增加儲存跳數,增加延遲。
設計串流工作負載時,需考慮可能失敗的情況
Databricks 建議你一定要設定串流工作失敗時自動重啟。 部分功能,包括結構演化,要求結構化串流工作負載自動重試。 請參閱 設定結構化串流作業,以在失敗時重新啟動串流查詢。
某些運算,例如 foreachBatch,保證至少執行一次,而非精確執行一次。 對於這些操作,請確保你的處理流程是冪元的。 請參閱 使用 foreachBatch 寫入任意數據接收器。
注意
當查詢重新啟動時,會處理在先前執行中安排的微批次。 如果您的工作因記憶體不足錯誤而失敗,或因為微批次過大而手動取消工作,您可能需要提升計算資源以成功處理該微批次。
如果您在執行之間變更設定,這些設定將會套用至第一個新批次計畫。 請參閱結構化串流查詢變更後的復原。
當工作重試時
你可以在 Azure Databricks 工作中排程多個任務。 當您使用連續觸發器設定作業時,無法設定工作之間的相依性。
您可以選擇使用下列其中一種方法,在單一工作中安排多個資料流:
- 多個任務:定義一個作業,其中包含多個使用連續觸發器執行串流工作負載的任務。
- 多個查詢:在單一工作的原始程式碼中定義多個串流查詢。
您也可以合併這些策略。 下表比較這些方法。
| 策略 | 多項工作 | 多個查詢 |
|---|---|---|
| 如何共用計算? | Databricks 建議依據每個串流任務適當地分配計算資源。 您可以選用跨工作共用計算。 | 所有查詢都會共用相同的計算。 你可以選擇將查詢指派到 排程池。 |
| 重試如何處理? | 所有任務都必須在工作重試之前失敗。 | 如果有任何查詢失敗,工作會重試。 |
欲了解更多如何處理多個任務或查詢的細節,請參閱 「在同一叢集上執行多個結構化串流查詢」。
設定結構化串流工作以在失敗時重新啟動串流查詢
Databricks 建議你用連續觸發器來設定所有串流工作負載。 請參閱 持續執行作業。
連續觸發預設的行為如下:
- 防止同一工作多次並行運行。
- 在上一次執行失敗時開始新的執行。
- 使用指數退避進行重試。
Databricks 建議在排程工作流程時一律使用作業計算,而不是通用計算。 在工作失敗並重試時,會部署新的計算資源。
注意
Databricks 建議你不要使用 streamingQuery.awaitTermination() 或 spark.streams.awaitAnyTermination()。 請參閱 何時使用 awaitTermination()。
何時使用 awaitTermination()
streamingQuery.awaitTermination() 並 spark.streams.awaitAnyTermination() 封鎖目前執行緒,直到串流查詢終止。 是否使用這些函式取決於你的執行環境。
對於 Lakeflow 工作,請勿使用 streamingQuery.awaitTermination() 或 spark.streams.awaitAnyTermination()。 這些功能並非必要,因為 Jobs 服務會在串流查詢啟動時自動阻止執行完成。 這兩個功能都會阻擋筆記本儲存格的完成,並阻止工作服務追蹤串流查詢,這會干擾待辦事項的指標和工作通知。
以下情況的使用 awaitTermination() :
| 應用案例 | 行為 |
|---|---|
| 多功能運算上的互動筆記本 |
awaitTermination() 這樣可以保持儲存格運作,讓你能觀察查詢狀態,並確保錯誤顯示在筆記本的輸出中。 |
| 地方與開發環境 | 當本地執行 Spark 程式時,當主執行緒完成時,程序會退出。 呼叫 awaitTermination() 讓程式持續運作,直到串流查詢結束或失敗。 |
| 失效傳播至驅動器 | 若無 awaitTermination(),非工作上下文中的串流查詢失敗可能無法傳達至呼叫執行緒。 查詢可能會悄無聲息地失敗,使失敗更難被偵測和診斷。 呼叫 awaitTermination() 時,驅動程式又會觸發查詢異常。 |