Cloud & Infrastruktur

Neuaufbau der Incident-Erkennung mit Kafka, Flink und OpenTelemetry

Ein Team reduzierte die Latenz bei der Incident-Erkennung von 40 Sekunden auf unter 10 Sekunden, indem es einen Node.js-Aggregator durch Apache Flink auf Kubernetes und OpenTelemetry ersetzte.

Automatisch aus dem englischen Original übersetzt.

Ein kleines Engineering-Team hat seine automatisierte Plattform zur Incident-Erkennung neu aufgebaut und dabei die Latenz von Ereignis zu Metrik von über 40 Sekunden auf unter 10 Sekunden gesenkt. Die neue Architektur basiert auf Apache Kafka, Apache Flink (ausgeführt auf Kubernetes) und OpenTelemetry, um Milliarden von täglichen clientseitigen Betriebsereignissen zu verarbeiten. Während das System Geschwindigkeit und Isolation verbesserte, berichtet das Team, dass die Recall-Raten zwischen 64 % und 86 % schwankten, was die anhaltende Komplexität einer präzisen Anomalieerkennung unterstreicht.

Was passiert ist

Die Organisation betreibt mehr als zehn Cloud-Produkte, die Millionen von Mandanten bedienen. Zuvor nutzte ihr Monitoring-Stack einen Node.js-Aggregator, der von einer gemeinsamen Cloud-Warteschlange gespeist wurde und auf etwa 90 virtuellen Maschinen lief. Dieses Legacy-System litt unter hoher Latenz, „Noisy Neighbor“-Problemen (bei denen der Traffic anderer Mandanten Verzögerungen verursachte) sowie Kosten, die linear mit jedem neuen onboardeten Produkt stiegen. Die jährlichen Betriebskosten waren von 120.000 USD auf 230.000 USD angestiegen, und die Cache-Schicht erreichte während routinemäßiger Änderungen häufig eine CPU-Auslastung von 100 %.

Um diese Einschränkungen zu beheben, entwarf das Team eine neue Pipeline mit fünf Zielen: Latenz unter 10 Sekunden, dedizierte Consumer-Pfade für Isolation, idempotente Schreibvorgänge für Korrektheit bei Wiederholungen, Skalierung der Kosten nach Volumen statt nach Funktionsumfang sowie konfigurierungsgetriebene Betriebbarkeit. Das Ergebnis ist ein System, bei dem das Onboarding einer neuen Experience nur einen Pull Request zur Aktualisierung der Konfiguration erfordert, nicht eine vollständige Bereitstellung. Das Team maß die Performance monatlich über achtzehn Monate hinweg und stellte fest, dass sich zwar die Geschwindigkeit dramatisch verbesserte, die Präzision jedoch weiterhin unter dem Zielwert lag.

Wie es funktioniert

Die Upstream-Ebene nutzt Apache Kafka als Event-Bus. Anstatt den gesamten Datenstrom („Firehose“) zu konsumieren, implementierte das Team einen serverseitigen Abonnementfilter, der im Code definiert ist. Dieser Filter lässt nur bestimmte Produkte und Experiences zu und verwirft experimentellen sowie synthetischen Traffic. Die gefilterten Daten landen in einem dedizierten Kafka-Topic mit einer Aufbewahrungsdauer von sieben Tagen, das als Replay-Fenster für Debugging oder Wiederherstellung dient. Zwei Side Inputs speisen die Pipeline: ein Tenant-Context-Dienst, der Metadaten wie Shard und Region bereitstellt, sowie ein Konfigurations-Repository.

Im Kern sitzt ein einzelner Apache Flink 1.20 Job, der über den Flink Kubernetes Operator bereitgestellt wird. Dieser Job filtert alte Ereignisse und Fehler aus, angereichert Daten mit Tenant-Kontext unter Verwendung eines asynchronen Sidecars mit Circuit Breakern und aggregiert Metriken. Er verwendet HyperLogLog-Skizzen, um die Anzahl der betroffenen eindeutigen Benutzer ohne Doppelzählung zu schätzen, und erreicht eine Fehlerrate von etwa 1,5 %. Der Zustand wird in RocksDB verwaltet, mit Checkpoints alle 30 Sekunden in den Objektspeicher, was Exactly-Once-Semantik für Parquet-Dateien und idempotente Schreibvorgänge in einen Key-Value-Store gewährleistet. Metriken werden über OpenTelemetry in eine Prometheus-kompatible Zeitreihendatenbank exportiert.

Die Entscheidungsebene, genannt AutoHOT, läuft als Go-Dienst in zwei Regionen. Sie konsumiert Alarme von Detektoren, quantifiziert die Auswirkung durch Abfragen des aggregierten Speichers und wendet eine Schweregrad-Matrix an. Um False Positives zu verhindern, prüft sie auf kurzzeitige Ausschläge („Blips“) in der Benutzeraktivität über fünfzehn Minuten, bevor ein Ticket erstellt wird. Die Engine verwendet verteilte Sperren, um sicherzustellen, dass nur eine Region einen Alarm gleichzeitig verarbeitet, was doppelte Incidents während Failovers verhindert. Zudem überwacht sie Stille in der Telemetrie, da fehlende Daten auf einen vollständig ausgefallenen Datenbank-Shard hinweisen können.

Wichtige Details

  • Die Latenz für die Verarbeitung von Ereignis zu Metrik sank von über 40 Sekunden auf unter 10 Sekunden.
  • Das System verarbeitet täglich Milliarden von Ereignissen mit einem einzigen Apache Flink 1.20 Job auf Kubernetes.
  • HyperLogLog-Skizzen ermöglichen zusammenführbare Zählungen eindeutiger Benutzer über Mandanten und Regionen hinweg mit ~1,5 % Fehlerquote.
  • Die Konfiguration wird über eine YAML-Datei verwaltet, die aus der Produktkonfiguration generiert wird, was Hot-Loading ohne erneute Bereitstellungen ermöglicht.
  • Die Recall-Rate innerhalb des Gültigkeitsbereichs erreichte einen Höchstwert von 86 %, pendelte sich aber in schwierigen Monaten bei etwa 64 % ein, was Verbesserungspotenzial aufzeigt.
  • Ein synthetischer Deep-Check läuft alle 15 Minuten, um die gesamte Alarmierungspipeline von der Injektion bis zur Schließung zu validieren.

Warum es wichtig ist

Für Ingenieure, die Observability-Plattformen entwickeln, demonstriert diese Fallstudie die Kompromisse zwischen Batch-artiger Aggregation und Echtzeit-Stream-Verarbeitung. Der Wechsel zu Flink ermöglichte es dem Team, die Kosten von der Anzahl der onboardeten Funktionen zu entkoppeln – ein kritischer Faktor für wachsende SaaS-Produkte. Der Einsatz von OpenTelemetry für End-to-End-Sichtbarkeit stellt sicher, dass das Monitoring-System selbst überwacht werden kann, was stille Ausfälle verhindert, die das Vertrauen in Dashboards erodieren. Allerdings erinnern die schwankenden Recall-Raten Entwickler daran, dass schnellere Daten nicht automatisch bessere Erkennungslogik bedeuten.

Die architektonische Entscheidung, idempotente Sinks und dedizierte Kafka-Topics zu verwenden, adressiert häufige Schmerzpunkte in verteilten Systemen: Datenduplikation und Ressourcenkonkurrenz. Indem das Team die UIDs der Operatoren als stabile API behandelte und die Grenzen des Autoscalers feinabstimmte, eliminierte es häufige Neustarts, die zuvor den Zustand störten. Dieser Ansatz bietet einen Bauplan für Teams, die mit lauten, teuren und langsamen Legacy-Monitoring-Stacks kämpfen, die auf gemeinsamen Warteschlangen und In-Memory-Caches basieren.

Was Sie tun können

  • Bewerten Sie Ihre aktuelle Event-Pipeline auf Abhängigkeiten von gemeinsamen Warteschlangen, die bei Spitzenlasten zu „Noisy Neighbor“-Problemen führen könnten.
  • Erwägen Sie die Verwendung von HyperLogLog oder ähnlichen probabilistischen Datenstrukturen, wenn Sie eindeutige Zählungen über verteilte Systeme hinweg benötigen.
  • Implementieren Sie serverseitiges Filtering auf Ebene des Message Buses, um das nachgelagerte Verarbeitungsvolumen und die Kosten zu reduzieren.
  • Entwerfen Sie Ihre Stream-Processing-Jobs mit idempotenten Sinks, damit sichere Wiederholungen möglich sind, ohne Auswirkungen doppelt zu zählen.
  • Fügen Sie synthetische End-to-End-Tests hinzu, die Fake-Alarme injizieren, um zu verifizieren, dass Ihre Erkennungspipeline funktioniert, auch wenn echte Incidents selten sind.
  • Trennen Sie Operatoren in Ihrem Stream-Graphen basierend auf ihren Ressourcenprofilen und teilen Sie CPU-lastige Transformationen von I/O-lastigen Anreicherungen ab.

Tools aus dem Bytechap-Shop

Weiterlesen

Alle Artikel