Todos os produtos
Search
Central de documentação

Simple Log Service:Use o Spark Streaming para consumir dados de log

Última atualização: Jul 03, 2026

O Simple Log Service (SLS) oferece suporte ao Spark Streaming para consumo de logs em tempo real. Use o Alibaba Cloud Spark SDK para consumir dados de log do SLS no modo Receiver ou Direct.

Limites

O consumo com Spark Streaming é compatível apenas com o Spark 2.x.

O Alibaba Cloud Spark SDK oferece dois modos de consumo: Receiver e Direct. Adicione a seguinte dependência Maven ao seu projeto:

<dependency>
  <groupId>com.aliyun.emr</groupId>
  <artifactId>emr-logservice_2.11</artifactId>
  <version>1.7.2</version>
</dependency>

Consumir dados de log no modo Receiver

No modo Receiver, um grupo de consumidores obtém dados do SLS e os armazena temporariamente em um executor do Spark. Após o início do job do Spark Streaming, o executor lê e processa esses dados. Cada entrada de log retorna como uma string JSON. O grupo de consumidores salva checkpoints automaticamente no SLS. Para mais informações, consulte Usar grupos de consumidores para consumir dados de log.

  • Parâmetros

    Parâmetro

    Tipo

    Descrição

    project

    String

    Nome do projeto no SLS.

    logstore

    String

    Nome do logstore no SLS.

    consumerGroup

    String

    Nome do grupo de consumidores.

    endpoint

    String

    Endpoint da região onde reside o projeto do SLS. Para mais informações, consulte Endpoints.

    accessKeyId

    String

    AccessKey ID usado para acessar o SLS.

    accessKeySecret

    String

    AccessKey secret usado para acessar o SLS.

  • Exemplo

    Nota

    No modo Receiver, pode ocorrer perda de dados com as configurações padrão. Para evitar isso, ative os Write-Ahead Logs (disponíveis no Spark 1.2 e versões posteriores). Para mais detalhes, consulte Spark.

    import org.apache.spark.storage.StorageLevel
    import org.apache.spark.streaming.aliyun.logservice.LoghubUtils
    import org.apache.spark.streaming.{ Milliseconds, StreamingContext}
    import org.apache.spark.SparkConf
    
    object TestLoghub {
      def main(args: Array[String]): Unit = {
        if (args.length < 7) {
          System.err.println(
            """Usage: TestLoghub <project> <logstore> <loghub group name> <endpoint>
              |         <access key id> <access key secret> <batch interval seconds>
            """.stripMargin)
          System.exit(1)
        }
    
        val project = args(0)
        val logstore = args(1)
        val consumerGroup = args(2)
        val endpoint = args(3)
        val accessKeyId = args(4)
        val accessKeySecret = args(5)
        val batchInterval = Milliseconds(args(6).toInt * 1000)
    
        def functionToCreateContext(): StreamingContext = {
          val conf = new SparkConf().setAppName("Test Loghub")
          val ssc = new StreamingContext(conf, batchInterval)
          val loghubStream = LoghubUtils.createStream(
            ssc,
            project,
            logstore,
            consumerGroup,
            endpoint,
            accessKeyId,
            accessKeySecret,
            StorageLevel.MEMORY_AND_DISK)
    
          loghubStream.checkpoint(batchInterval * 2).foreachRDD(rdd =>
            rdd.map(bytes => new String(bytes)).top(10).foreach(println)
          )
          ssc.checkpoint("hdfs:///tmp/spark/streaming") // set checkpoint directory
          ssc
        }
    
        val ssc = StreamingContext.getOrCreate("hdfs:///tmp/spark/streaming", functionToCreateContext _)
    
        ssc.start()
        ssc.awaitTermination()
      }
    }

Consumir dados de log no modo Direct

O modo Direct não exige um grupo de consumidores. Ele chama operações de API para solicitar dados diretamente do SLS. Esse modo oferece os seguintes benefícios:

  • Concorrência simplificada: o número de partições do Spark corresponde ao número de shards no Logstore. Divida os shards para aumentar a concorrência das tarefas.

  • Maior eficiência: não é necessário usar Write-Ahead Logs para evitar perda de dados.

  • Semântica exactly-once: os dados são lidos diretamente do SLS e os checkpoints são confirmados somente após o sucesso da tarefa.

    Se o Spark encerrar inesperadamente, alguns dados poderão ser consumidos mais de uma vez.

O modo Direct requer um serviço ZooKeeper para armazenar o estado intermediário. Defina um diretório de checkpoint no ZooKeeper para persistir os dados intermediários. Para reconsumir dados após reiniciar uma tarefa, exclua o diretório correspondente no ZooKeeper e altere o nome do grupo de consumidores.

  • Parâmetros

    Parâmetro

    Tipo

    Descrição

    project

    String

    Nome do projeto no SLS.

    logstore

    String

    Nome do logstore no SLS.

    consumerGroup

    String

    Nome do grupo de consumidores. Este nome serve apenas para salvar checkpoints de consumo.

    endpoint

    String

    Endpoint da região onde reside o projeto do SLS. Para mais informações, consulte Endpoints.

    accessKeyId

    String

    AccessKey ID usado para acessar o SLS.

    accessKeySecret

    String

    AccessKey secret usado para acessar o SLS.

    zkAddress

    String

    URL de conexão do serviço ZooKeeper.

  • Limitação de taxa

    O Spark Streaming processa dados em microlotes. Especifique a quantidade de entradas de log consumidas por lote e por shard.

    No SLS, cada solicitação de gravação é armazenada como um grupo de logs. Uma solicitação típica contém várias entradas em um único grupo. Quando o rastreamento web está ativo, cada solicitação inclui apenas uma entrada por grupo. Os parâmetros abaixo controlam o volume de dados consumidos em um único lote.

    Parâmetro

    Descrição

    Padrão

    spark.loghub.batchGet.step

    Quantidade máxima de grupos de logs buscados por solicitação de consumo.

    100

    spark.streaming.loghub.maxRatePerShard

    Limite máximo de entradas de log consumidas por shard em cada lote.

    10000

    O parâmetro spark.streaming.loghub.maxRatePerShard define o alvo máximo de entradas de log por shard e por lote. O SDK busca grupos de logs em incrementos de spark.loghub.batchGet.step e acumula a contagem de entradas. Ao atingir ou superar o valor de spark.streaming.loghub.maxRatePerShard, o SDK interrompe a busca. Como a granularidade do consumo ocorre no nível do grupo de logs, o limite de spark.streaming.loghub.maxRatePerShard é aproximado. A contagem real por lote depende de spark.loghub.batchGet.step e da quantidade de entradas em cada grupo de logs.

  • Exemplo

    import com.aliyun.openservices.loghub.client.config.LogHubCursorPosition
    import org.apache.spark.SparkConf
    import org.apache.spark.streaming.{ Milliseconds, StreamingContext}
    import org.apache.spark.streaming.aliyun.logservice.{ CanCommitOffsets, LoghubUtils}
    
    object TestDirectLoghub {
      def main(args: Array[String]): Unit = {
        if (args.length < 7) {
          System.err.println(
            """Usage: TestDirectLoghub <project> <logstore> <loghub group name> <endpoint>
              |         <access key id> <access key secret> <batch interval seconds> <zookeeper host:port=localhost:2181>
            """.stripMargin)
          System.exit(1)
        }
    
        val project = args(0)
        val logstore = args(1)
        val consumerGroup = args(2)
        val endpoint = args(3)
        val accessKeyId = args(4)
        val accessKeySecret = args(5)
        val batchInterval = Milliseconds(args(6).toInt * 1000)
        val zkAddress = if (args.length >= 8) args(7) else "localhost:2181"
    
        def functionToCreateContext(): StreamingContext = {
          val conf = new SparkConf().setAppName("Test Direct Loghub")
          val ssc = new StreamingContext(conf, batchInterval)
          val zkParas = Map("zookeeper.connect" -> zkAddress,
            "enable.auto.commit" -> "false")
          val loghubStream = LoghubUtils.createDirectStream(
            ssc,
            project,
            logStore,
            consumerGroup,
            accessKeyId,
            accessKeySecret,
            endpoint,
            zkParas,
            LogHubCursorPosition.END_CURSOR)
    
          loghubStream.checkpoint(batchInterval).foreachRDD(rdd => {
            println(s"count by key: ${rdd.map(s => {
              s.sorted
              (s.length, s)
            }).countByKey().size}")
            loghubStream.asInstanceOf[CanCommitOffsets].commitAsync()
          })
          ssc.checkpoint("hdfs:///tmp/spark/streaming") // set checkpoint directory
          ssc
        }
    
        val ssc = StreamingContext.getOrCreate("hdfs:///tmp/spark/streaming", functionToCreateContext _)
        ssc.start()
        ssc.awaitTermination()
      }
    }

Para acessar o código-fonte completo, consulte o GitHub.