Migración del pipeline de datos a arquitectura de streaming en tiempo real
Justificación técnica y de negocio para la transición de ETL batch nocturno (Apache Airflow) a procesamiento en streaming continuo (Apache Kafka + Apache Flink) en todos los hospitales cliente.
Resumen ejecutivo
NovaMed Analytics ha identificado que la latencia actual del pipeline de datos (~18h) es el principal bloqueador comercial en el segmento enterprise. Los hospitales cliente no pueden operar dashboards clínicos con datos del día anterior, ni activar alertas automatizadas sobre eventos en tiempo real.
Después de un spike técnico de 2 semanas y análisis de 4 alternativas, el equipo recomienda Apache Kafka como bus de eventos + Apache Flink para procesamiento stateful en streaming, desplegados en infraestructura propia sobre AWS EKS. Esta combinación ofrece el equilibrio óptimo entre capacidad técnica, coste, cumplimiento RGPD y viabilidad de migración en el plazo disponible.
Contexto y problema
Arquitectura actual
El sistema actual procesa datos de 23 hospitales mediante 47 pipelines Apache Airflow que ejecutan transformaciones ETL en ventanas batch de 6 horas, con entrega final al data warehouse (Snowflake) entre las 02:00 y las 06:00 AM. Los dashboards de los hospitales se actualizan una vez al día, a primera hora de la mañana.
Impacto en negocio
Tres propuestas enterprise están en etapa final de negociación con hospitales de más de 500 camas. Los tres han condicionado la firma a la disponibilidad de datos en tiempo real para sus casos de uso específicos:
Origen de la limitación técnica
Airflow fue elegido en 2023 como solución pragmática para el MVP. En ese momento, NovaMed operaba 3 hospitales piloto con baja frecuencia de datos y dashboards de reporting mensual. La arquitectura no fue diseñada para streaming y su extensión natural no es técnicamente viable sin una reescritura completa de los operadores y el modelo de datos subyacente.
Análisis de alternativas evaluadas
Se evaluaron cuatro enfoques durante el sprint de investigación (semanas 22-23, 2026). El equipo de data engineering (3 ingenieros) dedicó 2 semanas a spikes técnicos y análisis de coste.
| Alternativa | Latencia | Coste mensual | RGPD | Esfuerzo migración | Puntuación |
|---|---|---|---|---|---|
| ✅ Kafka + Flink (self-hosted) Opción elegida |
< 30s | ~2.400€ | Control total | Alta (12 sem.) | 87/100 |
| Confluent Cloud Kafka gestionado |
< 30s | ~7.200€ | Datos en Confluent | Media (6 sem.) | 61/100 |
| Databricks DLT Delta Live Tables |
~2 min | ~8.100€ | Certificado RGPD | Baja (3 sem.) | 58/100 |
| Postgres LISTEN/NOTIFY Solución casera |
< 5s | ~400€ | Control total | Baja (4 sem.) | 29/100 |
Por qué se descartaron
Confluent Cloud: Los datos de pacientes (categoría especial RGPD, Art. 9) no pueden residir en infraestructura de terceros sin DPA específico y evaluación de impacto (DPIA) completa. El equipo legal estimó 4+ meses para la aprobación regulatoria, incompatible con el timeline.
Databricks DLT: Coste 3x superior al autoalojado. La latencia mínima de ~2 minutos no satisface el SLA requerido por los hospitales (<5 min para alertas clínicas). Vendor lock-in con modelo de precios por DBU difícil de predecir.
Postgres LISTEN/NOTIFY: Validado en laboratorio hasta ~50.000 eventos/hora. NovaMed proyecta superar 2M eventos/hora con 50 hospitales en el plan de crecimiento 2027. La solución no escala horizontalmente y carece de replay, consumer groups y garantías de entrega exactamente-una-vez (EOS).
Decisión: Arquitectura Kafka + Flink
Componentes del stack
Apache Kafka 3.7 (self-managed sobre AWS MSK Serverless para arranque rápido, migración a EC2 dedicado en Q4): bus de eventos con retención de 7 días, 3 brokers, replicación factor 3. Cifrado TLS en tránsito + KMS en reposo.
Apache Flink 1.19 (operador Kubernetes en EKS): procesamiento stateful con exactamente-una-vez (EOS) mediante checkpoints en S3 cada 60 segundos. Jobs separados para alertas clínicas (baja latencia) y aggregations analíticas (alta throughput).
Schema Registry (Karapace, open source): control de contratos de datos, evolución de esquemas sin downtime, auditoría completa de cambios.
RGPD compliance: datos de pacientes tokenizados en origen (hospital) mediante vault de tokens propio, claves por hospital, audit log completo en OpenSearch, DPIA ya archivada con la AEPD.
Exactly-once semantics y RGPD
El spike técnico identificó un problema con EOS cuando los datos incluyen campos de identidad de paciente (NHC, DNI). La solución adoptada: tokenización en el conector de ingesta antes de publicar en Kafka. El NHC real nunca viaja en el bus de eventos; solo el token pseudoanonimizado. El vault de tokens está fuera del cluster Kafka, accesible únicamente por el servicio de tokenización. El equipo legal revisó y aprobó este diseño el 5 de junio de 2026.
Plan de migración de los 47 pipelines
Clasificación de pipelines
Timeline de ejecución
Gestión de riesgos
Criterios de éxito (OKRs del proyecto)
| Métrica | Baseline actual | Target semana 30 | Target final (sem. 39) |
|---|---|---|---|
| Latencia P95 (alertas clínicas) | 18 horas | < 5 minutos | < 30 segundos |
| Disponibilidad del pipeline | 99,1% (Airflow) | 99,5% | 99,9% |
| Contratos enterprise firmados | 0 de 3 | 3 de 3 | 3 de 3 + onboarding |
| Pipelines migrados | 0 de 47 | 20 de 47 | 47 de 47 |
| Pérdida de datos (eventos perdidos) | N/A | 0 eventos perdidos | 0 eventos perdidos |
| Coste infraestructura mensual | ~800€ | < 4.000€ | < 3.000€ |
Plan de rollback
Cada pipeline tiene un feature flag en la capa de configuración central. El interruptor streaming_enabled=false redirige el tráfico de vuelta al DAG de Airflow correspondiente sin downtime. Los dashboards de los hospitales no perciben el cambio; simplemente vuelven a la cadencia de actualización anterior.
Condiciones de rollback automático: El sistema monitoriza continuamente (1) divergencia entre outputs Kafka y Airflow superior al 0,001%, (2) latencia P95 de Flink superior a 30 segundos durante más de 5 minutos consecutivos, (3) tasa de error en jobs Flink superior al 0,1% en ventana de 10 minutos. Cualquier condición activa rollback automático y página al on-call.
Los tres pipelines críticos (facturación, alertas, reporting) requieren autorización manual del CTO para el cutover final y para desactivar el fallback a Airflow. Esta restricción se elimina 30 días después del cutover si las métricas son correctas.
Reader Testing — Validación del documento
El documento fue testado con un subagente Claude sin contexto previo de la conversación. Se verificó que cualquier lector puede responder las preguntas clave sin necesidad de información adicional.