Nube e infraestructura

Reconstrucción de la detección de incidentes con Kafka, Flink y OpenTelemetry

Un equipo redujo la latencia de detección de incidentes de 40 segundos a menos de 10 al reemplazar un agregador en Node.js por Apache Flink en Kubernetes y OpenTelemetry.

Traducido automáticamente del original en inglés.

Un pequeño equipo de ingeniería reconstruyó su plataforma automatizada de detección de incidentes, reduciendo la latencia de evento a métrica de más de 40 segundos a menos de 10. La nueva arquitectura se basa en Apache Kafka, Apache Flink ejecutándose en Kubernetes y OpenTelemetry para procesar miles de millones de eventos operativos diarios del lado del cliente. Si bien el sistema mejoró la velocidad y el aislamiento, el equipo informa que las tasas de recall fluctuaron entre el 64% y el 86%, lo que destaca la complejidad continua de una detección precisa de anomalías.

Qué ocurrió

La organización opera más de diez productos en la nube que sirven a millones de inquilinos (tenants). Anteriormente, su stack de monitoreo utilizaba un agregador en Node.js alimentado por una cola compartida en la nube, ejecutándose en aproximadamente 90 máquinas virtuales. Este sistema heredado sufría de alta latencia, problemas de "vecino ruidoso" donde el tráfico de otros inquilinos causaba retrasos, y costos que escalaban linealmente con cada nuevo producto incorporado. Los costos anuales de operación habían aumentado de $120,000 a $230,000, y la capa de caché frecuentemente alcanzaba el 100% de uso de CPU durante cambios rutinarios.

Para abordar estas limitaciones, el equipo diseñó una nueva pipeline centrada en cinco objetivos: latencia inferior a 10 segundos, rutas de consumo dedicadas para el aislamiento, escrituras idempotentes para garantizar la corrección durante las reproducciones, escalado de costos según el volumen y no según el número de funciones, y operabilidad basada en configuración. El resultado es un sistema donde incorporar una nueva experiencia solo requiere una pull request para actualizar la configuración, no un despliegue completo. El equipo midió el rendimiento mes a mes durante dieciocho meses, observando que, si bien la velocidad mejoró drásticamente, la precisión sigue por debajo de su objetivo.

Cómo funciona

La capa upstream utiliza Apache Kafka como bus de eventos. En lugar de consumir todo el flujo masivo de datos, el equipo implementó un filtro de suscripción del lado del servidor definido en código. Este filtro permite solo productos y experiencias específicos, descartando el tráfico experimental y sintético. Los datos filtrados llegan a un topic dedicado de Kafka con un período de retención de siete días, que sirve como ventana de reproducción para depuración o recuperación. Dos entradas laterales alimentan la pipeline: un servicio de contexto de inquilino que proporciona metadatos como shard y región, y un repositorio de configuración.

En el núcleo se encuentra un único job de Apache Flink 1.20 desplegado mediante el Flink Kubernetes Operator. Este job filtra eventos antiguos y errores, enriquece los datos con el contexto del inquilino utilizando un sidecar asíncrono con circuit breakers, y agrega métricas. Utiliza sketches HyperLogLog para estimar usuarios distintos afectados sin contarlos dos veces, logrando un margen de error de aproximadamente 1.5%. El estado se gestiona en RocksDB con checkpoints de 30 segundos hacia el almacenamiento de objetos, garantizando semántica exactly-once para archivos Parquet y escrituras idempotentes a un almacén clave-valor. Las métricas se exportan vía OpenTelemetry a una base de datos de series temporales compatible con Prometheus.

El plano de decisión, llamado AutoHOT, se ejecuta como un servicio en Go en dos regiones. Consume alertas de los detectores, cuantifica el impacto consultando el almacén agregado y aplica una matriz de severidad. Para evitar falsos positivos, verifica picos en la actividad de usuario durante quince minutos antes de crear un ticket. El motor utiliza bloqueos distribuidos para asegurar que solo una región procese una alerta a la vez, previniendo incidentes duplicados durante los failovers. También monitorea silencios en la telemetría, reconociendo que la falta de datos puede indicar un shard de base de datos completamente caído.

Detalles clave

  • La latencia disminuyó de más de 40 segundos a menos de 10 segundos para el procesamiento de evento a métrica.
  • El sistema procesa miles de millones de eventos diarios utilizando un único job de Apache Flink 1.20 en Kubernetes.
  • Los sketches HyperLogLog permiten conteos de usuarios distintos fusionables entre inquilinos y regiones con un error de ~1.5%.
  • La configuración se gestiona mediante un archivo YAML generado desde la configuración del producto, permitiendo la carga en caliente sin redespliegues.
  • El recall dentro del alcance alcanzó un pico del 86%, pero se estabilizó alrededor del 64% en meses difíciles, indicando margen de mejora.
  • Se ejecuta una verificación profunda sintética cada 15 minutos para validar toda la pipeline de alertamiento, desde la inyección hasta el cierre.

Por qué importa

Para los ingenieros que construyen plataformas de observabilidad, este caso de estudio demuestra los trade-offs entre la agregación estilo batch y el procesamiento de streams en tiempo real. Migrar a Flink permitió al equipo desacoplar el costo del número de funciones incorporadas, un factor crítico para productos SaaS en crecimiento. El uso de OpenTelemetry para la visibilidad end-to-end asegura que el propio sistema de monitoreo pueda ser monitorizado, previniendo fallos silenciosos que erosionan la confianza en los dashboards. Sin embargo, las tasas de recall fluctuantes recuerdan a los desarrolladores que datos más rápidos no significan automáticamente una lógica de detección mejor.

La elección arquitectónica de utilizar sinks idempotentes y topics dedicados de Kafka aborda puntos dolorosos comunes en sistemas distribuidos: duplicación de datos y contención de recursos. Al tratar los UIDs de los operadores como una API estable y ajustar los límites del autoscaler, el equipo eliminó reinicios frecuentes que previamente interrumpían el estado. Este enfoque ofrece una guía para equipos que luchan con stacks de monitoreo heredados, ruidosos, costosos y lentos que dependen de colas compartidas y cachés en memoria.

Qué puedes hacer

  • Evalúa tu actual pipeline de eventos por dependencias de colas compartidas que puedan causar problemas de vecino ruidoso durante cargas pico.
  • Considera usar HyperLogLog o estructuras de datos probabilísticas similares si necesitas conteos distintos a través de sistemas distribuidos.
  • Implementa filtrado del lado del servidor a nivel de bus de mensajes para reducir el volumen de procesamiento downstream y los costos.
  • Diseña tus jobs de stream processing con sinks idempotentes para permitir reproducciones seguras.

Herramientas de la Tienda de Bytechap

Seguir leyendo

Todos los artículos