Puedes leer sin crear un espacio. Créalo solo cuando quieras guardar.
Guardar progresoLección 5 de 5
State store, métricas y presión de memoria
Las métricas del state store revelan si el compromiso entre tardanza, cardinalidad y capacidad sigue siendo sostenible.
17 min aprox.
05Decisión de diseñoState store, métricas y presión de memoria
Las métricas del state store revelan si el compromiso entre tardanza, cardinalidad y capacidad sigue siendo sostenible.
+
State store, métricas y presión de memoria
Las métricas del state store revelan si el compromiso entre tardanza, cardinalidad y capacidad sigue siendo sostenible.
`stateOperators` informa filas totales y actualizadas, memoria usada y filas retiradas por watermark. Un crecimiento continuo después de varios horizontes puede indicar claves de cardinalidad extrema, timestamps futuros, ausencia de watermark efectivo o una condición de join no acotada.
El diagnóstico debe reproducir datos puntuales: a tiempo, duplicados, tardíos dentro de tolerancia, demasiado tardíos y timestamps inválidos. Una prueba de recuperación reinicia desde el mismo checkpoint; cambiar access mode o esquema de estado durante el experimento puede introducir otra variable.
progress = query.lastProgress or {}
for operator in progress.get("stateOperators", []):
print({
"rows_total": operator.get("numRowsTotal"),
"rows_updated": operator.get("numRowsUpdated"),
"rows_removed": operator.get("numRowsRemoved"),
"memory_bytes": operator.get("memoryUsedBytes"),
})Guarda la serie temporal; una única observación no permite distinguir una ventana grande de una fuga lógica de estado.
¿Qué patrón indica que un watermark no está retirando estado como se esperaba?
Profundiza
El state store es una base de datos incremental local y checkpointada que materializa la memoria de agregaciones, deduplicaciones y joins. Su tamaño no depende solo de filas por segundo: depende del número de claves vivas, ventanas simultáneas, versiones por clave, retraso aceptado y velocidad con que el watermark permite evacuar. Por eso un stream de bajo volumen pero millones de dispositivos únicos puede ser más difícil que uno denso con pocas claves. Las métricas deben leerse como un balance: filas actualizadas entran, filas eliminadas salen y el total retenido refleja la deuda de estado. Memoria, latencia de commit y tamaño de checkpoint crecen de manera distinta según el proveedor. RocksDB con changelog checkpointing puede reducir presión y duración de checkpoints en cargas stateful, pero cambiar el mecanismo de estado de una consulta ya iniciada puede exigir un checkpoint nuevo y reconstrucción.
Almacén versionado por operador que mantiene claves y valores necesarios entre microbatches y participa en la recuperación.
Es el principal determinante de estabilidad para operaciones stateful y debe dimensionarse como parte del diseño, no al final.Número de claves y ventanas que todavía no pueden eliminarse según la política temporal del operador.
Predice mejor el tamaño del estado que el throughput bruto de entrada.Estrategia que persiste cambios del estado entre snapshots periódicos en lugar de cargar un snapshot completo en cada checkpoint.
Puede reducir latencia de checkpoint para grandes estados RocksDB, pero requiere una decisión de arquitectura y recuperación compatible.Resumen
Puntos clave
- El número de filas de estado debe estabilizarse para una carga estacionaria y un horizonte finito.
- Los timestamps futuros pueden empujar el watermark y provocar pérdidas silenciosas de eventos normales.
- La política global `max` favorece latencia y puede descartar datos del stream más lento; `min` prioriza seguridad.
Evita
- Medir solo uso de CPU y omitir filas/memoria del estado.
- Cambiar a la política `max` para ocultar un stream lento sin aceptar explícitamente los descartes resultantes.