Spark on MaxCompute は、ローカルモードとクラスターモードの両方でジョブの実行をサポートしています。DataWorks では、クラスターモードで Spark on MaxCompute のオフラインジョブを実行し、他のタイプのノードと統合してスケジューリングできます。このドキュメントでは、DataWorks を使用して Spark on MaxCompute ジョブを設定およびスケジューリングする方法について説明します。
概要
Spark on MaxCompute は、MaxCompute が提供するオープンソースの Spark と互換性のあるコンピューティングサービスです。統一されたコンピューティングリソースと権限システムに基づいて Spark コンピューティングフレームワークを提供します。これにより、使い慣れた開発ワークフローを使用して Spark ジョブを送信・実行し、幅広いデータ処理および分析ニーズに対応できます。DataWorks では、MaxCompute Spark ノードを使用して Spark on MaxCompute ジョブをスケジューリング・実行し、他のジョブと統合できます。
Spark on MaxCompute は、Java、Scala、Python での開発をサポートし、ローカルモードまたはクラスターモードでジョブを実行します。DataWorks でオフラインの Spark on MaxCompute ジョブを実行する場合、クラスターモードで実行されます。Spark on MaxCompute の実行モードの詳細については、「実行モード」をご参照ください。
権限要件
ジョブを開発するには、RAM ユーザーを対応するワークスペースに追加し、開発者 または ワークスペース管理者 のいずれかのロールを付与する必要があります。ワークスペース管理者ロールには広範な権限が含まれているため、慎重に付与する必要があります。ワークスペースにメンバーを追加する方法の詳細については、「ワークスペースメンバーの追加」をご参照ください。
Alibaba Cloud アカウントを使用している場合は、この手順をスキップできます。
制限事項
Spark 3.x を使用する MaxCompute Spark ノードの送信時にエラーが発生した場合は、サーバーレスリソースグループを購入して使用する必要があります。詳細については、「サーバーレスリソースグループの使用」をご参照ください。
事前準備
MaxCompute Spark ノードは、Java/Scala および Python を使用した Spark on MaxCompute のオフラインジョブの実行をサポートします。開発手順と設定 UI は言語ごとに異なります。ビジネス要件に基づいて言語を選択してください。
Java/Scala
MaxCompute Spark ノードで Java または Scala コードを実行するには、事前に Spark on MaxCompute ジョブコードをローカルで開発し、MaxCompute リソースとして DataWorks にアップロードしておく必要があります。次の手順に従ってください:
-
開発環境をセットアップします。
お使いのオペレーティングシステムに基づいて、Spark on MaxCompute ジョブを実行するための開発環境を準備します。詳細については、「Linux 開発環境のセットアップ」および「Windows 開発環境のセットアップ」をご参照ください。
-
Java/Scala コードを開発します。
MaxCompute Spark ノードでコードを実行する前に、ローカルまたは既存の環境で Spark on MaxCompute コードを開発します。Spark on MaxCompute が提供するサンプルプロジェクトテンプレートの使用を推奨します。
-
コードをパッケージ化し、DataWorks にアップロードします。
開発が完了したら、コードをパッケージ化し、MaxCompute リソースとして DataWorks にアップロードします。詳細については、「リソース管理」をご参照ください。
Python (デフォルト環境)
DataWorks では、Python リソースに直接コードを記述して PySpark ジョブを開発できます。その後、MaxCompute Spark ノードを使用してコードを送信および実行できます。開発例については、「PySpark 開発例」をご参照ください。
デフォルト環境がジョブの依存関係要件を満たさない場合は、「Python (カスタム環境を使用)」セクションを参照して、カスタム Python 環境を準備してください。または、Python リソースのサポートがより充実している PyODPS 2 ノードまたは PyODPS 3 ノードを使用することもできます。
Python (カスタム環境)
デフォルトの Python 環境がビジネス要件を満たさない場合は、次の手順に従ってカスタム Python 環境を使用して Spark on MaxCompute ジョブを実行します。
-
Python 環境をローカルで準備します。
「PySpark の Python バージョンと依存関係のサポート」を参照して、必要な Python 環境を設定します。
-
環境をパッケージ化し、DataWorks にアップロードします。
Python 環境を .zip パッケージに圧縮し、MaxCompute リソースとして DataWorks にアップロードします。このパッケージが Spark on MaxCompute ジョブの実行環境になります。
パラメーター
DataWorks は、クラスターモードで MaxCompute 上の Spark オフラインジョブを実行します。クラスターモードでは、カスタムプログラムのエントリーポイント main を指定する必要があります。対応する Spark ジョブは、main 関数が Success または Fail のステータスで終了すると、終了します。さらに、spark-defaults.conf 内の設定を、MaxCompute Spark ノードの設定に 1 つずつ追加する必要があります。たとえば、executor の数、メモリサイズ、spark.hadoop.odps.runtime.end.point 設定などです。
spark-defaults.conf ファイルをアップロードする必要はありません。代わりに、spark-defaults.conf ファイル内の設定を 1 つずつ MaxCompute Spark ノードの設定項目に追加する必要があります。
Java/Scala
|
パラメーター |
説明 |
spark-submit コマンド |
|
[Spark バージョン] |
Spark のバージョン。有効な値は、Spark 1.x、Spark 2.x、Spark 3.x です。 説明
Spark 3.x バージョンを使用する MaxCompute Spark ノードの送信時にエラーが発生した場合は、サーバーレスリソースグループを購入して使用する必要があります。詳細については、「サーバーレスリソースグループの使用」をご参照ください。 |
— |
|
[言語] |
プログラミング言語。Spark on MaxCompute ジョブの開発に使用した言語に基づいて Java/Scala または [Python] を選択してください。 |
— |
|
[メイン jar リソースの選択] |
ジョブのメイン JAR リソースファイルを指定してください。 リソースファイルは DataWorks にアップロードしてコミットする必要があります。詳細については、「リソース管理」をご参照ください。 |
|
|
[設定項目] |
ジョブを送信するための設定項目を指定します。以下に注意してください:
|
|
|
[Main Class] |
メインクラス名を設定します。このパラメーターは、開発言語が |
|
|
[パラメータ] |
必要に応じてパラメーターを追加し、複数のパラメーターをスペースで区切ることができます。 DataWorks はスケジューリングパラメーターをサポートしています。 パラメータ の形式は スケジューリングパラメーター値のサポートされている形式については、「スケジューリングパラメーターのソースと式」をご参照ください。 |
|
|
[jar リソースの選択] |
プログラミング言語が リソースファイルは DataWorks にアップロードしてコミットする必要があります。詳細については、「リソース管理」をご参照ください。 |
リソースコマンド:
|
|
[ファイルリソースの選択] |
ジョブのファイルリソースを指定します。 |
|
|
[archives リソースの選択] |
ジョブのアーカイブリソースを指定します。.zip アーカイブのみがサポートされています。 |
|
Python
|
パラメーター |
説明 |
spark-submit コマンド |
|
[Spark バージョン] |
Spark のバージョン。有効な値は、Spark 1.x、Spark 2.x、Spark 3.x です。 説明
Spark 3.x バージョンを使用する MaxCompute Spark ノードの送信時にエラーが発生した場合は、サーバーレスリソースグループを購入して使用する必要があります。詳細については、「サーバーレスリソースグループの使用」をご参照ください。 |
— |
|
[言語] |
プログラミング言語。Spark on MaxCompute ジョブの開発に使用した言語に基づいて [Python] を選択してください。 |
— |
|
[メインの Python リソースの選択] |
ジョブのメイン Python リソースファイルを指定してください。 リソースファイルは DataWorks にアップロードしてコミットする必要があります。詳細については、「リソース管理」をご参照ください。 |
|
|
[設定項目] |
ジョブを送信するための設定項目を指定します。以下に注意してください:
|
|
|
[パラメータ] |
必要に応じて、パラメーターをスペースで区切って追加できます。 DataWorks はスケジューリングパラメーターをサポートしており、パラメータ フィールドに スケジューリングパラメーター値のサポートされている形式については、「スケジューリングパラメーターのソースと式」をご参照ください。 |
|
|
[Python リソースの選択] |
開発言語が リソースファイルは DataWorks にアップロードしてコミットする必要があります。詳細については、「リソース管理」をご参照ください。 |
|
|
[ファイルリソースの選択] |
ジョブのファイルリソースを指定します。 |
|
|
[archives リソースの選択] |
ジョブのアーカイブリソースを指定します。 |
|
操作手順
-
リソースを作成します。
-
Data Studio ページの左側のナビゲーションバーでリソース管理を見つけ、新規作成をクリックします。 MaxCompute Spark タイプの Python リソースを作成し、
spark_is_number.pyという名前を付けます。 詳細については、「リソース管理」をご参照ください。 コードは以下のとおりです。# -*- coding: utf-8 -*- import sys from pyspark.sql import SparkSession try: # Python 2 の場合 reload(sys) sys.setdefaultencoding('utf8') except: # Python 3 では不要 pass if __name__ == '__main__': spark = SparkSession.builder\ .appName("spark sql")\ .config("spark.sql.broadcastTimeout", 20 * 60)\ .config("spark.sql.crossJoin.enabled", True)\ .config("odps.exec.dynamic.partition.mode", "nonstrict")\ .config("spark.sql.catalogImplementation", "odps")\ .getOrCreate() def is_number(s): try: float(s) return True except ValueError: pass try: import unicodedata unicodedata.numeric(s) return True except (TypeError, ValueError): pass return False print(is_number('foo')) print(is_number('1')) print(is_number('1.3')) print(is_number('-1.37')) print(is_number('1e3')) -
リソースを保存します。
-
-
作成した MaxCompute Spark ノードで、ノードパラメーターとスケジューリングパラメーターを設定します。詳細については、「パラメーター」をご参照ください。
-
スケジュールに基づいてジョブを実行するには、ビジネス要件に基づいてスケジューリングプロパティを設定します。詳細については、「ノードのスケジューリング設定」をご参照ください。
-
ノードジョブを設定したら、ノードをデプロイします。詳細については、「ノード/ワークフローのデプロイ」をご参照ください。
-
ジョブがデプロイされた後、オペレーションセンターに移動して定期ジョブの実行ステータスを表示できます。詳細については、「オペレーションセンター入門」をご参照ください。
説明-
MaxCompute Spark ノードには、Data Studio での実行エントリポイントはありません。開発環境のオペレーションセンターで Spark ジョブを実行する必要があります。
-
データバックフィルインスタンスが正常に実行された後、インスタンスの実行ログにあるトラッキング URL を開いて結果を表示できます。
-
関連ドキュメント
-
他のユースケースでの Spark on MaxCompute ジョブの開発に関する詳細については、次のトピックをご参照ください:
-
Spark FAQ:一般的な Spark 実行の問題と解決策を掲載しています。トラブルシューティングの迅速化にお役立てください。詳細については、「Spark FAQ」をご参照ください。
-
Spark ジョブの診断:MaxCompute は、Logview ツールと Spark Web UI を提供します。ジョブログを使用して、ジョブが正しく送信および実行されているかを確認できます。詳細については、「Spark ジョブの診断」をご参照ください。