Spark

Junior

O problema dos arquivos pequenos: por que seu job Spark gasta mais tempo em I/O do que em computação

Um job pode ser limitado por I/O e parecer limitado por computação em qualquer dashboard que só acompanhe CPU, porque o gargalo está em abrir arquivos, não em ler bytes.

Ler em:
Overhead de abertura Overhead de agendamento Tempo de leitura real
Para onde o tempo vai num stage afetado por arquivos pequenos, aproximadamente - ilustrativo, não medido.

Abaixo de um certo tamanho de arquivo, o custo fixo de abri-lo - busca de metadados, agendamento da tarefa, configuração do codec - custa mais do que ler seu conteúdo real.

Como você acaba com milhares de arquivos minúsculos

As fontes usuais: particionar uma escrita por uma coluna de alta cardinalidade (cada valor de partição ganha seu próprio arquivo pequeno), jobs de streaming confirmando micro-batches pequenos com um intervalo de trigger curto, ou um stage de shuffle com muito mais partições de saída do que os dados realmente precisam. Nenhuma dessas é exatamente um erro - geralmente é a decisão certa para o job que está escrevendo os dados - mas deixa uma bagunça para quem for ler depois.

Por que custa mais do que a contagem de bytes sugere

Cada arquivo, por menor que seja, ainda precisa de uma tarefa completa: agendá-la, abrir o arquivo (uma viagem de rede real no S3 ou GCS), ler um footer ou header, depois fechar. Com 10.000 arquivos minúsculos, esse overhead por arquivo pode dominar o tempo total do stage mesmo que os dados reais mal preencham a memória de alguns executors - o job parece limitado por computação num gráfico de CPU porque a CPU está ocupada agendando e cuidando de I/O, não fazendo trabalho real.

# tamanho médio aproximado de arquivo para um caminho de saída avg_file_size = total_bytes_written / file_count # bem abaixo de ~128MB (Parquet) vale a pena investigar

O ajuste que realmente resolve

coalesce() ou repartition() logo antes da escrita, mirando algo entre 128MB-1GB por arquivo de saída dependendo do formato e do motor que vai ler depois. Um ajuste único do lado da escrita é mais barato do que cada job seguinte pagando repetidamente o imposto de abrir arquivos - compactar depois também funciona, mas é um custo recorrente que uma escrita corretamente dimensionada evita por completo.

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 que exigem sua aprovação - não mais uma regra geral.