このトピックでは、DMSAnalyticDBSparkOperator のパラメーターについて説明します。
概要
DMSAnalyticDBSparkOperator は、AnalyticDB Spark 用の汎用オペレーターです。Warehouse と Job の両方のリソースグループタイプ、および SQL、STREAMING、BATCH の各アプリケーションタイプをサポートしています。詳細については、「アプリケーションシナリオ」をご参照ください。
パラメーター
パラメーター sql、conf、cluster_id、resource_group、instance、schema、app_name、app_type、および execute_time_limit_in_seconds は、Jinja テンプレートを使用できます。
パラメーター | タイプ | 必須 | 説明 |
cluster_id | string | はい |
説明
|
instance | string | ||
resource_group | string | はい | AnalyticDB for MySQL クラスターのリソースグループの名前。 |
region | string | いいえ | リージョン ID。このパラメーターは、クロスリージョンコールに必要です。 |
app_name | string | いいえ | アプリケーション名。このパラメーターは、リソースグループタイプが Job の場合に使用されます。 |
app_type | string | いいえ | アプリケーションタイプ。このパラメーターは、リソースグループタイプが Job の場合に使用されます。有効値:
|
sql | string | はい | 実行する Spark SQL ステートメント。 |
conf | dict | いいえ | カスタム Spark 設定。詳細については、「ExecuteSparkWarehouseBatchSQL」、「GetSparkWarehouseBatchSQL」、および「CancelSparkWarehouseBatchSQL」をご参照ください。 説明
|
schema | string | いいえ | 使用するスキーマ。デフォルトは |
polling_interval | int | いいえ | 実行ステータスのポーリング間隔 (秒単位)。デフォルトは 10 です。0 または負の数に設定すると、オペレーターはタスクを送信し、完了を待たずに終了します。ステータスポーリングには、組み込みのリトライメカニズムが含まれています。 |
execute_time_limit_in_seconds | int | いいえ | 実行タイムアウト (秒単位)。デフォルトは 36000 (10 時間) です。 |
callback | function | いいえ | SQL の結果を処理するコールバック関数。入力パラメーターは SparkBatchSQL です。 説明 このパラメーターは、 |
例
task_id と dag は Airflow の特定のパラメーターです。詳細については、「Airflow 公式ドキュメント」をご参照ください。
from typing import Any
from airflow import DAG
from airflow.decorators import task
from airflow.models.param import Param
from airflow.operators.empty import EmptyOperator
from datetime import datetime
from airflow.providers.alibabadms.cloud.operators.dms_analyticdb_spark import DMSAnalyticDBSparkOperator
# alibabacloud_adb20211201.models.SparkBatchSQL オブジェクト
def print_result(result):
print(f"{result}")
with DAG(
"dms_adb_spark_dblink",
params={
"sql": "show databases;show databases;"
}
) as dag:
warehouse_operator: Any = DMSAnalyticDBSparkOperator(
task_id="warehouse_sql",
instance="dbl_adbmysql_89",
resource_group="hansheng_spark_test",
sql="{{ params.sql }}",
polling_interval=5,
conf={
'spark.adb.sqlOutputFormat':'CSV',
'spark.adb.sqlOutputPartitions':1,
'spark.adb.sqlOutputLocation':'oss://hansheng-bj/airflow/adb_spark/test',
'sep':'|'
},
callback=print_result,
dag=dag
)
run_this_last = EmptyOperator(
task_id="run_this_last",
dag=dag
)
warehouse_operator >> run_this_last
if __name__ == "__main__":
dag.test(
run_conf={}
)すべての DMS Airflow オペレーターは、タスクのキャンセルや自動リトライなどの共通機能をサポートしています。詳細については、「Airflow DMS Operator」をご参照ください。