ACK に Spark Operator をインストールして、宣言的な Kubernetes マニフェストで Spark ジョブのライフサイクル全体を管理します。
前提条件
-
Kubernetes 1.24 以降を実行している ACK Pro クラスターまたはACK Serverless Pro クラスター。「ACK マネージドクラスターの作成」、「ACK Serverless クラスターの作成」、および「ACK クラスターの手動アップグレード」をご参照ください。
kubectlクライアントがACKクラスターに接続されています。 詳細については、「クラスターのkubeconfigファイルを取得し、kubectlを使用してクラスターに接続する」をご参照ください。
仕組み
Spark Operator は、SparkApplication や ScheduledSparkApplication などの CustomResourceDefinition (CRD) リソースを使用して、Kubernetes 上で Spark ジョブのライフサイクルを自動化します。自動スケーリング、ヘルスチェック、リソース管理などのネイティブ Kubernetes 機能を活用します。ACK は、kubeflow/spark-operator をベースとした ack-spark-operator を提供しています。「Spark Operator | Kubeflow」をご参照ください。
利点:
-
管理の簡素化:宣言的な Kubernetes 設定により、Spark ジョブのデプロイとライフサイクルを自動化します。
-
マルチテナンシーのサポート:Kubernetes の名前空間とリソースクォータを使用してリソースを分離します。ノードセレクションを使用して、専用リソース上で Spark ワークロードを実行します。
-
弾力的なリソースプロビジョニング:ピーク負荷時に Elastic Container Instance (ECI) や弾性ノードプールなどの弾性リソースを使用してスケーリングし、パフォーマンスとコストのバランスを取ります。
ユースケース:
-
データ分析:Spark をインタラクティブな分析やデータクレンジングに使用します。
-
バッチコンピューティング:スケジュールされたバッチジョブを実行して、大規模なデータセットを処理します。
-
リアルタイム処理:Spark Streaming により、リアルタイムのデータストリーム処理を実現します。
手順の概要
このワークフローでは、Spark Operator のデプロイ、ジョブの投入、実行の監視、およびジョブのライフサイクル管理について説明します。
-
ack-spark-operator コンポーネントのデプロイ:ACK クラスターに Spark Operator をインストールします。
-
Spark ジョブの投入:Spark ジョブのマニフェストを作成して投入します。
-
Spark ジョブの監視:ジョブのステータス、Pod のステータस、およびログを確認します。
-
Spark Web UI へのアクセス:ブラウザでジョブ実行の詳細を表示します。
-
Spark ジョブの更新:ジョブマニフェストを変更して再適用します。
-
Spark ジョブの削除:ジョブを削除してリソースを解放します。
ステップ 1: ack-spark-operator コンポーネントのデプロイ
ACKコンソールにログインします。 左側のナビゲーションウィンドウで、 を選択します。
-
Marketplaceページで、アプリカタログタブをクリックし、[ack-spark-operator]を検索して選択します。
-
ack-spark-operator ページで、デプロイ をクリックします。
-
作成する パネルで、クラスターと名前空間を選択し、次へ をクリックします。
-
パラメーター ページで、パラメーターを設定し、OK をクリックします。
主要なパラメーターは以下のとおりです。完全なリストについては、[ack-spark-operator] ページの ConfigMap タブをご参照ください。
パラメーター
説明
デフォルト
controller.replicasコントローラーレプリカの数。
1
webhook.replicasWebhook レプリカの数。
1
spark.jobNamespacesSpark ジョブを実行できる名前空間。空の文字列を指定すると、すべての名前空間が許可されます。複数の値を指定する場合は、カンマ (
,) で区切ります。-
["default"](デフォルト) -
[""](すべての名前空間) -
["ns1","ns2","ns3"](複数の名前空間)
spark.serviceAccount.nameSpark Operator は、
spark.jobNamespacesで指定された各名前空間に、spark-operator-sparkという名前の ServiceAccount と必要な RBAC リソースを作成します。カスタマイズした場合は、Spark ジョブを投入する際に新しい名前を指定してください。spark-operator-spark -
ステップ 2: Spark ジョブの投入
SparkApplication マニフェストを作成して、Spark ジョブを投入します。
-
次の内容でマニフェストを作成し、
spark-pi.yamlとして保存します。apiVersion: sparkoperator.k8s.io/v1beta2 kind: SparkApplication metadata: name: spark-pi namespace: default # 名前空間は spark.jobNamespaces で指定された名前空間のリストに含まれている必要があります。 spec: type: Scala mode: cluster image: registry-cn-hangzhou.ack.aliyuncs.com/ack-demo/spark:3.5.4 imagePullPolicy: IfNotPresent mainClass: org.apache.spark.examples.SparkPi mainApplicationFile: local:///opt/spark/examples/jars/spark-examples_2.12-3.5.4.jar arguments: - "1000" sparkVersion: 3.5.4 driver: cores: 1 coreLimit: 1200m memory: 512m serviceAccount: spark-operator-spark # ServiceAccount 名をカスタマイズした場合は、それに応じて値を変更してください。 executor: instances: 1 cores: 1 coreLimit: 1200m memory: 512m restartPolicy: type: Never -
Spark ジョブを投入します:
kubectl apply -f spark-pi.yaml期待される出力:
sparkapplication.sparkoperator.k8s.io/spark-pi created
ステップ 3: Spark ジョブの監視
Spark ジョブのステータス、Pod、およびログを確認します。
-
Spark ジョブのステータスを確認します:
kubectl get sparkapplication spark-pi期待される出力:
NAME STATUS ATTEMPTS START FINISH AGE spark-pi SUBMITTED 1 2024-06-04T03:17:11Z <no value> 15s -
Pod のステータスを確認します。これにより、ラベル
sparkoperator.k8s.io/app-name=spark-piで Pod をフィルタリングします。kubectl get pod -l sparkoperator.k8s.io/app-name=spark-pi期待される出力:
NAME READY STATUS RESTARTS AGE spark-pi-driver 1/1 Running 0 49s spark-pi-7272428fc8f5f392-exec-1 1/1 Running 0 13sジョブが完了すると、driver はすべての executor Pod を削除します。
-
Spark ジョブの詳細を表示します:
kubectl describe sparkapplication spark-pi -
driver Pod のログの最後の 20 行を表示します。
kubectl logs --tail=20 spark-pi-driver期待される出力:
24/05/30 10:05:30 INFO TaskSchedulerImpl: Removed TaskSet 0.0, whose tasks have all completed, from pool 24/05/30 10:05:30 INFO DAGScheduler: ResultStage 0 (reduce at SparkPi.scala:38) finished in 7.942 s 24/05/30 10:05:30 INFO DAGScheduler: Job 0 is finished. Cancelling potential speculative or zombie tasks for this job 24/05/30 10:05:30 INFO TaskSchedulerImpl: Killing all running tasks in stage 0: Stage finished 24/05/30 10:05:30 INFO DAGScheduler: Job 0 finished: reduce at SparkPi.scala:38, took 8.043996 s Pi is roughly 3.1419522314195225 24/05/30 10:05:30 INFO SparkContext: SparkContext is stopping with exitCode 0. 24/05/30 10:05:30 INFO SparkUI: Stopped Spark web UI at http://spark-pi-1e18858fc8f56b14-driver-svc.default.svc:4040 24/05/30 10:05:30 INFO KubernetesClusterSchedulerBackend: Shutting down all executors 24/05/30 10:05:30 INFO KubernetesClusterSchedulerBackend$KubernetesDriverEndpoint: Asking each executor to shut down 24/05/30 10:05:30 WARN ExecutorPodsWatchSnapshotSource: Kubernetes client has been closed. 24/05/30 10:05:30 INFO MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped! 24/05/30 10:05:30 INFO MemoryStore: MemoryStore cleared 24/05/30 10:05:30 INFO BlockManager: BlockManager stopped 24/05/30 10:05:30 INFO BlockManagerMaster: BlockManagerMaster stopped 24/05/30 10:05:30 INFO OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: OutputCommitCoordinator stopped! 24/05/30 10:05:30 INFO SparkContext: Successfully stopped SparkContext 24/05/30 10:05:30 INFO ShutdownHookManager: Shutdown hook called 24/05/30 10:05:30 INFO ShutdownHookManager: Deleting directory /var/data/spark-14ed60f1-82cd-4a33-b1b3-9e5d975c5b1e/spark-01120c89-5296-4c83-8a20-0799eef4e0ee 24/05/30 10:05:30 INFO ShutdownHookManager: Deleting directory /tmp/spark-5f98ed73-576a-41be-855d-dabdcf7de189
ステップ 4: Spark Web UI へのアクセス
Web UI は、driver Pod が Running 状態の場合にのみ利用可能です。
デフォルトでは、controller.uiService.enable は true に設定されており、ポートフォワーディング用に Web UI を公開する Service が作成されます。false に設定した場合、Service は作成されず、driver Pod から直接ポートフォワーディングを行う必要があります。
kubectl port-forward はテストには適していますが、セキュリティリスクがあるため、本番環境での使用は推奨されません。
-
Web UI ポートをローカルマシンに転送します。
-
Service 経由のポートフォワーディング
kubectl port-forward services/spark-pi-ui-svc 4040 -
Pod 経由のポートフォワーディング
kubectl port-forward pods/spark-pi-driver 4040期待される出力:
Forwarding from 127.0.0.1:4040 -> 4040 Forwarding from [::1]:4040 -> 4040
-
-
ブラウザで http://127.0.0.1:4040 を開きます。
(オプション) ステップ 5: Spark ジョブの更新
ジョブマニフェストを更新して、Spark ジョブのパラメーターを変更します。
-
spark-pi.yamlを編集します。たとえば、argumentsを10000に、executorインスタンスを2に設定します。apiVersion: sparkoperator.k8s.io/v1beta2 kind: SparkApplication metadata: name: spark-pi spec: type: Scala mode: cluster image: registry-cn-hangzhou.ack.aliyuncs.com/ack-demo/spark:3.5.4 imagePullPolicy: IfNotPresent mainClass: org.apache.spark.examples.SparkPi mainApplicationFile: local:///opt/spark/examples/jars/spark-examples_2.12-3.5.4.jar arguments: - "10000" sparkVersion: 3.5.4 driver: cores: 1 coreLimit: 1200m memory: 512m serviceAccount: spark-operator-spark # ServiceAccount 名をカスタマイズした場合は、それに応じて値を変更してください。 executor: instances: 2 cores: 1 coreLimit: 1200m memory: 512m restartPolicy: type: Never -
変更を適用します。
kubectl apply -f spark-pi.yaml -
ジョブのステータスを確認します。
kubectl get sparkapplication spark-piSpark ジョブが再度実行されます。期待される出力:
NAME STATUS ATTEMPTS START FINISH AGE spark-pi RUNNING 1 2024-06-04T03:37:34Z <no value> 20m
(オプション) ステップ 6: Spark ジョブの削除
Spark ジョブを削除して、関連するリソースを解放します。
kubectl delete -f spark-pi.yaml
または:
kubectl delete sparkapplication spark-pi
関連ドキュメント
-
「Spark History Server を使用した Spark ジョブ情報の表示」をご参照ください。
-
「Log Service を使用した Spark ジョブログの収集」をご参照ください。
-
「Spark ジョブでの OSS データの読み取りと書き込み」をご参照ください。
-
「ECI 弾性リソースを使用した Spark ジョブの実行」をご参照ください。
-
「Spark ジョブでの RSS としての Celeborn の使用」をご参照ください。