Custo

Nível médio

Como reduzo o custo do meu cluster Flink sem quebrar o checkpointing?

A conta não dispara por causa de um job ruim — ela sobe aos poucos por causa de algumas configurações ajustadas uma vez durante o rollout e nunca revisadas. É por aí que vale começar.

Ler em:
Intervalo de checkpoint Paralelismo do TaskManager TTL no estado com chave TaskManagers para spot
Ordenado pela frequência com que é a causa real, não pelo tamanho aparente da solução.

A conta de um cluster Flink não se move como a de um job batch — você paga por TaskManagers que ficam ativos independentemente de haver backpressure, e por checkpoints que disparam por temporizador independentemente do quanto realmente mudou. Isso torna isso menos sobre um job ruim e mais sobre configurações que ninguém revisou desde o lançamento.

1. O intervalo de checkpoint é a alavanca de custo que ninguém revisa

Cada checkpoint significa uma passagem completa pelo state backend — no RocksDB, isso é I/O de disco mais, em uma implantação na nuvem, um lote de requisições PUT para o armazenamento de objetos que guarda state.checkpoints.dir. Um intervalo configurado de forma conservadora e curta durante o rollout inicial, quando ninguém queria arriscar o progresso de um job ainda não comprovado, continua custando o mesmo indefinidamente mesmo meses depois de o job estar estável. Aumentar execution.checkpointing.interval reduz esse custo recorrente diretamente — o trade-off é mais dados para reprocessar na próxima falha, então vale a pena dimensionar pela tolerância real de tempo de recuperação, não apenas por quão curto o número pode ser com segurança.

2. Paralelismo dimensionado para o pior dia, pago todo dia

O número de slots do TaskManager geralmente é configurado uma vez, com base em uma estimativa de pico de carga, e nunca mais é tocado — então o tráfego em estado estável paga por capacidade que não usa na maior parte do tempo. O Reactive Mode ou um autoscaler no nível da plataforma podem fechar essa lacuna automaticamente; mesmo sem isso, comparar o backpressure por subtask com os slots provisionados já pega o caso óbvio: paralelismo dimensionado para um pico que acontece duas vezes por ano, pago a preço cheio todo santo dia entre um e outro.

3. O crescimento ilimitado de estado infla os checkpoints silenciosamente

Estado com chave que nunca recebeu um TTL cresce enquanto o job estiver rodando, e cada checkpoint precisa serializar tudo isso — então a duração do checkpoint e o espaço de armazenamento sobem mês a mês sem que um único deploy seja o culpado. Nada quebra, e é exatamente por isso que passa despercebido. Aplicar um TTL que corresponda a quanto estado a lógica de negócio realmente precisa costuma ser um ganho maior do que parece, justamente porque corrige uma deriva lenta, não um valor mal configurado.

4. Nem todo TaskManager é um candidato seguro para spot

Se um TaskManager pode rodar em capacidade spot depende do que ele está segurando: um job que se recupera de forma limpa a partir do último checkpoint, com tolerância de tempo de recuperação suficiente para absorver uma preempção, é um bom candidato — o mesmo orçamento de tempo de recuperação da alavanca nº 1. Um job com um estado com chave grande que leva minutos para reconstruir do zero é um candidato pior, já que uma preempção ali custa mais do que a computação que ela economizou.

O que o opti-pipe lê de uma exportação de métricas do Flink: duração do checkpoint, tamanho do checkpoint e backpressure por subtask — nunca a lógica de negócio do seu grafo de jobs — para apontar qual dessas quatro alavancas realmente vale a pena puxar na sua implantação específica, com um número concreto em vez de uma regra geral.

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 — não mais uma regra de bolso.