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
- O agendador adquire um lock no objeto
FairSchedulerpara garantir a consistência dos dados. - 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
preCheckpor 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.