Saltar a contenido

ADR-0005 — Procesamiento de streams: Faust

Contexto

Con Kafka como capa de streaming, hace falta un marco de procesamiento que consuma de los temas y aplique lógica en tiempo real: limpieza de lecturas, agregación en ventanas móviles, puntuación de anomalías cargando el modelo entrenado y reacciones a patrones sobre el flujo.

Los requisitos son claros: throughput suficiente con margen para diez veces la carga; procesamiento con estado (ventanas móviles, tablas, cruces entre flujos); semántica de entrega razonable (al menos una vez, compensada con escrituras idempotentes en el destino); y, de forma explícita, mantenerse dentro del ecosistema Python sin introducir una pila JVM que duplicaría la demanda de habilidades y la carga operacional.

Decisión

Se adopta Faust como marco de procesamiento de streams. Faust es el equivalente en Python a las bibliotecas de streams sobre Kafka: combina consumidores, productores, tablas con estado y ventanas temporales en una API idiomática basada en decoradores, muy en la línea del resto del backend. El estado con tolerancia a fallos se mantiene localmente y se replica a través de Kafka. Los trabajadores escalan horizontalmente añadiendo procesos, con reequilibrado automático. Los modelos de ML se cargan al arrancar y se recargan periódicamente. Las escrituras al destino son idempotentes, de modo que la entrega "al menos una vez" produce en la práctica un efecto de "exactamente una vez".

Se utiliza la variante mantenida por la comunidad del proyecto.

Consecuencias

Positivas:

  • Pila Python de extremo a extremo, sin JVM ni operación políglota.
  • Los patrones con estado (tablas, ventanas, cruces) vienen resueltos, en lugar de reimplementarse a mano.
  • Huella pequeña: cada trabajador es un proceso Python ligero.
  • La inferencia de ML se integra de forma natural dentro del procesamiento del flujo.
  • Modelo mental coherente con el resto del backend (decoradores y asincronía).

Negativas o costes aceptados:

  • Throughput por trabajador inferior al de los motores sobre JVM; se compensa con paralelismo (varios trabajadores por grupo, alineados con el número de particiones).
  • La semántica estricta de "exactamente una vez" mediante transacciones es posible pero se prefiere la vía más simple de idempotencia en el destino.
  • Ecosistema de conectores menor que el de los motores más grandes.
  • Se depende de la variante comunitaria del proyecto, una dependencia asumida de forma consciente.

Alternativas descartadas

  • Consumidores Python simples: máxima simplicidad, pero obligarían a reimplementar a mano tablas, ventanas y cruces, es decir, a reconstruir Faust peor. Se reservan solo para casos puntuales sin estado.
  • Apache Spark Structured Streaming: estándar en Big Data, pero pesado en JVM, con latencia por micro-lotes y una curva de aprendizaje enorme; sobredimensionado para el volumen objetivo.
  • Apache Flink: el mejor motor de streaming (verdadero tiempo real, "exactamente una vez" estricto), pero también pesado en JVM y con la API de Python menos madura; sobredimensionado aquí.
  • Bibliotecas de streams sobre JVM (Java/Kotlin): máxima eficiencia con Kafka, pero chocan de frente con una pila 100 % Python.
  • SQL sobre streams: muy declarativo, pero la lógica compleja (inferencia de un modelo de ML) no se expresa bien en SQL y exige otro servicio.