このトピックでは、DMSLindormSparkOperator のパラメーターについて説明します。
概要
このオペレーターは、Lindorm Spark エンジンでタスクを実行します。SQL と JAR の 2 つのジョブタイプをサポートしています。
パラメーター
instance、sql、region、configs パラメーターは Jinja テンプレート をサポートしています。
パラメーター | タイプ | 必須 | 説明 |
instance | 文字列 | はい | DMS で設定した Lindorm インスタンスへの接続 (DBLink) の名前です。 |
job_type | 文字列 | いいえ | ジョブタイプです。有効な値は次のとおりです:
|
sql | 文字列 | いいえ | 実行する Spark SQL ステートメントです。 |
region | 文字列 | いいえ | Lindorm インスタンスのリージョン ID です。クロスリージョン呼び出しでは必須です。 |
configs | dict | いいえ | 設定オブジェクトです。
|
polling_interval | int | いいえ | タスクステータスを確認するためのポーリング間隔 (秒) です。デフォルト値は 10 です。0 または負の数に設定した場合、オペレーターはタスクをサブミットしますが、完了を待機しません。ポーリング処理には、リトライメカニズムが組み込まれています。 |
例
task_id と dag は Airflow の標準パラメーターです。詳細については、Airflow の公式ドキュメントをご参照ください。
SQL ジョブ
from airflow import DAG
from airflow.providers.alibabadms.cloud.operators.dms_lindorm_spark import DMSLindormSparkOperator
with DAG(
"dms_lindorm_spark_sql",
) as dag:
lindorm_sql = DMSLindormSparkOperator(
task_id="lindorm_spark_sql",
instance="my_lindorm_link",
job_type="sql",
sql="SELECT * FROM wide_table WHERE dt = '{{ ds }}' LIMIT 1000;",
polling_interval=15,
dag=dag,
)JAR ジョブ
from airflow import DAG
from airflow.providers.alibabadms.cloud.operators.dms_lindorm_spark import DMSLindormSparkOperator
with DAG(
"dms_lindorm_spark_jar",
) as dag:
lindorm_jar = DMSLindormSparkOperator(
task_id="lindorm_spark_jar",
instance="my_lindorm_link",
job_type="jar",
configs={
"mainClass": "com.example.LindormETL",
"mainResource": "oss://my-bucket/jars/lindorm-etl.jar",
"args": ["--date", "{{ ds }}"],
"configs": {
"spark.executor.memory": "4g",
"spark.executor.instances": "4",
},
},
dag=dag,
)すべての DMS Airflow オペレーターは、タスクのキャンセルや自動リトライなどの共通機能をサポートしています。詳細については、Airflow DMS Operator のドキュメントをご参照ください。