在工作流中執行管線

您可以使用 Lakeflow 作業、Apache Airflow 或 Azure Data Factory,作為資料處理工作流程的一部分來執行管線。

管線會自動解決資料集間的相依關係,因此它能自行處理簡單的管線內協調。 對於管線本身不具備的編排,例如條件執行、任務結果分支、重試,或將管線與其他類型的工作協調,請使用專用的工作流程編排器,而非將邏輯內建於管線中。

準備好您的管線,以利編排

當每個管線各自涵蓋你想要排程、驗證或獨立執行的明確工作單元時,編排最能發揮效果。 圍繞這些邊界設計你的管線,讓工作流程能將它們作為獨立任務協調,包括上游與下游任務間適當的控制流程。

如果你已經有一個大型管線,其中結合了你想要分開編排的工作,可以透過將資料表移至新管線,將其拆分為較小的管線。 請參閱 在管線之間移動資料表

Lakeflow 作業

你可以在 Lakeflow Jobs 中協調多個任務,實作資料處理工作流程。 若要在任務中包含管線,請在建立任務時使用 管線 任務。 請參閱 作業的管線工作

Apache Airflow

Apache Airflow 是一種用於管理和排程資料工作流程的開放原始碼解決方案。 Airflow 將工作流程表示為操作的有向無環圖 (DAG)。 您可以在 Python 檔案中定義工作流程,Airflow 會管理排程和執行。 如需搭配 Azure Databricks 安裝和使用 Airflow 的相關資訊,請參閱 使用 Apache Airflow 協調 Lakeflow 作業

若要在 Airflow 工作流程中運行管線,請使用 DatabricksSubmitRunOperator

需求

使用Lakeflow管線的氣流支援需具備以下條件:

Example

下列範例建立了一個 Airflow DAG,以觸發 ID 為 8279d543-063c-4d63-9926-dae38e35ce8b 的管線更新:

from airflow import DAG
from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator
from airflow.utils.dates import days_ago

default_args = {
  'owner': 'airflow'
}

with DAG('ldp',
         start_date=days_ago(2),
         schedule_interval="@once",
         default_args=default_args
         ) as dag:

  opr_run_now=DatabricksSubmitRunOperator(
    task_id='run_now',
    databricks_conn_id='CONNECTION_ID',
    pipeline_task={"pipeline_id": "8279d543-063c-4d63-9926-dae38e35ce8b"}
  )

CONNECTION_ID 替換為工作區中的 Airflow 連線 的識別符。

將此範例儲存在目錄中 airflow/dags ,並使用 Airflow UI 來 檢視和觸發 DAG。 使用管線 UI 來檢視管線更新的詳細資料。

Azure Data Factory

備註

Lakeflow pipelines 和 Azure Data Factory 各自都包含設定失敗時重試次數的選項。 如果在管線 呼叫管線的 Azure Data Factory 活動上設定重試值,則重試次數是 Azure Data Factory 重試值乘以管線重試值。

例如,如果管線更新失敗,管線預設最多可重試五次。 如果 Azure Data Factory 重試設定為 3,且管線使用預設值 5 次重試,則失敗的管線最多可能會重試 15 次。 若要避免在管線更新失敗時嘗試過多,Databricks 建議在設定管線或呼叫管線的 Azure Data Factory 活動時限制重試次數。

若要變更管線的重試組態,請在設定管線時使用設定 pipelines.numUpdateRetryAttempts

Azure Data Factory 是雲端式 ETL 服務,可讓您協調資料整合和轉換工作流程。 Azure Data Factory 直接支援在工作流程中執行 Azure Databricks 工作,包括 筆記本、JAR 工作和 Python 腳本。 您也可以從 Azure Data Factory Web 活動呼叫管線 REST API,在工作流程中包含管線。 例如,若要從 Azure Data Factory 觸發管線更新:

  1. 建立資料處理站 或開啟現有的資料處理站。

  2. 建立完成後,請開啟資料處理站的頁面,然後按一下 [開啟 Azure Data Factory Studio] 磚。 Azure Data Factory 使用者介面隨即出現。

  3. 在 Azure Data Factory Studio 使用者介面的 [新增] 下拉式功能表中選擇 [管線],以建立新的 Azure Data Factory 管線。

  4. [活動] 工具箱中,展開 [ 一般 ],然後將 [網頁 ] 活動拖曳至管線畫布。 按一下 「設定 」標籤,然後輸入下列值:

    備註

    作為安全性最佳做法,當您使用自動化工具、系統、指令碼和應用程式進行驗證時,Databricks 建議您使用屬於服務主體的個人存取權杖,而不是工作區使用者。 若要建立服務主體的令牌,請參閱 管理服務主體的令牌

    • URLhttps://<databricks-instance>/api/2.0/pipelines/<pipeline-id>/updates

      取代 <get-workspace-instance>

      取代 <pipeline-id> 為管線識別碼。

    • 方法:從下拉式功能表中選擇 POST

    • 標頭:按一下 + 新增。 在 「名稱 」文字方塊中,輸入 Authorization。 在 文字方塊中,輸入 Bearer <personal-access-token>

      請將 <personal-access-token> 替換為 Azure Databricks 的 個人存取權杖

    • 內文:若要傳遞其他要求參數,請輸入包含參數的 JSON 文件。 例如,若要啟動更新程序並重新處理整條管線的所有資料:{"full_refresh": "true"}。 如果沒有其他請求參數,請輸入空大括弧 ({})。

若要測試 Web 活動,請按一下 Data Factory UI 中管線工具列上的 [偵錯 ]。 執行的輸出和狀態 (包括錯誤) 會顯示在 Azure Data Factory 管線的 [輸出] 索引標籤中。 使用管線 UI 來檢視管線更新的詳細資料。

小提示

常見的工作流程需求是在完成前一個任務之後開始任務。 由於管線updates請求是非同步的,在開始更新後但尚未完成前返回,因此依賴管線更新的 Azure Data Factory 管線任務必須等待更新完成。 若要等待更新完成,可在觸發管線更新的 Web 活動之後新增 Until 活動。 在「Until」活動中:

  1. 新增 等待 活動 ,以等候設定的秒數以完成更新。
  2. 在 [等待] 活動之後新增 Web 活動,該活動會使用管線更新詳細資料要求來取得更新的狀態。 state回應中的欄位會傳回更新的目前狀態,包括更新是否已完成。
  3. 使用欄位的 state 值來設定「直到」活動的終止條件。 您也可以使用 設定變數活動,依據 state 的值來新增管線變數,並將此變數用於終止條件。