Flink
Nível médioComo configuro o state backend e o armazenamento de checkpoints do Flink antes de ir para produção?
Os padrões do Flink são ajustados para colocar um job rodando rápido no desenvolvimento, não para sobreviver a um restart real com um tamanho de estado real. Duas configurações decidem quase tudo.
Todo job com estado no Flink precisa de um state backend (onde o estado do operador vive enquanto o job roda) e armazenamento de checkpoints (onde snapshots consistentes desse estado são escritos). O Flink roda bem nos padrões até o estado crescer além do que um heap de JVM aguenta confortavelmente.
HashMapStateBackend vs EmbeddedRocksDBStateBackend
HashMapStateBackend mantém todo o estado como objetos Java no heap - rápido, mas limitado
por quanto heap você deu ao task manager, e todo checkpoint precisa serializar o estado inteiro que está
sendo escrito. EmbeddedRocksDBStateBackend mantém o estado em disco local (com um cache em
memória configurável) e só serializa o que é realmente necessário por checkpoint, ao custo de alguma
latência por acesso ao estado. Passado algumas centenas de MB de estado por slot de tarefa, RocksDB
costuma ser o padrão certo, não uma otimização para mais tarde.
state.backend.incremental: true importa quase tanto quanto a escolha do backend em si -
sem ele, todo checkpoint reescreve o snapshot completo do estado em vez de só o que mudou desde o último,
o que fica caro rápido conforme o estado cresce.
O armazenamento de checkpoints precisa ser durável e compartilhado, não local
O armazenamento de checkpoints é uma configuração separada do state backend - é onde os arquivos de checkpoint realmente caem, e precisa ser alcançável por qualquer task manager que recupere uma tarefa que falhou, o que na prática significa object storage, não disco local:
num-retained mantém mais de um checkpoint de propósito - se o checkpoint mais recente
acabar corrompido ou no meio da escrita quando um task manager morre, o Flink recorre ao anterior em vez
de não ter nada para recuperar.
A configuração dolorosa de mudar depois
Trocar de state backend depois que um job já tem estado real de produção significa que o novo backend não consegue ler o formato de checkpoint do backend antigo - não existe migração in-place, só um savepoint tirado sob o backend antigo e restaurado sob o novo, com downtime real enquanto isso acontece. Escolher RocksDB antes do lançamento se há qualquer chance do estado crescer não custa nada; migrar para fora do HashMapStateBackend sob carga é uma janela de manutenção planejada.
Fora do escopo do opti-pipe: state backend e armazenamento de checkpoints são configuração de cluster/job, definida antes do job começar - o motor de regras só lê as métricas que um job rodando ou concluído exporta, não mexe em como ou onde o estado é persistido.
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.