O EMR Remote Shuffle Service (ESS) é uma extensão do E-MapReduce (EMR) que otimiza operações de shuffle em mecanismos de computação.
Contexto
Os mecanismos tradicionais de shuffle apresentam diversos desafios:
Em cenários com grandes volumes de dados, as operações de escrita de shuffle podem causar transbordamento para o disco, resultando em amplificação de escrita.
As operações de leitura de shuffle geram alto volume de pacotes de rede pequenos, o que pode provocar erros de redefinição de conexão.
A leitura de shuffle envolve inúmeras solicitações de I/O pequenas e leituras aleatórias, sobrecarregando discos e CPUs.
Quando a quantidade de mappers (M) e reducers (N) atinge milhares, o total de conexões de rede (M × N) pode impedir a conclusão do job.
O NodeManager e o Spark Shuffle Service executam no mesmo processo. Se o volume de dados de shuffle for extremamente grande, o NodeManager poderá reiniciar e afetar a estabilidade do agendamento do YARN.
O ESS oferece as seguintes vantagens:
Utiliza mecanismo de shuffle baseado em push em vez de pull, reduzindo a pressão de memória nos mappers.
Suporta agregação de I/O, diminuindo o número de conexões de leitura de shuffle de M × N para N e substituindo leituras aleatórias por sequenciais.
Oferece suporte a mecanismo de duas réplicas para reduzir a probabilidade de falhas de busca.
Permite arquitetura de separação entre computação e armazenamento, possibilitando implantar o Shuffle Service em ambiente de hardware separado e desacoplado do cluster de computação.
Elimina a dependência de discos locais ao executar o Spark no Kubernetes.
A figura a seguir ilustra a arquitetura do ESS.
Limitações
Este documento aplica-se apenas a versões do EMR anteriores à EMR-3.39.1, versões da série EMR-4.x e versões anteriores à EMR-5.5.0. Para EMR-3.39.1 ou posterior e EMR-5.5.0 ou posterior, consulte RSS.
Crie um cluster
Por exemplo, na versão EMR-4.5.0, crie um cluster com ESS de duas maneiras:
Crie um cluster E-MapReduce Shuffle Service. Na página Software Configuration, em Cluster Type, selecione Shuffle Service. O serviço obrigatório é o ESS (1.0.0).
Crie um cluster E-MapReduce Hadoop. Na página Software Configuration, na seção Cluster Type, selecione um tipo como Hadoop, Kafka ou Druid. Em seguida, configure as Cloud Native Options (por exemplo, no ECS) e a Product Version (por exemplo, EMR-4.5.0). A página exibirá os serviços obrigatórios correspondentes (como HDFS, YARN e Spark) e os serviços opcionais (como ESS, HBase e Flink) com suas respectivas versões.
Para obter mais informações sobre como criar um cluster, consulte Criar um cluster.
Uso do ESS
Para usar o ESS com o Spark, adicione os seguintes parâmetros ao envio do job Spark. Para mais detalhes sobre a configuração de parâmetros, consulte Edite jobs.
Para mais informações sobre os parâmetros do Spark, consulte Configuração do Spark.
|
Parâmetro |
Descrição |
|
spark.shuffle.manager |
O valor deve ser org.apache.spark.shuffle.ess.EssShuffleManager. |
|
spark.ess.master.address |
Especifique o endereço no formato <ess-master-ip>:<ess-master-port>. Os parâmetros são:
|
|
spark.shuffle.service.enabled |
Defina o valor como Desative o serviço externo de shuffle padrão para usar o EMR Remote Shuffle Service. |
|
spark.shuffle.useOldFetchProtocol |
Defina o valor como Habilita compatibilidade com o protocolo legado de shuffle. |
|
spark.sql.adaptive.enabled |
Defina o valor como O EMR Remote Shuffle Service não oferece suporte à Adaptive Execution. |
|
spark.sql.adaptive.skewJoin.enabled |
Parâmetros
A página de configuração do serviço ESS lista todos os parâmetros disponíveis.
|
Parâmetro |
Descrição |
Padrão |
|
ess.push.data.replicate |
Ativa ou desativa o recurso de duas réplicas. Valores válidos:
Nota
Recomendamos ativar este recurso em ambientes de produção. |
true |
|
ess.worker.flush.queue.capacity |
Quantidade de buffers de flush por diretório. Nota
Para melhorar o desempenho, configure vários discos. Para obter throughput ideal de leitura e escrita, use no máximo dois diretórios por disco. A memória heap consumida pelo buffer de flush de cada diretório corresponde a ess.worker.flush.buffer.size ess.worker.flush.queue.capacity, ou seja, |
512 |
|
ess.flush.timeout |
Tempo limite para flush de dados na camada de armazenamento. |
240s |
|
ess.application.timeout |
Tempo limite de heartbeat da aplicação. Se nenhum heartbeat for recebido nesse período, o ESS limpará os recursos da aplicação. |
240s |
|
ess.worker.flush.buffer.size |
Tamanho do buffer de flush. Quando o buffer ultrapassa esse tamanho, o ESS faz o flush dos dados para o disco. |
256k |
|
ess.metrics.system.enable |
Ativa ou desativa o monitoramento. Valores válidos:
|
false |
|
ess_worker_offheap_memory |
Tamanho da memória off-heap para um nó principal. |
4g |
|
ess_worker_memory |
Tamanho da memória heap para um nó principal. |
4g |
|
ess_master_memory |
Tamanho da memória heap para o nó mestre. |
4g |