É possível enviar um job para um cluster do Serviço Gerenciado para Apache Spark usando uma
jobs.submit ou uma solicitação programática, usando a ferramenta de linha de comando
gcloudGoogle Cloud CLI em uma
janela do terminal local ou no Cloud Shell ou
pelo Google Cloud console aberto em
um navegador local. Também é possível executar SSH na instância mestre
do cluster e executar um job diretamente na instância sem usar o Serviço Gerenciado para Apache Spark.
Simultaneidade do job
é possível configurar o número máximo de jobs simultâneos do Serviço Gerenciado para Apache Spark
com a propriedade
dataproc:dataproc.scheduler.max-concurrent-jobs
quando você cria um cluster. Se esse valor de propriedade não estiver definido, o limite superior em jobs simultâneos será calculado como max((masterMemoryMb - 3584) / masterMemoryMbPerJob, 5).
masterMemoryMb é determinado pelo tipo de máquina da VM mestre.
masterMemoryMbPerJob é 1024 por padrão, mas é
configurável na criação do cluster com a
dataproc:dataproc.scheduler.driver-size-mb.
Como enviar um job
Isolamento de classpath e JARs personalizados: não copie "JARs gordos" personalizados ou JARs de pacote (como os orbundles de execução do Apache Iceberg ou Google Cloud bundles) diretamente em diretórios do sistema de cluster, como /usr/lib/spark/jars/. A colocação de JARs personalizados em diretórios do sistema polui o classpath do agente do Serviço Gerenciado para Apache Spark com dependências transitivas (como bibliotecas Guava ou Hadoop) que entram em conflito com as bibliotecas integradas do agente. Esse conflito pode causar erros de resolução de classe (como ClassNotFoundException ou NoClassDefFoundError) ao realizar operações de gerenciamento de jobs, como cancelar um job, fazendo com que o agente falhe e deixando aplicativos YARN órfãos.
Em vez disso, use uma das seguintes abordagens compatíveis:
- Dependências específicas do job: especifique o caminho do Cloud Storage para seus JARs ao enviar o job usando a flag
--jarsna CLI do Google Cloud, o campo Arquivos JAR no Google Cloud console ou o campojarFileUrisna API. O Spark distribui as dependências para o driver e os executores do job sem poluir o classpath do agente. - Dependências em todo o cluster:especifique as dependências durante a criação do cluster usando a
spark:spark.jarspropriedade de cluster: Isso instrui o Serviço Gerenciado para Apache Spark a configurar--properties="spark:spark.jars=gs://YOUR_BUCKET/jar-1.jar,gs://YOUR_BUCKET/jar-2.jar"
/etc/spark/conf/spark-defaults.conf, distribuindo automaticamente as dependências para os classpaths do driver e do executor do Spark para todos os jobs, deixando o agente do Serviço Gerenciado para Apache Spark local do nó isolado.
Console
Abra a página Enviar um job do Serviço Gerenciado para Apache Spark no Google Cloud console no seu navegador.
Exemplo de job do Spark
Para enviar um job do Spark de exemplo, preencha os campos na página Enviar um job da seguinte maneira:
- Selecione o nome do Cluster na lista de clusters.
- Defina o tipo de serviço como
Spark. - Defina Classe principal ou jar como
org.apache.spark.examples.SparkPi. - Defina Argumentos como o argumento único
1000. - Adicione
file:///usr/lib/spark/examples/jars/spark-examples.jara Arquivos JAR (ou o campojarFileUrisda API):file:///indica um esquema de LocalFileSystem do Hadoop. O Serviço Gerenciado para Apache Spark instalou/usr/lib/spark/examples/jars/spark-examples.jarno nó mestre do cluster quando criou o cluster. Esse caminho é usado apenas para exemplos pré-instalados fornecidos pelo Serviço Gerenciado para Apache Spark. Não copie JARs personalizados em diretórios do sistema (consulte o aviso em Como enviar um job).- Como alternativa, você pode especificar um caminho do Cloud Storage
(
gs://your-bucket/your-jarfile.jar) ou um caminho do sistema de arquivos distribuídos do Hadoop (hdfs://path-to-jar.jar) para um dos seus jars. Se você enviar o job usando a API, especifique esse caminho no campojarFileUris.
Clique em Enviar para iniciar o job. Depois de iniciado, o job será adicionado à lista.
Clique no código da tarefa para abrir a página Jobs e ver a saída do driver do job. Como esse job produz linhas de saída longas que excedem a largura da janela do navegador, marque a caixa Quebra de linha para exibir todo o texto de saída e mostrar o resultado calculado para pi.
Visualize a saída do driver de job na linha de comando usando o
comando gcloud dataproc jobs wait
mostrado abaixo. Para mais informações, consulte
Conferir a saída do job – COMANDO GCLOUD.
Copie e cole o ID do projeto como o valor para a flag --project e o código da tarefa (mostrado na lista de jobs) como o argumento final.
gcloud dataproc jobs wait job-id \ --project=project-id \ --region=region
Confira os snippets da saída do driver do job de exemplo SparkPi:
... 2015-06-25 23:27:23,810 INFO [dag-scheduler-event-loop] scheduler.DAGScheduler (Logging.scala:logInfo(59)) - Stage 0 (reduce at SparkPi.scala:35) finished in 21.169 s 2015-06-25 23:27:23,810 INFO [task-result-getter-3] cluster.YarnScheduler (Logging.scala:logInfo(59)) - Removed TaskSet 0.0, whose tasks have all completed, from pool 2015-06-25 23:27:23,819 INFO [main] scheduler.DAGScheduler (Logging.scala:logInfo(59)) - Job 0 finished: reduce at SparkPi.scala:35, took 21.674931 s Pi is roughly 3.14189648 ... Job [c556b47a-4b46-4a94-9ba2-2dcee31167b2] finished successfully. driverOutputUri: gs://sample-staging-bucket/google-cloud-dataproc-metainfo/cfeaa033-749e-48b9-... ...
gcloud
Para enviar um job a um cluster do Serviço Gerenciado para Apache Spark, execute o comando gcloud dataproc submit da CLI gcloud gcloud dataproc jobs submit localmente em uma janela de terminal ou no Cloud Shell.
gcloud dataproc jobs submit job-command \ --cluster=cluster-name \ --region=region \ other dataproc-flags \ -- job-args
- Liste o
hello-world.pyacessível publicamente localizado no Cloud Storage. Listagem de arquivos:gcloud storage cat gs://dataproc-examples/pyspark/hello-world/hello-world.py
#!/usr/bin/python import pyspark sc = pyspark.SparkContext() rdd = sc.parallelize(['Hello,', 'world!']) words = sorted(rdd.collect()) print(words)
- Envie o job do PySpark para o Serviço Gerenciado para Apache Spark.
Saída do terminal:gcloud dataproc jobs submit pyspark \ gs://dataproc-examples/pyspark/hello-world/hello-world.py \ --cluster=cluster-name \ --region=region
Waiting for job output... … ['Hello,', 'world!'] Job finished successfully.
- Execute o exemplo SparkPi pré-instalado no nó mestre do cluster do Serviço Gerenciado para Apache Spark. O caminho
file:///usr/lib/spark/examples/jars/spark-examples.jaré apenas para exemplos pré-instalados. Para dependências de JAR personalizadas, consulte o aviso em Como enviar um job. Saída do terminal:gcloud dataproc jobs submit spark \ --cluster=cluster-name \ --region=region \ --class=org.apache.spark.examples.SparkPi \ --jars=file:///usr/lib/spark/examples/jars/spark-examples.jar \ -- 1000
Job [54825071-ae28-4c5b-85a5-58fae6a597d6] submitted. Waiting for job output… … Pi is roughly 3.14177148 … Job finished successfully. …
REST
Nesta seção, mostramos como enviar um job do Spark para calcular o valor aproximado
de pi usando o Serviço Gerenciado para Apache Spark
jobs.submit API.
Antes de usar qualquer um dos dados da solicitação, faça as seguintes substituições:
- project-id: Google Cloud ID do projeto
- region: região do cluster
- clusterName: nome do cluster
Método HTTP e URL:
POST https://dataproc.googleapis.com/v1/projects/project-id/regions/region/jobs:submit
Corpo JSON da solicitação:
{
"job": {
"placement": {
"clusterName": "cluster-name"
},
"sparkJob": {
"args": [
"1000"
],
"mainClass": "org.apache.spark.examples.SparkPi",
"jarFileUris": [
"file:///usr/lib/spark/examples/jars/spark-examples.jar"
]
}
}
}
Para enviar a solicitação, expanda uma destas opções:
Você receberá uma resposta JSON semelhante a esta:
{
"reference": {
"projectId": "project-id",
"jobId": "job-id"
},
"placement": {
"clusterName": "cluster-name",
"clusterUuid": "cluster-Uuid"
},
"sparkJob": {
"mainClass": "org.apache.spark.examples.SparkPi",
"args": [
"1000"
],
"jarFileUris": [
"file:///usr/lib/spark/examples/jars/spark-examples.jar"
]
},
"status": {
"state": "PENDING",
"stateStartTime": "2020-10-07T20:16:21.759Z"
},
"jobUuid": "job-Uuid"
}
Java
Python
Go
Node.js
Enviar um trabalho diretamente no cluster
Se você quiser executar um job diretamente no cluster sem usar o Serviço Gerenciado para Apache Spark, use o SSH no nó mestre do cluster e execute o job no nó mestre.
Depois de estabelecer uma conexão SSH com a instância mestre de VM, execute comandos em uma janela de terminal no nó mestre do cluster para:
- abrir um shell do Spark;
- executar um job do Spark para contar o número de linhas em um arquivo "hello-world" do Python (sete linhas) localizado em um arquivo do Cloud Storage acessível publicamente;
sair do shell.
user@cluster-name-m:~$ spark-shell ... scala> sc.textFile("gs://dataproc-examples" + "/pyspark/hello-world/hello-world.py").count ... res0: Long = 7 scala> :quit
Executar jobs do Bash no Serviço Gerenciado para Apache Spark
Execute um script bash como job do Serviço Gerenciado para Apache Spark, seja porque os mecanismos usados não são compatíveis com um tipo de serviço de nível superior ou porque você precisa configurar ou calcular os argumentos adicionais antes de iniciar um job usando hadoop ou spark-submit do seu script.
Exemplo de Python
Suponha que você tenha copiado um script hello.sh bash no Cloud Storage:
gcloud storage cp hello.sh gs://${BUCKET}/hello.shComo o comando pig fs usa caminhos do Hadoop, copie o script do
Cloud Storage para um destino especificado como file:/// para garantir
que ele esteja no sistema de arquivos local em vez do HDFS. Os comandos sh subsequentes fazem referência ao sistema de arquivos local automaticamente e não exigem o prefixo file:///.
gcloud dataproc jobs submit pig --cluster=${CLUSTER} --region=${REGION} \
-e='fs -cp -f gs://${BUCKET}/hello.sh file:///tmp/hello.sh; sh chmod 750 /tmp/hello.sh; sh /tmp/hello.sh'Como alternativa, como os jobs do Serviço Gerenciado para Apache Spark enviam o argumento --jars e em um diretório temporário criado durante a vida útil do job, especifique o script de shell do Cloud Storage como um argumento --jars:
gcloud dataproc jobs submit pig --cluster=${CLUSTER} --region=${REGION} \
--jars=gs://${BUCKET}/hello.sh \
-e='sh chmod 750 ${PWD}/hello.sh; sh ${PWD}/hello.sh'Observe que o argumento --jars também pode fazer referência a um script local:
gcloud dataproc jobs submit pig --cluster=${CLUSTER} --region=${REGION} \
--jars=hello.sh \
-e='sh chmod 750 ${PWD}/hello.sh; sh ${PWD}/hello.sh'