Todos os produtos
Search
Central de documentação

E-MapReduce:Exemplo de consumo offline com Spark

Última atualização: Jun 27, 2026

Use o Spark para ler dados do Simple Log Service (SLS) em modo batch (offline), tratando-os como um conjunto de dados delimitado com horários de início e fim definidos, em vez de um fluxo contínuo. O EMR oferece suporte a duas abordagens: Spark RDD e Spark SQL.

Pré-requisitos

Antes de começar, verifique se você tem:

  • Um cluster EMR com o Spark instalado

  • Um projeto e um Logstore do Simple Log Service com os dados a serem lidos

  • Um AccessKey ID e um AccessKey secret de um usuário RAM com permissão de leitura no Logstore

  • (Apenas para Spark SQL) Acesso ao nó do cluster EMR onde estão os arquivos JAR

Usar Spark RDD para ler do Simple Log Service

Código de exemplo

O programa Scala a seguir lê dados de log de um Logstore para um intervalo de tempo especificado e salva a saída em um caminho do sistema de arquivos.

// TestBatchLoghub.scala

object TestBatchLoghub {
  def main(args: Array[String]): Unit = {
    if (args.length < 6) {
      System.err.println(
        """Usage: TestBatchLoghub <sls project> <sls logstore> <sls endpoint>
          |  <access key id> <access key secret> <output path> <start time> <end time=now>
        """.stripMargin)
      System.exit(1)
    }

    val loghubProject    = args(0)
    val logStore         = args(1)
    val endpoint         = args(2)
    val accessKeyId      = args(3)     // Read from environment variable at submit time
    val accessKeySecret  = args(4)     // Read from environment variable at submit time
    val outputPath       = args(5)
    val startTime        = args(6).toLong

    val sc = new SparkContext(new SparkConf().setAppName("test batch loghub"))
    var rdd: JavaRDD[String] = null

    if (args.length > 7) {
      // Read log data between startTime and endTime
      rdd = LoghubUtils.createRDD(sc, loghubProject, logStore, accessKeyId, accessKeySecret, endpoint, startTime, args(7).toLong)
    } else {
      // Read log data from startTime to now
      rdd = LoghubUtils.createRDD(sc, loghubProject, logStore, accessKeyId, accessKeySecret, endpoint, startTime)
    }

    rdd.saveAsTextFile(outputPath)
  }
}

A função LoghubUtils.createRDD() retorna um JavaRDD[String] em que cada elemento corresponde a uma entrada de log. Passe sete argumentos para ler até o horário atual ou oito argumentos para definir explicitamente um horário de término.

Para obter a configuração do Maven POM, consulte aliyun-emapreduce-demo.

Parâmetros de conexão

Parâmetro

Descrição

Exemplo

<sls project>

Nome do projeto SLS

my-project

<sls logstore>

Nome do Logstore no projeto

my-logstore

<sls endpoint>

Endpoint do SLS para sua região

cn-hangzhou.log.aliyuncs.com

<access key id>

AccessKey ID do seu usuário RAM

Use $ALIBABA_CLOUD_ACCESS_KEY_ID

<access key secret>

AccessKey secret do seu usuário RAM

Use $ALIBABA_CLOUD_ACCESS_KEY_SECRET

<output path>

Caminho de saída para o resultado do RDD

oss://my-bucket/output/

<start time>

Início do intervalo de tempo

1700000000

<end time>

(Opcional) Fim do intervalo de tempo. O padrão é o horário atual.

1700003600

Compilar e executar

Etapa 1: Compile o código.

mvn clean package -DskipTests

O JAR compilado é salvo no diretório target/shaded/.

Etapa 2: Envie o job.

Importante

Use o par de AccessKey de um usuário RAM em vez do par de AccessKey da sua conta Alibaba Cloud. Um par de AccessKey no nível da conta concede acesso a todas as operações de API. Para saber como criar um usuário RAM, consulte Criar um usuário RAM. Armazene as credenciais como variáveis de ambiente; não as codifique diretamente nos scripts.

Defina as variáveis de ambiente antes de enviar:

export ALIBABA_CLOUD_ACCESS_KEY_ID=<access_key_id>
export ALIBABA_CLOUD_ACCESS_KEY_SECRET=<access_key_secret>

Em seguida, envie o job:

spark-submit \
  --master yarn-cluster \
  --executor-cores 2 \
  --executor-memory 1g \
  --driver-memory 1g \
  --num-executors 2 \
  --class x.x.x.TestBatchLoghub xxx.jar \
  <sls project> <sls logstore> <sls endpoint> \
  $ALIBABA_CLOUD_ACCESS_KEY_ID $ALIBABA_CLOUD_ACCESS_KEY_SECRET \
  <output path> <start time> [<end time>]

Substitua x.x.x.TestBatchLoghub pelo nome totalmente qualificado da classe real e xxx.jar pelo caminho real do JAR. Ajuste --executor-cores, --executor-memory e --num-executors conforme o volume de dados e a capacidade do cluster.

Usar Spark SQL para ler do Simple Log Service

O Spark SQL usa a source de dados loghub, incluída no JAR de extensão do Spark do EMR.

Iniciar spark-sql com o JAR do LogHub

Importante

Use o par de AccessKey de um usuário RAM em vez do par de AccessKey da sua conta Alibaba Cloud. Armazene as credenciais como variáveis de ambiente passadas via --hiveconf; não as codifique diretamente nos scripts.

spark-sql \
  --jars /opt/apps/SPARK-EXTENSION/spark-extension-current/spark3-emrsdk/* \
  --hiveconf accessKeyId=$ALIBABA_CLOUD_ACCESS_KEY_ID \
  --hiveconf accessKeySecret=$ALIBABA_CLOUD_ACCESS_KEY_SECRET

Se o cluster EMR usar o Spark 2, substitua spark3 por spark2 no caminho do JAR:

/opt/apps/SPARK-EXTENSION/spark-extension-current/spark2-emrsdk/*

Criar uma tabela e consultar dados

CREATE TABLE test_sls
USING loghub
OPTIONS (
  endpoint         = 'cn-hangzhou-intranet.log.aliyuncs.com',
  access.key.id    = '${hiveconf:accessKeyId}',
  access.key.secret= '${hiveconf:accessKeySecret}',
  sls.project      = 'test_project',
  sls.store        = 'test_store',
  startingoffsets  = 'earliest'
);

SELECT * FROM test_sls;

Parâmetros de conexão

Parâmetro

Obrigatório

Descrição

Exemplo

endpoint

Sim

Endpoint do SLS para sua região

cn-hangzhou-intranet.log.aliyuncs.com

access.key.id

Sim

AccessKey ID, passado via --hiveconf

${hiveconf:accessKeyId}

access.key.secret

Sim

AccessKey secret, passado via --hiveconf

${hiveconf:accessKeySecret}

sls.project

Sim

Nome do projeto SLS

test_project

sls.store

Sim

Nome do Logstore

test_store

startingoffsets

Não

Posição inicial de leitura. Use earliest para ler todos os dados disponíveis desde o início do Logstore.

earliest

Usar o JAR do LogHub em um ambiente de desenvolvimento local

Para desenvolver e testar localmente com o Spark 3 (o procedimento é o mesmo para o Spark 2), instale o JAR da source de dados do EMR no repositório Maven local.

Etapa 1: Baixe o JAR do cluster EMR.

Copie o JAR do seguinte caminho no nó do cluster para a máquina local:

/opt/apps/SPARK-EXTENSION/spark-extension-current/spark3-emrsdk/emr-datasources_shaded_2.12

Etapa 2: Instale o JAR no repositório Maven local.

mvn install:install-file \
  -DgroupId=com.aliyun.emr \
  -DartifactId=emr-datasources_shaded_2.12 \
  -Dversion=3.0.2 \
  -Dpackaging=jar \
  -Dfile=<path-to-downloaded-jar>

Substitua <path-to-downloaded-jar> pelo caminho local onde o JAR foi salvo.

**Etapa 3: Adicione a dependência ao pom.xml.**

<dependency>
  <groupId>com.aliyun.emr</groupId>
  <artifactId>emr-datasources_shaded_2.12</artifactId>
  <version>3.0.2</version>
</dependency>

Referências

Para mais informações sobre como usar o Spark para acessar o Kafka, consulte Guia de integração do Structured Streaming + Kafka.