Otimização de Desempenho do Agendador Hadoop YARN: Práticas e Evolução

Com o crescimento da escala do cluster e o aumento do volume de negócios, a capacidade de agendamento do YARN tornou-se um gargalo de desempenho crítico. Por exemplo, em um cluster de 1000 nós com 100 CPUs cada, demandando 100.000 CPUs em horários de pico e com tarefas de curta duração, o agendador pode conseguir processar apenas 50.000 tarefas por minuto. Isso resulta em uma subutilização significativa dos recursos do cluster (apenas 50% de utilização), que piora drasticamente com o aumento da escala do cluster. A meta é otimizar o agendador para suportar clusters com dezenas de milhares de nós e dezenas de milhares de trabalhos concorrentes.

Arquitetura e Conceitos do YARN

O YARN é responsável por gerenciar recursos do cluster, localizar recursos adequados para os trabalhos, iniciar contêineres de tarefas e gerenciar o ciclo de vida dos trabalhos.

Abstração de Recursos

O YARN abstrai os recursos do cluster em CPU e memória:


class Recurso {
   int numCPUs;       // Número de núcleos de CPU
   int memoriaMB;     // Memória em Megabytes
}
   

As requisições de recursos de um trabalho são expressas como uma lista de SolicitacaoRecurso:


class SolicitacaoRecurso {
   int numContêineres;   // Quantidade de contêineres necessários
   Recurso capacidade;    // Capacidade de recurso de cada contêiner
}
   

A resposta do YARN é uma lista de Contêiner alocados:


class Contêiner {
   IdentificadorContêiner containerId; // Identificador único global do contêiner
   Recurso capacidade;               // Informações de recurso do contêiner
   String enderecoHttpNodeManager;  // Endereço do NodeManager onde o contêiner pode ser iniciado
}
   

Arquitetura do Agendador YARN

Componentes chave do agendador incluem:

  • ResourceScheduler: A interface abstrata do agendador, responsável pela alocação de contêineres.
  • AsyncDispatcher: Um despachante de eventos single-thread para o agendador.
  • ResourceTrackerService: Gerencia os heartbeats dos NodeManagers.
  • ApplicationMasterService: O serviço RPC para ApplicationMasters, gerenciando seus heartbeats.
  • AppMaster: O controlador do trabalho, interagindo com o YARN para solicitar/liberar recursos.

O fluxo de agendamento é assíncrono. O AppMaster informa suas necessidades de recursos (List<ResourceRequest>) e recebe contêineres alocados em troca dos heartbeats.

O Agendador Justo (FairScheduler)

A Meituan utiliza o FairScheduler, que organiza os trabalhos (Apps) em uma estrutura de árvore de filas hierárquicas.

Fluxo Central de Agendamento

  1. O agendador adquire um lock no objeto FairScheduler para garantir a consistência dos dados.
  2. O agendador seleciona um nó do cluster e, começando pela raiz da árvore de filas, navega recursivamente pelas filas, aplicando a política de justiça para selecionar uma sub-fila e, finalmente, um trabalho (App) na fila folha. O objetivo é encontrar um recurso adequado no nó selecionado para o App escolhido.

Em cada nível da árvore de filas:

  • Uma verificação prévia é realizada para garantir que o uso de recursos da fila não exceda sua cota (Quota).
  • As sub-filas ou aplicativos são ordenados com base na política de justiça.
  • O processo é chamado recursivamente para as sub-filas.

class AgendadorJusto {
   // Locks e threads para garantir consistência e paralelismo
   synchronized Recurso tentarAgendar(IdentificadorNo no) {
       raiz.alocarContêiner(no);
   }
}

class Fila {
   Recurso alocarContêiner(IdentificadorNo no) {
       if (!preVerificar(no)) return null; // Verificação prévia
       ordenar(this.filhos);           // Ordenar sub-filas/apps
       if (this.ehPai) {
           for (Fila f : this.filhos)
               f.alocarContêiner(no);    // Chamada recursiva
       } else {
           for (App app : this.appsExecutaveis)
               app.alocarContêiner(no); // Alocar para o app
       }
       return alocado;
   }
}

class App {
   Recurso alocarContêiner(IdentificadorNo no) {
       // Lógica de alocação para o app
   }
}
   

A arquitetura do FairScheduler envolve múltiplos threads assíncronos, mas o processo central de alocação de contêineres é single-thread, o que representa o principal gargalo de desempenho.

Avaliação de Desempenho

A avaliação do desempenho de um sistema de agendamento difere de sistemas online tradicionais. Não é baseada na latência de RPC, pois estes são assíncronos em relação ao processo de agendamento.

Métrica de Negócio: Agendamento Válido (ValidSchedule)

O principal objetivo é atender às demandas de recursos dos negócios. ValidSchedulePerMin é uma métrica que indica se o agendamento foi bem-sucedido em um determinado minuto.


// Definições:
// validPending: Demanda efetiva de recursos do cluster.
// queuePending: Demanda total de recursos nas filas.
// QueueMaxQuota: Cota máxima de recursos da fila.
// usage: Recursos em uso no cluster.
// total: Recursos totais do cluster.

// Cálculo de validPending
validPending = min(queuePending, QueueMaxQuota);

// Condições para um agendamento válido por minuto
if (usage / total > 0.9 || validPending == 0) {
   validSchedulePerMin = 1; // Cluster com >90% de uso ou sem demanda efetiva.
} else if (validPending > 0 && usage / total < 0.9) {
   validSchedulePerMin = 0; // Cluster com <90% de uso e com demanda efetiva.
} else {
   validSchedulePerMin = 1; // Caso padrão, assumindo que o agendador está atendendo à demanda.
}

// Taxa de sucesso diária
validSchedulePerDay = ΣvalidSchedulePerMin / 1440;
   

A meta online é validSchedulePerMin > 0.9 e validSchedulePerDay > 0.99.

Métrica de Sistema: Contêineres por Segundo (CPS)

CPS mede a capacidade bruta do agendador de alocar contêineres.

CPS é influenciado por:

  • Número total de recursos do cluster.
  • Número de aplicativos em execução.
  • Número de filas.
  • Tempo de execução das tarefas (tarefas mais curtas geram mais pressão).

Exemplo: Em um cluster de 1000 nós com 1000 Apps concorrentes em 500 filas, cada contêiner executando por 1 minuto, o CPS pode ser de 1000/s.

Simulador de Carga de Agendamento (Scheduler Load Simulator - SLS)

Para testes offline eficientes sem a necessidade de um cluster físico massivo, o SLS foi adaptado. A versão modificada integra um ResourceManager real, permitindo que o SLS simule apenas as requisições de recursos e heartbeats dos nós. Isso garante que as métricas coletadas sejam comparáveis entre os ambientes de teste e produção.

Métricas de Monitoramento Granular

Para identificar gargalos específicos, métricas granulares são coletadas diretamente no fluxo de agendamento, medindo o tempo gasto em funções chave, como:

  • Tempo de preCheck por fila pai/filha.
  • Tempo de ordenação por fila pai/filha.
  • Tempo de alocação de recursos para um App.
  • Tempo gasto com Jobs sem demanda de recursos.

Pontos Chave de Otimização

A otimização é um processo iterativo, guiado por testes de carga e métricas.

Otimização da Função de Comparação de Ordenação

Identificou-se que a ordenação de filas/jobs consumia uma grande parte do tempo de agendamento (30 segundos em um ciclo de 50 segundos). A ineficiência vinha do cálculo recorrente do uso de recursos (resourceUsage) para cada comparação.

Estratégia de Otimização: Calcular o resourceUsage antecipadamente. Ao alocar ou liberar um contêiner, o resourceUsage das filas pais é atualizado em O(1). Isso reduziu o tempo de ordenação para menos de 5 segundos.

Otimização do Tempo de Pular Jobs

Após otimizar a ordenação, o tempo gasto pulando jobs sem demanda de recursos amuentou significativamente (de 2 para 20 segundos).

Estratégia de Otimização: Antes da ordenação, remover filas/jobs sem demanda de recursos da lista a ser processada. Isso reduziu o tempo para um valor desprezível.

Otimização da Ordenação Paralela de Filas

A ordenação ocorria a cada alocação de contêiner, tornando-se um gargalo à medida que o número de filas e jobs aumentava. Uma reflexão sobre a filosofia do FairScheduler revelou que a garantia de justiça em um único ponto no tempo pode não ser suficiente; a justiça ao longo do tempo requer mecanismos como preempção.

Estratégia de Otimização: Paralelizar a ordenação. Um pool de threads separado é usado para ordenar as filas, com cada thread processando uma fila. Antes da ordenação, os dados relevantes para a ordenação são clonados para evitar a modificação das estruturas de dados principais.

Resultado: A ordenação de 2000 filas pode ser realizada em menos de 5ms. Em testes de carga extrema (10.000 jobs, 1.200 nós, 2000 filas), o CPS atingiu 50.000, permitindo a utilização completa dos recursos do cluster.

Estratégias de Implantação Estável

A transição de otimizações do ambiente de teste para produção requer cautela.

Estratégia de Rollback On line

Cada otimização é controlada por parâmetros configuráveis. Modificar esses parâmetros permite alternar entre a lógica otimizada e a original sem reiniciar o serviço. Para evitar inconsistências de dados entre threads de agendamento e atualização de configuração, os parâmetros de otimização são copiados no início de cada ciclo de agendamento.

Estratégia de Validação Automática de Dados

Para garantir a correção dos algoritmos otimizados (como o cálculo de resourceUsage), comparações regulares são feitas entre os resultados do método antigo e do novo. Se as diferenças excederem um limite, um alerta é gerado e os dados incorretos são automaticamente substituídos pelos corretos, garantindo a integridade do agendamento.

Conclusão e Perspectivas Futuras

As otimizações focaram em definir métricas claras, usar ferramentas de teste eficazes, reduzir a complexidade algorítmica, eliminar cálculos redundantes e introduzir paralelismo. O trabalho contínuo visa escalar para clusters ainda maiores.

As direções futuras incluem a exploração do Global Scheduling do Hadoop 3.x para melhorar drasticamente o desempenho do agendamento em clusters únicos e a adoção e aprimoramento do YARN Federation para suportar múltiplos clusters YARN atuando como um único serviço de computação, escalando horizontalmente a capacidade de agendamento.

Tags: hadoop YARN FairScheduler otimização de desempenho agendamento

Publicado em 8-2 08:33