Flink

Nível médio

Como 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.

Ler em:
EmbeddedRocksDBStateBackend HashMapStateBackend State backend customizado
Como jobs Flink em produção com tamanho de estado relevante costumam se dividir entre backends, aproximadamente - não um levantamento medido.

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.

# flink-conf.yaml state.backend.type: rocksdb state.backend.incremental: true

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:

# flink-conf.yaml state.checkpoint-storage: filesystem state.checkpoints.dir: s3://my-flink-checkpoints/prod/ state.checkpoints.num-retained: 3

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.