Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Lecció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.
- Duración
- 17 min aprox.
- Objetivo
- Las métricas del state store revelan si el compromiso entre tardanza, cardinalidad y capacidad sigue siendo sostenible.
- Siguiente paso
- Continuar con el laboratorio
Ver detalles del módulo
Estado, ventanas, watermarks y datos tardíos
Controla el crecimiento de estado y la corrección temporal con eventos fuera de orden.
- Definir event time y processing time
- Aplicar watermark con intención
- Deduplicar y agregar ventanas acotando estado
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.
Modelo mental
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.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.
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.
Recuerdo activo