Container Service for Kubernetes (ACK) クラスターで、組み込みの PageRank ジョブと Hadoop OSS SDK、Hadoop S3 SDK、または JindoSDK を使用して、Object Storage Service (OSS) データを読み書きする Spark ジョブを実行します。
前提条件
次の要件を満たしていることを確認してください。
-
Kubernetes 1.24 以降を実行する ACK Pro マネージドクラスターまたは ACK Serverless Pro クラスターが作成されていること。詳細については、「ACK マネージドクラスターの作成」、「ACK Serverless クラスターの作成」、「ACK クラスターの手動アップグレード」をご参照ください。
-
ack-spark-operator アドオンがインストールされていること。詳細については、「ack-spark-operator コンポーネントのインストール」をご参照ください。
-
kubectl クライアントがクラスターに接続されていること。詳細については、「kubeconfig ファイルを取得し、kubectl を使用してクラスターに接続する」をご参照ください。
-
OSS バケットが作成されていること。詳細については、「バケットの作成」をご参照ください。
-
ossutil がインストールおよび設定されていること。詳細については、「ossutil コマンドリファレンス」をご参照ください。
OSS 統合の選択
OSS へのアクセスには 3 つの SDK が利用できます。先に進む前に、いずれか 1 つを選択してください。
| SDK | 使用ケース | URI スキーム |
|---|---|---|
| Hadoop OSS SDK | Alibaba Cloud ファイルシステムによるネイティブの OSS サポートです。ほとんどのワークロードでシンプルです。 | oss:// |
| Hadoop S3 SDK | クラスターやツールが S3 互換 API を使用する場合です。Alibaba Cloud 固有の依存関係はありません。 | s3a:// |
| JindoSDK | Jindo アクセラレーションによる最適化された OSS パフォーマンスです。 | oss:// |
すべてのステップで同じ SDK を使用してください。
ステップ 1: テストデータの準備とアップロード
PageRank データセットを生成し、OSS バケットにアップロードします。
-
次の内容で
generate_pagerank_dataset.shを作成します。#!/bin/bash # 引数の数を確認 if [ "$#" -ne 2 ]; then echo "Usage: $0 M N" echo "M: Number of web pages" echo "N: Number of records to generate" exit 1 fi M=$1 N=$2 # M と N が正の整数であるか検証 if ! [[ "$M" =~ ^[0-9]+$ ]] || ! [[ "$N" =~ ^[0-9]+$ ]]; then echo "Both M and N must be positive integers." exit 1 fi # データセットの生成 for ((i=1; i<=$N; i++)); do # 送信元と送信先のページが異なることを確認 while true; do src=$((RANDOM % M + 1)) dst=$((RANDOM % M + 1)) if [ "$src" -ne "$dst" ]; then echo "$src $dst" break fi done done -
データセットを生成します。
M=100000 # Web ページの数 N=10000000 # 生成するレコードの数 # データセットをランダムに生成し、pagerank_dataset.txt に保存 bash generate_pagerank_dataset.sh $M $N > pagerank_dataset.txt -
データセットを OSS バケットの
data/にアップロードします。ossutil cp pagerank_dataset.txt oss://<BUCKET_NAME>/data/
ステップ 2: Spark コンテナイメージのビルド
OSS アクセス用の JAR 依存関係を含むコンテナイメージをビルドします。詳細については、「Container Registry Enterprise Edition インスタンスを使用したイメージのビルド」をご参照ください。
サンプル Dockerfile のベース Spark イメージは、オープンソースコミュニティから提供されています。必要に応じて置き換え、SDK バージョンを Spark バージョンに合わせてください。
Hadoop OSS SDK の使用
この例では、Spark 3.5.5 と Hadoop OSS SDK 3.3.4 を使用します。
ARG SPARK_IMAGE=spark:3.5.5
FROM ${SPARK_IMAGE}
# Hadoop Alibaba Cloud OSS サポート用の依存関係を追加
ADD --chown=spark:spark --chmod=644 https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aliyun/3.3.4/hadoop-aliyun-3.3.4.jar ${SPARK_HOME}/jars
ADD --chown=spark:spark --chmod=644 https://repo1.maven.org/maven2/com/aliyun/oss/aliyun-sdk-oss/3.17.4/aliyun-sdk-oss-3.17.4.jar ${SPARK_HOME}/jars
ADD --chown=spark:spark --chmod=644 https://repo1.maven.org/maven2/org/jdom/jdom2/2.0.6.1/jdom2-2.0.6.1.jar ${SPARK_HOME}/jars
Hadoop S3 SDK の使用
この例では、Spark 3.5.5 と Hadoop S3 SDK 3.3.4 を使用します。
ARG SPARK_IMAGE=spark:3.5.5
FROM ${SPARK_IMAGE}
# Hadoop AWS S3 サポート用の依存関係を追加
ADD --chown=spark:spark --chmod=644 https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/3.3.4/hadoop-aws-3.3.4.jar ${SPARK_HOME}/jars
ADD --chown=spark:spark --chmod=644 https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/1.12.367/aws-java-sdk-bundle-1.12.367.jar ${SPARK_HOME}/jars
JindoSDK の使用
この例では、Spark 3.5.5 と JindoSDK 6.8.0 を使用します。
ARG SPARK_IMAGE=spark:3.5.5
FROM ${SPARK_IMAGE}
# JindoSDK サポート用の依存関係を追加
ADD --chown=spark:spark --chmod=644 https://jindodata-binary.oss-cn-shanghai.aliyuncs.com/mvn-repo/com/aliyun/jindodata/jindo-core/6.8.0/jindo-core-6.8.0.jar ${SPARK_HOME}/jars
ADD --chown=spark:spark --chmod=644 https://jindodata-binary.oss-cn-shanghai.aliyuncs.com/mvn-repo/com/aliyun/jindodata/jindo-sdk/6.8.0/jindo-sdk-6.8.0.jar ${SPARK_HOME}/jars
ステップ 3: Kubernetes Secret への OSS 認証情報の保存
AccessKey ID と AccessKey Secret は、SparkApplication マニフェストに埋め込むのではなく、Kubernetes Secret に保存します。Spark オペレーターは、Secret のキーを環境変数としてドライバー Pod とエグゼキューター Pod に挿入します。
ステップ 4 のhadoopConfフィールドはEnvironmentVariableCredentialsProviderを使用しており、これらの環境変数を自動的に読み取ります。認証情報を YAML ファイルに埋め込んだり、ソース管理にコミットしたりしないでください。
Hadoop OSS SDK の使用
-
spark-oss-secret.yamlを作成します。apiVersion: v1 kind: Secret metadata: name: spark-oss-secret namespace: default stringData: # を Alibaba Cloud アカウントの AccessKey ID に置き換えます。 OSS_ACCESS_KEY_ID: <ACCESS_KEY_ID> # を Alibaba Cloud アカウントの AccessKey Secret に置き換えます。 OSS_ACCESS_KEY_SECRET: <ACCESS_KEY_SECRET> -
Secret を適用します。
kubectl apply -f spark-oss-secret.yaml出力例:
secret/spark-oss-secret created
Hadoop S3 SDK の使用
Hadoop S3 SDK は AWS_ACCESS_KEY_ID と AWS_SECRET_ACCESS_KEY から認証情報を読み取るため、キー名は OSS SDK とは異なります。
-
spark-s3-secret.yamlを作成します。apiVersion: v1 kind: Secret metadata: name: spark-s3-secret namespace: default stringData: # を Alibaba Cloud アカウントの AccessKey ID に置き換えます。 AWS_ACCESS_KEY_ID: <ACCESS_KEY_ID> # を Alibaba Cloud アカウントの AccessKey Secret に置き換えます。 AWS_SECRET_ACCESS_KEY: <ACCESS_KEY_SECRET> -
Secret を適用します。
kubectl apply -f spark-s3-secret.yaml出力例:
secret/spark-s3-secret created
JindoSDK の使用
-
spark-oss-secret.yamlを作成します。apiVersion: v1 kind: Secret metadata: name: spark-oss-secret namespace: default stringData: # を Alibaba Cloud アカウントの AccessKey ID に置き換えます。 OSS_ACCESS_KEY_ID: <ACCESS_KEY_ID> # を Alibaba Cloud アカウントの AccessKey Secret に置き換えます。 OSS_ACCESS_KEY_SECRET: <ACCESS_KEY_SECRET> -
Secret を適用します。
kubectl apply -f spark-oss-secret.yaml出力例:
secret/spark-oss-secret created
ステップ 4: Spark ジョブの送信
SparkApplication マニフェストを作成して送信し、OSS データセットを対象に PageRank ジョブを実行します。
Hadoop OSS SDK の使用
spark-pagerank.yaml を作成します。Hadoop-Aliyun モジュールには、利用可能なすべての OSS パラメーターが記載されています。
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: spark-pagerank
namespace: default
spec:
type: Scala
mode: cluster
# をステップ 2 でビルドした Spark コンテナイメージに置き換えます。
image: <SPARK_IMAGE>
mainApplicationFile: local:///opt/spark/examples/jars/spark-examples_2.12-3.5.5.jar
mainClass: org.apache.spark.examples.SparkPageRank
arguments:
- oss://<OSS_BUCKET>/data/pagerank_dataset.txt # 入力テストデータセット。 を実際の OSS バケット名に置き換えます。
- "10" # 反復回数
sparkVersion: 3.5.5
hadoopConf:
fs.AbstractFileSystem.oss.impl: org.apache.hadoop.fs.aliyun.oss.OSS
fs.oss.impl: org.apache.hadoop.fs.aliyun.oss.AliyunOSSFileSystem
# OSS エンドポイント。 を実際の OSS エンドポイントに置き換えます。
# たとえば、中国 (北京) リージョンの内部エンドポイントは oss-cn-beijing-internal.aliyuncs.com です。
fs.oss.endpoint: <OSS_ENDPOINT>
fs.oss.credentials.provider: com.aliyun.oss.common.auth.EnvironmentVariableCredentialsProvider
driver:
cores: 1
coreLimit: 1200m
memory: 512m
envFrom:
- secretRef:
name: spark-oss-secret # OSS へのアクセスに使用する Secret
serviceAccount: spark-operator-spark
executor:
instances: 2
cores: 1
coreLimit: "2"
memory: 8g
envFrom:
- secretRef:
name: spark-oss-secret # OSS へのアクセスに使用する Secret
restartPolicy:
type: Never
Hadoop S3 SDK の使用
spark-pagerank.yaml を作成します。Hadoop-AWS モジュールには、利用可能なすべての S3 パラメーターが記載されています。
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: spark-pagerank
namespace: default
spec:
type: Scala
mode: cluster
# をステップ 2 でビルドした Spark コンテナイメージに置き換えます。
image: <SPARK_IMAGE>
mainApplicationFile: local:///opt/spark/examples/jars/spark-examples_2.12-3.5.5.jar
mainClass: org.apache.spark.examples.SparkPageRank
arguments:
- s3a://<OSS_BUCKET>/data/pagerank_dataset.txt # 入力テストデータセット。 を実際の OSS バケット名に置き換えます。
- "10" # 反復回数
sparkVersion: 3.5.5
hadoopConf:
fs.s3a.impl: org.apache.hadoop.fs.s3a.S3AFileSystem
# OSS エンドポイント。 を実際の OSS エンドポイントに置き換えます。
# たとえば、中国 (北京) リージョンの内部エンドポイントは oss-cn-beijing-internal.aliyuncs.com です。
fs.s3a.endpoint: <OSS_ENDPOINT>
# OSS エンドポイントのリージョン。たとえば、中国 (北京) リージョンの場合は cn-beijing です。
fs.s3a.endpoint.region: <OSS_REGION>
driver:
cores: 1
coreLimit: 1200m
memory: 512m
envFrom:
- secretRef:
name: spark-s3-secret # OSS へのアクセスに使用する Secret
serviceAccount: spark-operator-spark
executor:
instances: 2
cores: 1
coreLimit: "2"
memory: 8g
envFrom:
- secretRef:
name: spark-s3-secret # OSS へのアクセスに使用する Secret
restartPolicy:
type: Never
JindoSDK の使用
spark-pagerank.yaml を作成します。
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: spark-pagerank
namespace: default
spec:
type: Scala
mode: cluster
# をステップ 2 でビルドした Spark コンテナイメージに置き換えます。
image: <SPARK_IMAGE>
mainApplicationFile: local:///opt/spark/examples/jars/spark-examples_2.12-3.5.5.jar
mainClass: org.apache.spark.examples.SparkPageRank
arguments:
- oss://<OSS_BUCKET>/data/pagerank_dataset.txt # 入力テストデータセット。 を実際の OSS バケット名に置き換えます。
- "10" # 反復回数
sparkVersion: 3.5.5
hadoopConf:
fs.AbstractFileSystem.oss.impl: com.aliyun.jindodata.oss.JindoOSS
fs.oss.impl: com.aliyun.jindodata.oss.JindoOssFileSystem
# OSS エンドポイント。 を実際の OSS エンドポイントに置き換えます。
# たとえば、中国 (北京) リージョンの内部エンドポイントは oss-cn-beijing-internal.aliyuncs.com です。
fs.oss.endpoint: <OSS_ENDPOINT>
fs.oss.credentials.provider: com.aliyun.jindodata.oss.auth.EnvironmentVariableCredentialsProvider
driver:
cores: 1
coreLimit: 1200m
memory: 512m
serviceAccount: spark-operator-spark
envFrom:
- secretRef:
name: spark-oss-secret # OSS へのアクセスに使用する Secret
executor:
instances: 2
cores: 1
coreLimit: "2"
memory: 8g
envFrom:
- secretRef:
name: spark-oss-secret # OSS へのアクセスに使用する Secret
restartPolicy:
type: Never
ジョブの送信と確認
これらのコマンドは、3 つの SDK オプションすべてに共通です。
-
ジョブを送信します。
kubectl apply -f spark-pagerank.yaml -
ジョブのステータスを監視します。
kubectl get sparkapplications spark-pagerankジョブ完了時の出力例:
NAME STATUS ATTEMPTS START FINISH AGE spark-pagerank COMPLETED 1 2024-10-09T12:54:25Z 2024-10-09T12:55:46Z 90s -
ドライバー Pod のログの最後の 20 行を表示します。
kubectl logs spark-pagerank-driver --tail=20Hadoop OSS SDK — 出力例:
30024 has rank: 1.0709659078941967 . 21390 has rank: 0.9933356174074005 . 28500 has rank: 1.0404018494028928 . 2137 has rank: 0.9931000490520374 . 3406 has rank: 0.9562543137167121 . 20904 has rank: 0.8827028621652337 . 25604 has rank: 1.0270134041934191 . 24/10/09 12:48:36 INFO SparkUI: Stopped Spark web UI at http://spark-pagerank-dd0d4d927151c9d0-driver-svc.default.svc:4040 24/10/09 12:48:36 INFO KubernetesClusterSchedulerBackend: Shutting down all executors 24/10/09 12:48:36 INFO KubernetesClusterSchedulerBackend$KubernetesDriverEndpoint: Asking each executor to shut down 24/10/09 12:48:36 WARN ExecutorPodsWatchSnapshotSource: Kubernetes client has been closed. 24/10/09 12:48:36 INFO MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped! 24/10/09 12:48:36 INFO MemoryStore: MemoryStore cleared 24/10/09 12:48:36 INFO BlockManager: BlockManager stopped 24/10/09 12:48:36 INFO BlockManagerMaster: BlockManagerMaster stopped 24/10/09 12:48:36 INFO OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: OutputCommitCoordinator stopped! 24/10/09 12:48:36 INFO SparkContext: Successfully stopped SparkContext 24/10/09 12:48:36 INFO ShutdownHookManager: Shutdown hook called 24/10/09 12:48:36 INFO ShutdownHookManager: Deleting directory /tmp/spark-e8b8c2ab-c916-4f84-b60f-f54c0de3a7f0 24/10/09 12:48:36 INFO ShutdownHookManager: Deleting directory /var/data/spark-c5917d98-06fb-46fe-85bc-199b839cb885/spark-23e2c2ae-4754-43ae-854d-2752eb83b2c5Hadoop S3 SDK — 出力例:
3406 has rank: 0.9562543137167121 . 20904 has rank: 0.8827028621652337 . 25604 has rank: 1.0270134041934191 . 25/04/07 03:54:11 INFO SparkContext: SparkContext is stopping with exitCode 0. 25/04/07 03:54:11 INFO SparkUI: Stopped Spark web UI at http://spark-pagerank-0f7dec960e615617-driver-svc.spark.svc:4040 25/04/07 03:54:11 INFO KubernetesClusterSchedulerBackend: Shutting down all executors 25/04/07 03:54:11 INFO KubernetesClusterSchedulerBackend$KubernetesDriverEndpoint: Asking each executor to shut down 25/04/07 03:54:11 WARN ExecutorPodsWatchSnapshotSource: Kubernetes client has been closed. 25/04/07 03:54:11 INFO MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped! 25/04/07 03:54:11 INFO MemoryStore: MemoryStore cleared 25/04/07 03:54:11 INFO BlockManager: BlockManager stopped 25/04/07 03:54:11 INFO BlockManagerMaster: BlockManagerMaster stopped 25/04/07 03:54:11 INFO OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: OutputCommitCoordinator stopped! 25/04/07 03:54:11 INFO SparkContext: Successfully stopped SparkContext 25/04/07 03:54:11 INFO ShutdownHookManager: Shutdown hook called 25/04/07 03:54:11 INFO ShutdownHookManager: Deleting directory /var/data/spark-20d425bb-f442-4b0a-83e2-5a0202959a54/spark-ff5bbf08-4343-4a7a-9ce0-3f7c127cf4a9 25/04/07 03:54:11 INFO ShutdownHookManager: Deleting directory /tmp/spark-a421839a-07af-49c0-b637-f15f76c3e752 25/04/07 03:54:11 INFO MetricsSystemImpl: Stopping s3a-file-system metrics system... 25/04/07 03:54:11 INFO MetricsSystemImpl: s3a-file-system metrics system stopped. 25/04/07 03:54:11 INFO MetricsSystemImpl: s3a-file-system metrics system shutdown complete.JindoSDK — 出力例:
21390 has rank: 0.9933356174074005 . 28500 has rank: 1.0404018494028928 . 2137 has rank: 0.9931000490520374 . 3406 has rank: 0.9562543137167121 . 20904 has rank: 0.8827028621652337 . 25604 has rank: 1.0270134041934191 . 24/10/09 12:55:44 INFO SparkContext: SparkContext is stopping with exitCode 0. 24/10/09 12:55:44 INFO SparkUI: Stopped Spark web UI at http://spark-pagerank-6a5e3d9271584856-driver-svc.default.svc:4040 24/10/09 12:55:44 INFO KubernetesClusterSchedulerBackend: Shutting down all executors 24/10/09 12:55:44 INFO KubernetesClusterSchedulerBackend$KubernetesDriverEndpoint: Asking each executor to shut down 24/10/09 12:55:44 WARN ExecutorPodsWatchSnapshotSource: Kubernetes client has been closed. 24/10/09 12:55:45 INFO MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped! 24/10/09 12:55:45 INFO MemoryStore: MemoryStore cleared 24/10/09 12:55:45 INFO BlockManager: BlockManager stopped 24/10/09 12:55:45 INFO BlockManagerMaster: BlockManagerMaster stopped 24/10/09 12:55:45 INFO OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: OutputCommitCoordinator stopped! 24/10/09 12:55:45 INFO SparkContext: Successfully stopped SparkContext 24/10/09 12:55:45 INFO ShutdownHookManager: Shutdown hook called 24/10/09 12:55:45 INFO ShutdownHookManager: Deleting directory /var/data/spark-87e8406e-06a7-4b4a-b18f-2193da299d35/spark-093a1b71-121a-4367-9d22-ad4e397c9815 24/10/09 12:55:45 INFO ShutdownHookManager: Deleting directory /tmp/spark-723e2039-a493-49e8-b86d-fff5fd1bb168
(任意) ステップ 5: クリーンアップ
不要になった Spark ジョブと Secret を削除します。
Spark ジョブを削除します。
kubectl delete -f spark-pagerank.yaml
Secret を削除します。
Hadoop OSS SDK または JindoSDK:
kubectl delete -f spark-oss-secret.yaml
Hadoop S3 SDK:
kubectl delete -f spark-s3-secret.yaml