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
NotaNo 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.