Saltar a contenido

Procesamiento en streaming

El corazon en tiempo real de la plataforma. Toma lo que tradicionalmente se haria por lotes cada hora (limpiar lecturas, calcular agregados, puntuar anomalias) y lo lleva a near-realtime sobre el flujo de eventos, procesando cada lectura segundos despues de que entre. Esta construido con Faust, una libreria que trae el modelo de Kafka Streams al ecosistema Python, lo que permite mantener estado distribuido y ventanas temporales sin salir del lenguaje del resto del stack.

Su rol es consumir el topic de lecturas crudas y, sobre ese mismo flujo, hacer tres cosas independientes a la vez. Es un unico proceso que registra tres agentes desacoplados, cada uno con su propio grupo de consumo, de forma que el fallo de uno no arrastra a los demas.

Vision de conjunto de los tres agentes sobre los topics de eventos

El primer agente limpia. Aplica en cadena una bateria de detectores de calidad (desfase de reloj, valores fuera de rango, duplicados y outliers estadisticos) reutilizando exactamente los mismos detectores que usaba el proceso batch, para que la limpieza online y la offline coincidan. Lo que pasa el filtro se emite limpio a un nuevo topic y se persiste; lo que esta corrupto de origen se desvia a una cola de mensajes muertos para poder inspeccionarlo despues, en lugar de perderlo en silencio.

El segundo agente agrega. Mantiene ventanas deslizantes de un minuto por sensor y, al cerrarse cada ventana, calcula estadisticos (media, desviacion, extremos) y los guarda. Una decision de diseno importante es que las ventanas se cierran por tiempo del evento y no por el reloj de la maquina: se usa la marca temporal que trae cada lectura, con un margen de gracia para admitir llegadas tardias. Asi los agregados son correctos aunque los datos lleguen con retraso o desordenados.

El tercer agente puntua anomalias. Agrupa las lecturas en micro-lotes, reconstruye las mismas caracteristicas que se usaron para entrenar y las evalua con un modelo de deteccion de anomalias (IsolationForest) que carga desde el registro de modelos. Si no hay ningun modelo disponible todavia, no se cae: cuenta el hueco en una metrica y sigue, de modo que el resto de la plataforma funciona mientras la capa de ML madura.

El estado (ventanas deslizantes, caches de duplicados) es persistente y se replica junto con las particiones, asi que el procesamiento puede escalar horizontalmente y recuperarse tras un reinicio sin perder el contexto acumulado.