Spark
Nível médioComo configuro dynamic allocation do Spark do zero, não só ajusto depois?
Dynamic allocation costuma ser discutido como um botão de ajuste, mas tem pré-requisitos rígidos - sem eles, silenciosamente não faz nada, o que parece idêntico a um job mal ajustado.
spark.dynamicAllocation.enabled=true sozinho não faz nada na maioria dos cluster managers. Dynamic allocation primeiro precisa de um shuffle service externo configurado e rodando - sem ele, executors não podem ser removidos com segurança no meio do job, então o Spark simplesmente... não os remove.
O pré-requisito que ninguém menciona primeiro
Remover um executor no meio do job só é seguro se os dados de shuffle que ele guarda continuarem
disponíveis para outros executors depois que ele se for - é isso que o shuffle service externo faz,
rodando independentemente do ciclo de vida de qualquer executor específico. Sem ele habilitado, o Spark
aceita dynamicAllocation.enabled=true sem erro e depois simplesmente nunca reduz a escala,
porque fazer isso arriscaria perder dados de shuffle. No YARN e no Kubernetes isso é um serviço separado
que precisa ser implantado e apontado a partir da configuração do Spark - não aparece automaticamente só
porque dynamic allocation está ligado.
Definindo limites reais, não só um teto
Deixar minExecutors no padrão de 0 significa que um job pode reduzir totalmente a zero
executors durante um momento de calmaria e depois pagar a latência completa de alocação do cluster manager
para escalar de volta - tudo bem para um job em lote com folga, custoso para qualquer coisa sensível a
latência. maxExecutors é o botão mais conhecido (é o teto de custo), mas é só metade do
quadro sem um piso correspondente:
Timeout de ociosidade, e onde ainda é preciso bom senso
spark.dynamicAllocation.executorIdleTimeout (padrão 60s) controla quanto tempo um
executor fica ocioso antes de ser liberado - curto demais e um job com estágios irregulares e em rajadas
fica sacudindo executors para cima e para baixo o tempo todo, adicionando latência de alocação a cada
rajada; longo demais e você está pagando por capacidade ociosa entre rajadas. Não existe fórmula para o
valor certo - depende de quão irregular a carga de trabalho realmente é, que é exatamente o tipo de padrão
que vale a pena olhar na linha do tempo de um run real antes de escolher um número.
Onde o opti-pipe entra: o motor de regras lê o event log de um job concluído e sinaliza
indícios de sub ou superprovisionamento contra os limites que estão realmente configurados - ele não
define minExecutors/maxExecutors por você, já que isso é uma troca entre custo
e latência que só você pode fazer.
Veja como isso fica no seu próprio pipeline.
Envie um event log real do Spark, um run_results.json do dbt ou uma exportação de métricas do Flink e receba recomendações concretas, para aprovar antes de aplicar - não mais uma regra geral.