Flink
Nivel medio¿Cómo configuro el state backend y el almacenamiento de checkpoints de Flink antes de ir a producción?
Los valores por defecto de Flink están afinados para que un job arranque rápido en desarrollo, no para sobrevivir a un reinicio real con un tamaño de estado real. Dos configuraciones deciden casi todo.
Todo job con estado de Flink necesita un state backend (dónde vive el estado del operador mientras el job corre) y almacenamiento de checkpoints (dónde se escriben snapshots consistentes de ese estado). Flink corre bien con los valores por defecto justo hasta que el estado crece más allá de lo que un heap de JVM soporta cómodamente.
HashMapStateBackend vs EmbeddedRocksDBStateBackend
HashMapStateBackend mantiene todo el estado como objetos Java en el heap - rápido, pero
acotado por cuánto heap le haya dado al task manager, y cada checkpoint tiene que serializar todo el estado
que se está escribiendo. EmbeddedRocksDBStateBackend mantiene el estado en disco local (con una
caché en memoria configurable) y solo serializa lo que realmente se necesita por checkpoint, al costo de
algo de latencia por acceso al estado. Pasados unos cientos de MB de estado por slot de tarea, RocksDB
suele ser el valor por defecto correcto, no una optimización para más adelante.
state.backend.incremental: true importa casi tanto como la elección del backend en sí -
sin él, cada checkpoint reescribe el snapshot completo del estado en vez de solo lo que cambió desde el
último, lo cual se vuelve costoso rápidamente a medida que el estado crece.
El almacenamiento de checkpoints tiene que ser durable y compartido, no local
El almacenamiento de checkpoints es una configuración separada del state backend - es donde realmente caen los archivos de checkpoint, y tiene que ser alcanzable por cualquier task manager que recupere una tarea fallida, lo que en la práctica significa object storage, no disco local:
num-retained mantiene más de un checkpoint a propósito - si el checkpoint más reciente
resulta estar corrupto o a medio escribir cuando un task manager muere, Flink retrocede al anterior en vez
de no tener nada de dónde recuperarse.
La configuración dolorosa de cambiar después
Cambiar de state backend después de que un job tiene estado real de producción significa que el nuevo backend no puede leer el formato de checkpoint del backend anterior - no hay migración in situ, solo un savepoint tomado bajo el backend anterior y restaurado bajo el nuevo, con downtime real mientras eso sucede. Elegir RocksDB antes del lanzamiento si hay alguna posibilidad de que el estado crezca no cuesta nada; migrar fuera de HashMapStateBackend bajo carga es una ventana de mantenimiento planificada.
Fuera del alcance de opti-pipe: el state backend y el almacenamiento de checkpoints son configuración de cluster/job, definida antes de que un job arranque - el motor de reglas solo lee las métricas que exporta un job corriendo o completado, no toca cómo o dónde se persiste el estado.
Vea cómo se ve esto en su propio pipeline.
Suba un event log real de Spark, un run_results.json de dbt o una exportación de métricas de Flink y obtenga recomendaciones concretas, que debe aprobar antes de aplicar - no otra regla general.