このトピックでは、DMSNotebookOperator の設定について説明します。
概要
DMS で管理されているノートブックファイル (.ipynb) を実行します。
前提条件
セッションを再利用する場合は、ノートブックセッションが作成されていることを確認してください。
新しいセッションを作成する場合は、ノートブックセッションの設定を含むテンプレートが必要です。
パラメーター
file_path、run_params、profile、session_id、profile_id、cluster_id、session_name、profile_name、cluster_name の各パラメーターは Jinja テンプレートをサポートします。
パラメーター | 型 | 必須 | 説明 |
file_path | string | はい | ノートブックファイル (.ipynb) のパスを指定します。 |
profile | dict | いいえ | ノートブックセッションのプロファイルを指定します。セッションを新規作成する場合、
|
profile_id | string | いいえ 説明 このパラメーターは、セッションを再利用しない場合に必須です。 |
|
profile_name | string | ||
cluster_type | string | DMS ワークスペース内のコンピューティングクラスタータイプを指定します。有効な値は次のとおりです:
| |
cluster_id | string |
説明 いずれか一方を指定する必要があります。 | |
cluster_name | string | ||
spec | string | ドライバーのリソース仕様を指定します。有効な値は次のとおりです:
| |
runtime_name | string | イメージ名を指定します。 | |
session_id | string | いいえ 説明 このパラメーターは、セッションを再利用する場合に必須です。 | 再利用するセッションを指定します。既存のセッションを再利用する場合、
説明 いずれか一方を指定する必要があります。 |
session_name | string | ||
run_params | dict | いいえ | ノートブックファイル内の変数を置き換えるランタイムパラメーターを指定します。 |
timeout | int | いいえ | ノートブックファイルの実行タイムアウト期間を秒単位で指定します。 |
polling_interval | int | いいえ | 実行ステータスをポーリングする間隔を秒単位で指定します。デフォルト値は 10 です。このパラメーターを 0 または負の数に設定すると、結果を待たずにタスクが送信されます。ポーリングには組み込みの再試行メカニズムが含まれます。 |
例
task_id と dag は Airflow の特定のパラメーターです。詳細については、「Airflow 公式ドキュメント」をご参照ください。
from airflow import DAG
from airflow.decorators import task
from airflow.models.param import Param
from airflow.operators.bash import BashOperator
from airflow.operators.empty import EmptyOperator
import json
from airflow.providers.alibabadms.cloud.operators.dms_notebook import DMSNotebookOperator
with DAG(
"dms_notebook_test",
params={
"x":3
},
) as dag:
notebook_operator = DMSNotebookOperator(
task_id='notebook_test_hz_name',
profile_name='hansheng_profile.48',
profile={},
cluster_type='spark',
cluster_name='spark_general2.218',
spec='4C32G',
runtime_name='Spark3.5_Scala2.12_Python3.9_General:1.0.9',
file_path='/Workspace/code/default/test.ipynb',
run_params={
'a':"{{ params.x }}"
},
polling_interval=5,
dag=dag
)
run_this_last = EmptyOperator(
task_id="run_this_last",
dag=dag,
)
notebook_operator >> run_this_last
if __name__ == "__main__":
dag.test(
run_conf={}
)すべての DMS Airflow オペレーターは、タスクのキャンセルや自動再試行などの共通機能をサポートします。詳細については、「Airflow DMS Operator」をご参照ください。