Saltar a contenido

ADR-0003 — Capa de streaming: Apache Kafka

Contexto

Un flujo puramente batch, en el que la ingesta escribe directamente al almacén y un trabajo horario limpia los datos, arrastra una latencia de hasta una hora entre el dato crudo y el dato limpio. Introducir una capa de streaming baja esa latencia a un par de segundos y, sobre todo, habilita el fan-out: que varios consumidores independientes (limpieza, agregación, detección de anomalías) lean el mismo flujo sin competir entre sí y sin reescribir a los productores.

Los requisitos de esta capa son exigentes: throughput sostenido con margen para diez veces la carga, orden garantizado dentro de cada sensor (la detección de anomalías y la predicción dependen de secuencias temporales estrictas), varios grupos de consumidores con sus propios avances de lectura, retención suficiente para poder reprocesar, una cola de mensajes muertos para lo que llega corrupto y una operativa manejable tanto en desarrollo como en producción.

Decisión

Se adopta Apache Kafka como capa de streaming, desplegado en modo KRaft (sin el componente externo de coordinación clásico) con varios brokers replicados. Los mensajes se particionan por identificador de sensor, lo que garantiza el orden por sensor. Se definen temas separados para las lecturas crudas, las limpias, las anomalías, las predicciones y los mensajes muertos, cada uno con su política de retención. Cada consumidor (limpieza, anomalías, agregación) forma su propio grupo con avances de lectura independientes.

Se autohospeda por completo, sin recurrir a Kafka gestionado.

Por qué encaja

El patrón que se necesita (un log durable, particionado, con varios consumidores independientes y capacidad de reproducir) es exactamente el diseño nativo de Kafka. El resto de alternativas son aproximaciones con compromisos.

Consecuencias

Positivas:

  • Fan-out sin tocar a los productores: cada consumidor se suscribe al mismo tema con su propio avance de lectura.
  • Reprocesamiento nativo: un consumidor puede rebobinar su avance y volver a procesar una ventana de datos sin lógica adicional.
  • Orden por sensor garantizado gracias al particionado por identificador.
  • La contrapresión se absorbe en disco gracias a la retención generosa, de modo que los productores casi nunca se bloquean.
  • Un ecosistema maduro habilita evoluciones futuras (registro de esquemas, conectores, almacenamiento por niveles), todas aditivas.

Negativas o costes aceptados:

  • Varios servicios adicionales y una huella de memoria apreciable.
  • Curva de aprendizaje operacional real: ajuste, monitorización del retraso de consumidores y diagnóstico de réplicas.
  • Latencia de extremo a extremo mayor que la de una solución en memoria, irrelevante para el objetivo de un par de segundos.
  • El coste de almacenamiento se multiplica por el factor de replicación.

Alternativas descartadas

  • Redis Streams: latencia mínima y ya presente en la pila, pero de un solo hilo por instancia, con retención ligada a la memoria (inviable para conservar días de datos) y sin el ecosistema de ingeniería de datos que el proyecto busca demostrar. Entra en zona incómoda justo en el rango de escala objetivo.
  • RabbitMQ: excelente broker de colas, pero no es un log de streaming: no hay retención para reproducir (mensaje consumido, mensaje perdido) ni particionado nativo para orden por clave.
  • NATS JetStream: moderno y simple de operar, pero con un ecosistema mucho menor y sin integración natural con el resto de la pila de procesamiento.
  • Servicios gestionados en la nube: dependencia de proveedor y coste recurrente, y ocultan precisamente la operativa de clúster que es parte del aprendizaje.