在 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()。
在無伺服器運算中,以連續模式使用觸發模式與連續流水線模式。
請參閱 串流限制。
:::
設計串流工作負載時,需考慮可能失敗的情況
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() 時,驅動程式又會觸發查詢異常。 |