Este tópico responde às perguntas mais comuns sobre o uso do Spark.
-
Spark Core
-
Spark SQL
-
PySpark
-
Spark Streaming
-
spark-submit
-
No nó em que você executa as ferramentas de CLI (por exemplo, o nó master), crie um arquivo de configuração log4j.properties. Para copiar o arquivo de configuração padrão, execute o seguinte comando:
cp /etc/emr/spark-conf/log4j.properties /new/path/to/log4j.properties -
Modifique o nível de log no novo arquivo de configuração.
log4j.rootCategory=WARN, console -
No arquivo spark-defaults.conf do serviço Spark, atualize a propriedade spark.driver.extraJavaOptions. Substitua -Dlog4j.configuration=file:/etc/emr/spark-conf/log4j.properties por -Dlog4j.configuration=file:/new/path/to/log4j.properties.
ImportanteO caminho deve ter o prefixo file:.
-
Para o Spark 2, utilize uma das seguintes abordagens:
Filtre dados irrelevantes, como valores nulos, durante a leitura da tabela.
-
Faça broadcast da tabela menor.
select /*+ BROADCAST (table1) */ * from table1 join table2 on table1.id = table2.id -
Separe os dados desbalanceados com base na chave de skew.
select * from table1_1 join table2 on table1_1.id = table2.id union all select /*+ BROADCAST (table1_2) */ * from table1_2 join table2 on table1_2.id = table2.id -
Se a chave de skew for conhecida, disperse os dados.
select id, value, concat(id, (rand() * 10000) % 3) as new_id from A select id, value, concat(id, suffix) as new_id from ( select id, value, suffix from B Lateral View explode(array(0, 1, 2)) tmp as suffix) -
Caso a chave de skew seja desconhecida, disperse os dados.
select t1.id, t1.id_rand, t2.name from ( select id , case when id = null then concat('SkewData_', cast(rand() as string)) else id end as id_rand from test1 where statis_date = '20221130') t1 left join test2 t2 on t1.id_rand = t2.id
Para o Spark 3, acesse a aba Configure do serviço Spark 3 no console do EMR e defina os parâmetros spark.sql.adaptive.enabled e spark.sql.adaptive.skewJoin.enabled como true.
Conecte-se ao cluster via SSH. Para mais informações, consulte Fazer login em um cluster.
-
Execute o comando abaixo para alterar a versão do Python:
export PYSPARK_PYTHON=/usr/bin/python3 -
Verifique a versão do Python com o seguinte comando:
pysparkSe a saída contiver a mensagem a seguir, a versão foi alterada para Python 3.
Using Python version 3.6.8 Conecte-se ao cluster via SSH. Para mais informações, consulte Fazer login em um cluster.
-
Modifique o arquivo de configuração.
-
Abra o arquivo profile com o seguinte comando:
vi /etc/profile Pressione
ipara entrar no modo de inserção.-
Ao final do arquivo profile, adicione a seguinte linha:
export PYSPARK_PYTHON=/usr/bin/python3 Pressione
Escpara sair do modo de inserção. Em seguida, insira:wqpara salvar e fechar o arquivo.
-
-
Recarregue o arquivo de configuração para aplicar as alterações imediatamente:
source /etc/profile -
Verifique a versão do Python com o seguinte comando:
pysparkSe a saída contiver a mensagem a seguir, a versão foi alterada para Python 3.
Using Python version 3.6.8 -
Caso sua versão do Spark seja anterior à 1,6, faça upgrade.
Versões do Spark anteriores à 1,6 têm um vazamento de memória capaz de causar o encerramento dos containers.
Certifique-se de que seu código esteja otimizado quanto ao uso de memória.
Onde visualizar jobs históricos do Spark?
No console do EMR, acesse a aba Access Links and Ports do cluster desejado e clique em no link da Spark UI para visualizar os jobs históricos. Para obter mais informações sobre como acessar as interfaces dos componentes, consulte Acessar as web UIs de componentes open-source no console do EMR.
É possível enviar jobs do Spark no modo standalone?
Não. O E-MapReduce suporta o envio de jobs apenas por meio do Spark on YARN e Spark on Kubernetes. Os modos standalone e Mesos não são suportados.
Como reduzir a saída de log da CLI do Spark 2
Por padrão, ferramentas de interface de linha de comando (CLI) como spark-sql e spark-shell em um cluster DataLake do EMR geram logs no nível INFO.
Como usar a mesclagem de arquivos pequenos do Spark 3
Defina o parâmetro spark.sql.adaptive.merge.output.small.files.enabled como true para mesclar arquivos pequenos automaticamente. Os arquivos mesclados serão compactados. Caso os arquivos resultantes fiquem muito pequenos, aumente o valor do parâmetro spark.sql.adaptive.advisoryOutputFileSizeInBytes. O valor padrão é 256 MB.
Como lidar com data skew no Spark SQL
Como especificar o Python 3 para o PySpark
Esta seção demonstra como especificar o Python 3 para o PySpark, utilizando como exemplo um cluster DataLake no EMR V5.7.0 com Spark 2.
Existem dois métodos para alterar a versão do Python:
Método temporário
Método permanente
Por que jobs do Spark Streaming param inesperadamente
Por que um job parado aparece como em execução
Isso pode ocorrer se você enviou o job no modo YARN-client, pois o E-MapReduce não monitora com precisão o status dos jobs nesse modo. Para garantir o relatório correto de status, envie o job no modo YARN-cluster.
Erro java.lang.ClassNotFoundException no spark-submit em modo YARN-cluster
Exemplo de mensagem de erro:
Process Output>>> 24/09/30 15:41:24 WARN HiveConf: HiveConf of name hive.metastore.type does not exist
Process Output>>> 24/09/30 15:41:24 ERROR Hive: Unable to instantiate a metastore client factory com.aliyun.datalake.metastore.hive2.DlfMetaStoreClientFactory: java.lang.ClassNotFoundException: Class com.aliyun.datalake.metastore.hive2.DlfMetaStoreClientFactory not found
Process Output>>> java.lang.ClassNotFoundException: Class com.aliyun.datalake.metastore.hive2.DlfMetaStoreClientFactory not found
Process Output>>> at org.apache.hadoop.conf.Configuration.getClassByName(Configuration.java:2542)
Process Output>>> at org.apache.hadoop.hive.ql.metadata.Hive.createMetaStoreClient(Hive.java:3711)
Process Output>>> at org.apache.hadoop.hive.ql.metadata.Hive.getMSC(Hive.java:3794)
Process Output>>> at org.apache.hadoop.hive.ql.metadata.Hive.getMSC(Hive.java:3774)
Process Output>>> at org.apache.hadoop.hive.ql.metadata.Hive.getDelegationToken(Hive.java:3924)
Process Output>>> at org.apache.spark.sql.hive.security.HiveDelegationTokenProvider.$anonfun$obtainDelegationTokens$4(HiveDelegationTokenProvider.scala:104)
Process Output>>> at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
Process Output>>> at org.apache.spark.sql.hive.security.HiveDelegationTokenProvider$$anon$1.run(HiveDelegationTokenProvider.scala:139)
Process Output>>> at java.security.AccessController.doPrivileged(Native Method)
Process Output>>> at javax.security.auth.Subject.doAs(Subject.java:422)
Process Output>>> at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1730)
Process Output>>> at org.apache.spark.sql.hive.security.HiveDelegationTokenProvider.doAsRealUser(HiveDelegationTokenProvider.scala:138)
Process Output>>> at org.apache.spark.sql.hive.security.HiveDelegationTokenProvider.obtainDelegationTokens(HiveDelegationTokenProvider.scala:102)
Process Output>>> at org.apache.spark.deploy.security.HadoopDelegationTokenManager.$anonfun$obtainDelegationTokens$2(HadoopDelegationTokenManager.scala:164)
Process Output>>> at scala.collection.TraversableLike.$anonfun$flatMap$1(TraversableLike.scala:293)
Process Output>>> at scala.collection.Iterator.foreach(Iterator.scala:943)
Process Output>>> at scala.collection.Iterator.foreach$(Iterator.scala:943)
Process Output>>> at scala.collection.AbstractIterator.foreach(Iterator.scala:1431)
Process Output>>> at scala.collection.MapLike$DefaultValuesIterable.foreach(MapLike.scala:214)
Process Output>>> at scala.collection.TraversableLike.flatMap(TraversableLike.scala:293)
Process Output>>> at scala.collection.TraversableLike.flatMap$(TraversableLike.scala:290)
Causa: Em um cluster EMR com Kerberos ativado, o classpath do driver não é preenchido automaticamente com os JARs necessários quando um job é executado no modo YARN-cluster. Isso causa o erro ClassNotFoundException.
Solução: Em um cluster EMR com Kerberos ativado, use obrigatoriamente o parâmetro --jars ao enviar um job com spark-submit no modo YARN-cluster. Além dos JARs da sua aplicação, inclua todos os pacotes JAR do diretório /opt/apps/METASTORE/metastore-current/hive2.
No modo YARN-cluster, todos os caminhos de arquivo no parâmetro --jars devem ser separados por vírgula. Diretórios não são suportados.
Por exemplo, se o JAR da sua aplicação for /opt/apps/SPARK3/spark3-current/examples/jars/spark-examples_2.12-3.5.3-emr.jar, execute o seguinte comando spark-submit:
spark-submit --deploy-mode cluster --class org.apache.spark.examples.SparkPi --master yarn \
--jars $(ls /opt/apps/METASTORE/metastore-current/hive2/*.jar | tr '\n' ',') \
/opt/apps/SPARK3/spark3-current/examples/jars/spark-examples_2.12-3.5.3-emr.jar