Ingestion API¶
La puerta de entrada de la plataforma. Es un servicio HTTP construido con FastAPI que recibe lecturas de sensores por lotes, las valida y las hace llegar al resto del sistema. Es tambien la capa que escala automaticamente cuando sube la carga, replicandose sola en el cluster segun el uso de CPU.
Su rol en el flujo es sencillo de enunciar y sutil de implementar: cada cliente autenticado envia un batch de lecturas y el servicio garantiza que esas lecturas queden a la vez en el bus de eventos (Kafka), donde las consumira el procesamiento en streaming, y en la base de datos de series temporales (TimescaleDB), donde queda el crudo consultable. Ademas del ingreso, expone la gestion de sensores, la consulta de anomalias detectadas y varios endpoints de administracion del pipeline.
La decision de diseno mas interesante es el orden del doble escritura. Kafka se escribe PRIMERO, de forma deliberada, y solo si Kafka confirma la recepcion se toca la base de datos. Kafka es la fuente de verdad del flujo aguas abajo; la tabla de crudos es reconstruible desde el propio topic. Si Kafka falla, la base de datos no se escribe y el cliente recibe una respuesta que le pide reintentar el batch entero. Asi nunca hay una lectura en la base de datos que no exista tambien en el flujo de eventos, y el reintento es seguro.
La idempotencia es la otra pieza clave. La insercion en base de datos ignora los duplicados por clave (tiempo mas sensor), de modo que reenviar un batch ya procesado no crea filas repetidas. Esto convierte los reintentos en una operacion inofensiva y elimina toda una clase de errores tipica de los sistemas de ingesta.
Bajo carga, este es el servicio que se estira. Corre con varias replicas y un autoescalado horizontal que las multiplica cuando la CPU media sube, y las reduce cuando baja. El esquema de la base de datos se gestiona con migraciones versionadas que se aplican al arrancar, de modo que el despliegue es reproducible y el estado del almacen siempre esta bajo control.
