Todos os produtos
Search
Central de documentação

Container Service for Kubernetes:Read and write OSS data in Spark jobs

Última atualização: Jun 27, 2026

Execute jobs do Spark para ler e gravar dados no OSS em clusters ACK. Utilize o job PageRank integrado com Hadoop OSS SDK, Hadoop S3 SDK ou JindoSDK.

Pré-requisitos

Verifique se você tem:

Escolha uma integração com o OSS

Três SDKs estão disponíveis para acesso ao OSS. Escolha um antes de prosseguir.

SDK

Quando usar

Esquema de URI

Hadoop OSS SDK

Suporte nativo ao OSS via sistema de arquivos Aliyun; simples para a maioria das cargas de trabalho

oss://

Hadoop S3 SDK

Seu cluster ou ferramentas utilizam APIs compatíveis com S3; sem dependência específica da Aliyun

s3a://

JindoSDK

Desempenho otimizado do OSS com aceleração Jindo

oss://

Use o mesmo SDK em todas as etapas.

Etapa 1: Preparar e carregar dados de teste

Gere um conjunto de dados PageRank e carregue-o no seu bucket do OSS.

  1. Crie o arquivo generate_pagerank_dataset.sh com o seguinte conteúdo:

    #!/bin/bash
    
    # Check the number of arguments
    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
    
    # Verify if M and N are positive integers
    if ! [[ "$M" =~ ^[0-9]+$ ]] || ! [[ "$N" =~ ^[0-9]+$ ]]; then
        echo "Both M and N must be positive integers."
        exit 1
    fi
    
    # Generate dataset
    for ((i=1; i<=$N; i++)); do
        # Ensure the source and target pages are different
        while true; do
            src=$((RANDOM % M + 1))
            dst=$((RANDOM % M + 1))
            if [ "$src" -ne "$dst" ]; then
                echo "$src $dst"
                break
            fi
        done
    done
  2. Gere o conjunto de dados:

    M=100000    # The number of web pages
    
    N=10000000  # The number of records
    
    # Generate dataset randomly and save as pagerank_dataset.txt
    bash generate_pagerank_dataset.sh $M $N > pagerank_dataset.txt
  3. Carregue o conjunto de dados na pasta data/ do seu bucket do OSS:

    ossutil cp pagerank_dataset.txt oss://<BUCKET_NAME>/data/

Etapa 2: Criar uma imagem de contêiner Spark

Crie uma imagem de contêiner com as dependências JAR necessárias para acessar o OSS. Consulte Use a Container Registry Enterprise Edition instance to build an image.

A imagem base do Spark nos Dockerfiles de exemplo vem da comunidade open source. Substitua-a conforme necessário e alinhe a versão do SDK à sua versão do Spark.

Use Hadoop OSS SDK

Este exemplo usa Spark 3.5.5 e Hadoop OSS SDK 3.3.4:

ARG SPARK_IMAGE=spark:3.5.5

FROM ${SPARK_IMAGE}

# Add dependencies for Hadoop Aliyun OSS support
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

Use Hadoop S3 SDK

Este exemplo usa Spark 3.5.5 e Hadoop S3 SDK 3.3.4:

ARG SPARK_IMAGE=spark:3.5.5

FROM ${SPARK_IMAGE}

# Add dependencies for Hadoop AWS S3 support
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

Use JindoSDK

Este exemplo usa Spark 3.5.5 e JindoSDK 6.8.0:

ARG SPARK_IMAGE=spark:3.5.5

FROM ${SPARK_IMAGE}

# Add dependencies for JindoSDK support
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

Etapa 3: Armazenar credenciais do OSS em um Secret do Kubernetes

Armazene seu AccessKey ID e AccessKey Secret em um Secret do Kubernetes em vez de incorporá-los diretamente no manifesto SparkApplication. O operador Spark injeta as chaves do Secret como variáveis de ambiente nos pods do driver e dos executores.

Os campos hadoopConf na Etapa 4 usam EnvironmentVariableCredentialsProvider , que lê essas variáveis de ambiente automaticamente. Nunca incorpore credenciais em arquivos YAML nem as envie para sistemas de controle de versão.

Use Hadoop OSS SDK

  1. Crie o arquivo spark-oss-secret.yaml:

    apiVersion: v1
    kind: Secret
    metadata:
      name: spark-oss-secret
      namespace: default
    stringData:
      # Replace <ACCESS_KEY_ID> with the AccessKey ID of your Alibaba Cloud account.
      OSS_ACCESS_KEY_ID: <ACCESS_KEY_ID>
      # Replace <ACCESS_KEY_SECRET> with the AccessKey Secret of your Alibaba Cloud account.
      OSS_ACCESS_KEY_SECRET: <ACCESS_KEY_SECRET>
  2. Aplique o Secret:

    kubectl apply -f spark-oss-secret.yaml

    Saída esperada:

    secret/spark-oss-secret created

Use Hadoop S3 SDK

O Hadoop S3 SDK lê credenciais das variáveis AWS_ACCESS_KEY_ID e AWS_SECRET_ACCESS_KEY, portanto os nomes das chaves diferem do OSS SDK.

  1. Crie o arquivo spark-s3-secret.yaml:

    apiVersion: v1
    kind: Secret
    metadata:
      name: spark-s3-secret
      namespace: default
    stringData:
      # Replace <ACCESS_KEY_ID> with the AccessKey ID of your Alibaba Cloud account.
      AWS_ACCESS_KEY_ID: <ACCESS_KEY_ID>
      # Replace <ACCESS_KEY_SECRET> with the AccessKey Secret of your Alibaba Cloud account.
      AWS_SECRET_ACCESS_KEY: <ACCESS_KEY_SECRET>
  2. Aplique o Secret:

    kubectl apply -f spark-s3-secret.yaml

    Saída esperada:

    secret/spark-s3-secret created

Use JindoSDK

  1. Crie o arquivo spark-oss-secret.yaml:

    apiVersion: v1
    kind: Secret
    metadata:
      name: spark-oss-secret
      namespace: default
    stringData:
      # Replace <ACCESS_KEY_ID> with the AccessKey ID of your Alibaba Cloud account.
      OSS_ACCESS_KEY_ID: <ACCESS_KEY_ID>
      # Replace <ACCESS_KEY_SECRET> with the AccessKey Secret of your Alibaba Cloud account.
      OSS_ACCESS_KEY_SECRET: <ACCESS_KEY_SECRET>
  2. Aplique o Secret:

    kubectl apply -f spark-oss-secret.yaml

    Saída esperada:

    secret/spark-oss-secret created

Etapa 4: Enviar o job do Spark

Crie e envie um manifesto SparkApplication para executar o job PageRank no seu conjunto de dados do OSS.

Use Hadoop OSS SDK

Crie o arquivo spark-pagerank.yaml. O módulo Hadoop-Aliyun descreve todos os parâmetros disponíveis para o OSS.

apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
  name: spark-pagerank
  namespace: default
spec:
  type: Scala
  mode: cluster
  # Replace <SPARK_IMAGE> with the Spark container image built in Step 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           # Specify the input test dataset. Replace <OSS_BUCKET> with your OSS bucket name.
  - "10"                                                   # The number of iterations.
  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 endpoint. Replace <OSS_ENDPOINT> with your OSS endpoint.
    # For example, the internal endpoint for the China (Beijing) region is 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               # Specify the Secret used to access OSS.
    serviceAccount: spark-operator-spark
  executor:
    instances: 2
    cores: 1
    coreLimit: "2"
    memory: 8g
    envFrom:
    - secretRef:
        name: spark-oss-secret               # Specify the Secret used to access OSS.
  restartPolicy:
    type: Never

Use Hadoop S3 SDK

Crie o arquivo spark-pagerank.yaml. O módulo Hadoop-AWS descreve todos os parâmetros disponíveis para S3.

apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
  name: spark-pagerank
  namespace: default
spec:
  type: Scala
  mode: cluster
  # Replace <SPARK_IMAGE> with the Spark container image built in Step 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           # Specify the input test dataset. Replace <OSS_BUCKET> with your OSS bucket name.
  - "10"                                                   # The number of iterations.
  sparkVersion: 3.5.5
  hadoopConf:
    fs.s3a.impl: org.apache.hadoop.fs.s3a.S3AFileSystem
    # OSS endpoint. Replace <OSS_ENDPOINT> with your OSS endpoint.
    # For example, the internal endpoint for the China (Beijing) region is oss-cn-beijing-internal.aliyuncs.com.
    fs.s3a.endpoint: <OSS_ENDPOINT>
    # The region where the OSS endpoint is located. For example, cn-beijing for the China (Beijing) region.
    fs.s3a.endpoint.region: <OSS_REGION>
  driver:
    cores: 1
    coreLimit: 1200m
    memory: 512m
    envFrom:
    - secretRef:
        name: spark-s3-secret               # Specify the Secret used to access OSS.
    serviceAccount: spark-operator-spark
  executor:
    instances: 2
    cores: 1
    coreLimit: "2"
    memory: 8g
    envFrom:
    - secretRef:
        name: spark-s3-secret               # Specify the Secret used to access OSS.
  restartPolicy:
    type: Never

Use JindoSDK

Crie o arquivo spark-pagerank.yaml:

apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
  name: spark-pagerank
  namespace: default
spec:
  type: Scala
  mode: cluster
  # Replace <SPARK_IMAGE> with the Spark container image built in Step 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    # Specify the input test dataset. Replace <OSS_BUCKET> with your OSS bucket name.
  - "10"                                            # The number of iterations.
  sparkVersion: 3.5.5
  hadoopConf:
    fs.AbstractFileSystem.oss.impl: com.aliyun.jindodata.oss.JindoOSS
    fs.oss.impl: com.aliyun.jindodata.oss.JindoOssFileSystem
    # OSS endpoint. Replace <OSS_ENDPOINT> with your OSS endpoint.
    # For example, the internal endpoint for the China (Beijing) region is 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                    # Specify the Secret used to access OSS.
  executor:
    instances: 2
    cores: 1
    coreLimit: "2"
    memory: 8g
    envFrom:
    - secretRef:
        name: spark-oss-secret                    # Specify the Secret used to access OSS.
  restartPolicy:
    type: Never

Enviar e verificar o job

Estes comandos aplicam-se às três opções de SDK.

  1. Envie o job:

    kubectl apply -f spark-pagerank.yaml
  2. Monitore o status do job:

    kubectl get sparkapplications spark-pagerank

    Saída esperada quando o job for concluído:

    NAME             STATUS      ATTEMPTS   START                  FINISH                 AGE
    spark-pagerank   COMPLETED   1          2024-10-09T12:54:25Z   2024-10-09T12:55:46Z   90s
  3. Visualize as últimas 20 linhas de log do pod driver:

    kubectl logs spark-pagerank-driver --tail=20

    Hadoop OSS SDK — saída esperada:

    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-2752eb83b2c5

    Hadoop S3 SDK — saída esperada:

    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 — saída esperada:

    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

(Opcional) Etapa 5: Limpar recursos

Exclua o job do Spark e o Secret quando não forem mais necessários.

Exclua o job do Spark:

kubectl delete -f spark-pagerank.yaml

Exclua o Secret:

Hadoop OSS SDK ou JindoSDK:

kubectl delete -f spark-oss-secret.yaml

Hadoop S3 SDK:

kubectl delete -f spark-s3-secret.yaml

Próximas etapas