Jalankan pekerjaan Spark yang membaca dan menulis data OSS di kluster ACK menggunakan pekerjaan PageRank bawaan dengan Hadoop OSS SDK, Hadoop S3 SDK, atau JindoSDK.
Prasyarat
Pastikan Anda telah:
-
Membuat kluster ACK Pro atau kluster ACK Serverless Pro yang menjalankan Kubernetes 1.24 atau lebih baru: Buat kluster ACK yang dikelola, Buat kluster ACK Serverless, dan Lakukan peningkatan manual kluster ACK
-
Menginstal add-on ack-spark-operator: Instal komponen ack-spark-operator
-
Menghubungkan klien kubectl ke kluster: Dapatkan file kubeconfig dan hubungkan kubectl ke kluster
-
Membuat bucket OSS: Buat bucket
-
Menginstal dan mengonfigurasi ossutil: Referensi perintah ossutil
Pilih integrasi OSS
Tiga SDK tersedia untuk akses OSS. Pilih salah satu sebelum melanjutkan.
| SDK | Gunakan saat | Skema URI |
|---|---|---|
| Hadoop OSS SDK | Dukungan OSS native melalui filesystem Aliyun; cocok untuk sebagian besar beban kerja | oss:// |
| Hadoop S3 SDK | Kluster atau tooling Anda menggunakan API yang kompatibel dengan S3; tanpa dependensi khusus Aliyun | s3a:// |
| JindoSDK | Performa OSS yang dioptimalkan dengan akselerasi Jindo | oss:// |
Gunakan SDK yang sama untuk semua langkah.
Langkah 1: Siapkan dan unggah data uji
Buat set data PageRank dan unggah ke bucket OSS Anda.
-
Buat
generate_pagerank_dataset.shdengan konten berikut:#!/bin/bash # Periksa jumlah argumen if [ "$#" -ne 2 ]; then echo "Usage: $0 M N" echo "M: Jumlah halaman web" echo "N: Jumlah record yang akan dibuat" exit 1 fi M=$1 N=$2 # Verifikasi apakah M dan N adalah bilangan bulat positif if ! [[ "$M" =~ ^[0-9]+$ ]] || ! [[ "$N" =~ ^[0-9]+$ ]]; then echo "M dan N harus berupa bilangan bulat positif." exit 1 fi # Buat set data for ((i=1; i<=$N; i++)); do # Pastikan halaman sumber dan tujuan berbeda while true; do src=$((RANDOM % M + 1)) dst=$((RANDOM % M + 1)) if [ "$src" -ne "$dst" ]; then echo "$src $dst" break fi done done -
Buat set data:
M=100000 # Jumlah halaman web N=10000000 # Jumlah record # Buat set data secara acak dan simpan sebagai pagerank_dataset.txt bash generate_pagerank_dataset.sh $M $N > pagerank_dataset.txt -
Unggah set data ke
data/di bucket OSS Anda:ossutil cp pagerank_dataset.txt oss://<BUCKET_NAME>/data/
Langkah 2: Bangun gambar kontainer Spark
Bangun gambar kontainer dengan dependensi JAR untuk akses OSS. Lihat Gunakan instans Container Registry Enterprise Edition untuk membangun gambar.
Gambar dasar Spark dalam Dockerfile contoh berasal dari komunitas open source. Gantilah sesuai kebutuhan dan sesuaikan versi SDK dengan versi Spark Anda.
Gunakan Hadoop OSS SDK
Contoh ini menggunakan Spark 3.5.5 dan Hadoop OSS SDK 3.3.4:
ARG SPARK_IMAGE=spark:3.5.5
FROM ${SPARK_IMAGE}
# Tambahkan dependensi untuk dukungan OSS Aliyun Hadoop
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
Gunakan Hadoop S3 SDK
Contoh ini menggunakan Spark 3.5.5 dan Hadoop S3 SDK 3.3.4:
ARG SPARK_IMAGE=spark:3.5.5
FROM ${SPARK_IMAGE}
# Tambahkan dependensi untuk dukungan S3 AWS Hadoop
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
Gunakan JindoSDK
Contoh ini menggunakan Spark 3.5.5 dan JindoSDK 6.8.0:
ARG SPARK_IMAGE=spark:3.5.5
FROM ${SPARK_IMAGE}
# Tambahkan dependensi untuk dukungan 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
Langkah 3: Simpan kredensial OSS dalam Secret Kubernetes
Simpan ID AccessKey dan AccessKey Secret Anda dalam Secret Kubernetes, bukan menyematkannya dalam manifes SparkApplication. Operator Spark menyuntikkan kunci Secret sebagai variabel lingkungan ke dalam Pod driver dan executor.
FieldhadoopConfpada Langkah 4 menggunakanEnvironmentVariableCredentialsProvider, yang secara otomatis membaca variabel lingkungan ini. Jangan pernah menyematkan kredensial dalam file YAML atau melakukan commit ke kontrol versi.
Gunakan Hadoop OSS SDK
-
Buat
spark-oss-secret.yaml:apiVersion: v1 kind: Secret metadata: name: spark-oss-secret namespace: default stringData: # Ganti <ACCESS_KEY_ID> dengan ID AccessKey Akun Alibaba Cloud Anda. OSS_ACCESS_KEY_ID: <ACCESS_KEY_ID> # Ganti <ACCESS_KEY_SECRET> dengan AccessKey Secret Akun Alibaba Cloud Anda. OSS_ACCESS_KEY_SECRET: <ACCESS_KEY_SECRET> -
Terapkan Secret:
kubectl apply -f spark-oss-secret.yamlOutput yang diharapkan:
secret/spark-oss-secret created
Gunakan Hadoop S3 SDK
Hadoop S3 SDK membaca kredensial dari AWS_ACCESS_KEY_ID dan AWS_SECRET_ACCESS_KEY, sehingga nama kuncinya berbeda dari SDK OSS.
-
Buat
spark-s3-secret.yaml:apiVersion: v1 kind: Secret metadata: name: spark-s3-secret namespace: default stringData: # Ganti <ACCESS_KEY_ID> dengan ID AccessKey Akun Alibaba Cloud Anda. AWS_ACCESS_KEY_ID: <ACCESS_KEY_ID> # Ganti <ACCESS_KEY_SECRET> dengan AccessKey Secret Akun Alibaba Cloud Anda. AWS_SECRET_ACCESS_KEY: <ACCESS_KEY_SECRET> -
Terapkan Secret:
kubectl apply -f spark-s3-secret.yamlOutput yang diharapkan:
secret/spark-s3-secret created
Gunakan JindoSDK
-
Buat
spark-oss-secret.yaml:apiVersion: v1 kind: Secret metadata: name: spark-oss-secret namespace: default stringData: # Ganti <ACCESS_KEY_ID> dengan ID AccessKey Akun Alibaba Cloud Anda. OSS_ACCESS_KEY_ID: <ACCESS_KEY_ID> # Ganti <ACCESS_KEY_SECRET> dengan AccessKey Secret Akun Alibaba Cloud Anda. OSS_ACCESS_KEY_SECRET: <ACCESS_KEY_SECRET> -
Terapkan Secret:
kubectl apply -f spark-oss-secret.yamlOutput yang diharapkan:
secret/spark-oss-secret created
Langkah 4: Kirim pekerjaan Spark
Buat dan kirim manifes SparkApplication untuk menjalankan pekerjaan PageRank terhadap set data OSS Anda.
Gunakan Hadoop OSS SDK
Buat spark-pagerank.yaml. Modul Hadoop-Aliyun menjelaskan semua parameter OSS yang tersedia.
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: spark-pagerank
namespace: default
spec:
type: Scala
mode: cluster
# Ganti <SPARK_IMAGE> dengan gambar kontainer Spark yang dibuat di Langkah 2.
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 # Tentukan set data uji input. Ganti <OSS_BUCKET> dengan nama bucket OSS Anda.
- "10" # Jumlah iterasi.
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
# Titik akhir OSS. Ganti <OSS_ENDPOINT> dengan titik akhir OSS Anda.
# Misalnya, titik akhir internal untuk wilayah China (Beijing) adalah 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 # Tentukan Secret yang digunakan untuk mengakses OSS.
serviceAccount: spark-operator-spark
executor:
instances: 2
cores: 1
coreLimit: "2"
memory: 8g
envFrom:
- secretRef:
name: spark-oss-secret # Tentukan Secret yang digunakan untuk mengakses OSS.
restartPolicy:
type: Never
Gunakan Hadoop S3 SDK
Buat spark-pagerank.yaml. Modul Hadoop-AWS menjelaskan semua parameter S3 yang tersedia.
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: spark-pagerank
namespace: default
spec:
type: Scala
mode: cluster
# Ganti <SPARK_IMAGE> dengan gambar kontainer Spark yang dibuat di Langkah 2.
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 # Tentukan set data uji input. Ganti <OSS_BUCKET> dengan nama bucket OSS Anda.
- "10" # Jumlah iterasi.
sparkVersion: 3.5.5
hadoopConf:
fs.s3a.impl: org.apache.hadoop.fs.s3a.S3AFileSystem
# Titik akhir OSS. Ganti <OSS_ENDPOINT> dengan titik akhir OSS Anda.
# Misalnya, titik akhir internal untuk wilayah China (Beijing) adalah oss-cn-beijing-internal.aliyuncs.com.
fs.s3a.endpoint: <OSS_ENDPOINT>
# Wilayah tempat titik akhir OSS berada. Misalnya, cn-beijing untuk wilayah China (Beijing).
fs.s3a.endpoint.region: <OSS_REGION>
driver:
cores: 1
coreLimit: 1200m
memory: 512m
envFrom:
- secretRef:
name: spark-s3-secret # Tentukan Secret yang digunakan untuk mengakses OSS.
serviceAccount: spark-operator-spark
executor:
instances: 2
cores: 1
coreLimit: "2"
memory: 8g
envFrom:
- secretRef:
name: spark-s3-secret # Tentukan Secret yang digunakan untuk mengakses OSS.
restartPolicy:
type: Never
Gunakan JindoSDK
Buat spark-pagerank.yaml:
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: spark-pagerank
namespace: default
spec:
type: Scala
mode: cluster
# Ganti <SPARK_IMAGE> dengan gambar kontainer Spark yang dibuat di Langkah 2.
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 # Tentukan set data uji input. Ganti <OSS_BUCKET> dengan nama bucket OSS Anda.
- "10" # Jumlah iterasi.
sparkVersion: 3.5.5
hadoopConf:
fs.AbstractFileSystem.oss.impl: com.aliyun.jindodata.oss.JindoOSS
fs.oss.impl: com.aliyun.jindodata.oss.JindoOssFileSystem
# Titik akhir OSS. Ganti <OSS_ENDPOINT> dengan titik akhir OSS Anda.
# Misalnya, titik akhir internal untuk wilayah China (Beijing) adalah 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 # Tentukan Secret yang digunakan untuk mengakses OSS.
executor:
instances: 2
cores: 1
coreLimit: "2"
memory: 8g
envFrom:
- secretRef:
name: spark-oss-secret # Tentukan Secret yang digunakan untuk mengakses OSS.
restartPolicy:
type: Never
Kirim dan verifikasi pekerjaan
Perintah-perintah berikut berlaku untuk ketiga opsi SDK.
-
Kirim pekerjaan:
kubectl apply -f spark-pagerank.yaml -
Monitor status pekerjaan:
kubectl get sparkapplications spark-pagerankOutput yang diharapkan saat pekerjaan selesai:
NAME STATUS ATTEMPTS START FINISH AGE spark-pagerank COMPLETED 1 2024-10-09T12:54:25Z 2024-10-09T12:55:46Z 90s -
Lihat 20 baris log terakhir dari Pod driver:
kubectl logs spark-pagerank-driver --tail=20Hadoop OSS SDK — output yang diharapkan:
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 — output yang diharapkan:
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 — output yang diharapkan:
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
(Opsional) Langkah 5: Bersihkan
Hapus pekerjaan Spark dan Secret jika tidak lagi diperlukan.
Hapus pekerjaan Spark:
kubectl delete -f spark-pagerank.yaml
Hapus Secret:
Hadoop OSS SDK atau JindoSDK:
kubectl delete -f spark-oss-secret.yaml
Hadoop S3 SDK:
kubectl delete -f spark-s3-secret.yaml