結構化串流的生產考量

在 Azure Databricks 上,將生產環境的 Structured Streaming 工作負載作為排定的 Lakeflow Jobs 執行。 請參閱 Lakeflow 職位

Databricks 建議你一定要設定以下事項:

  • 從會傳回結果的筆記本中移除不必要的程式碼,例如 displaycount
  • 不要用萬能運算來執行結構化串流工作負載。 務必使用作業運算,將串流排程為 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() 時,驅動程式又會觸發查詢異常。