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 |
|
|
Nome do projeto SLS |
|
|
|
Nome do Logstore no projeto |
|
|
|
Endpoint do SLS para sua região |
|
|
|
AccessKey ID do seu usuário RAM |
Use |
|
|
AccessKey secret do seu usuário RAM |
Use |
|
|
Caminho de saída para o resultado do RDD |
|
|
|
Início do intervalo de tempo |
|
|
|
(Opcional) Fim do intervalo de tempo. O padrão é o horário atual. |
|
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.
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
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 |
|
|
Sim |
Endpoint do SLS para sua região |
|
|
|
Sim |
AccessKey ID, passado via |
|
|
|
Sim |
AccessKey secret, passado via |
|
|
|
Sim |
Nome do projeto SLS |
|
|
|
Sim |
Nome do Logstore |
|
|
|
Não |
Posição inicial de leitura. Use |
|
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.