Cloud & infrastructure

Rebuilding incident detection with Kafka, Flink, and OpenTelemetry

A team reduced incident detection latency from 40 seconds to under 10 by replacing a Node.js aggregator with Apache Flink on Kubernetes and OpenTelemetry.

A small engineering team rebuilt their automated incident detection platform, cutting event-to-metric latency from over 40 seconds to under 10. The new architecture relies on Apache Kafka, Apache Flink running on Kubernetes, and OpenTelemetry to process billions of daily client-side operational events. While the system improved speed and isolation, the team reports that recall rates fluctuated between 64% and 86%, highlighting the ongoing complexity of accurate anomaly detection.

What happened

The organization runs more than ten cloud products serving millions of tenants. Previously, their monitoring stack used a Node.js aggregator fed by a shared cloud queue, running on approximately 90 virtual machines. This legacy system suffered from high latency, noisy neighbor issues where other tenants’ traffic caused lag, and costs that scaled linearly with each new product onboarded. Annual run costs had risen from $120,000 to $230,000, and the cache tier frequently hit 100% CPU during routine changes.

To address these limitations, the team designed a new pipeline focused on five goals: sub-10-second latency, dedicated consumer paths for isolation, idempotent writes for correctness during replays, cost scaling with volume rather than feature count, and configuration-driven operability. The result is a system where onboarding a new experience requires only a pull request to update configuration, not a full deployment. The team measured performance month by month for eighteen months, noting that while speed improved dramatically, precision remains below their target.

How it works

The upstream layer uses Apache Kafka as the event bus. Instead of consuming the entire firehose, the team implemented a server-side subscription filter defined in code. This filter allows only specific products and experiences while dropping experimental and synthetic traffic. The filtered data lands on a dedicated Kafka topic with a seven-day retention period, which serves as the replay window for debugging or recovery. Two side inputs feed the pipeline: a tenant-context service providing metadata like shard and region, and a configuration repository.

At the core sits a single Apache Flink 1.20 job deployed via the Flink Kubernetes Operator. This job filters out old events and errors, enriches data with tenant context using an async sidecar with circuit breakers, and aggregates metrics. It uses HyperLogLog sketches to estimate distinct impacted users without double-counting, achieving about 1.5% error margin. The state is managed in RocksDB with 30-second checkpoints to object storage, ensuring exactly-once semantics for Parquet files and idempotent writes to a key-value store. Metrics are exported via OpenTelemetry to a Prometheus-compatible time-series database.

The decision plane, called AutoHOT, runs as a Go service in two regions. It consumes alerts from detectors, quantifies impact by querying the aggregated store, and applies a severity matrix. To prevent false positives, it checks for blips in user activity over fifteen minutes before creating a ticket. The engine uses distributed locks to ensure only one region processes an alert at a time, preventing duplicate incidents during failovers. It also monitors for silence in telemetry, recognizing that a lack of data can indicate a hard-down database shard.

Key details

  • Latency dropped from over 40 seconds to under 10 seconds for event-to-metric processing.
  • The system processes billions of events daily using a single Apache Flink 1.20 job on Kubernetes.
  • HyperLogLog sketches enable mergeable distinct-user counts across tenants and regions with ~1.5% error.
  • Configuration is managed via a YAML file generated from product config, allowing hot-loading without redeployments.
  • In-scope recall peaked at 86% but settled around 64% in difficult months, indicating room for improvement.
  • A synthetic deep-check runs every 15 minutes to validate the entire alerting pipeline from injection to closure.

Why it matters

For engineers building observability platforms, this case study demonstrates the trade-offs between batch-style aggregation and real-time stream processing. Moving to Flink allowed the team to decouple cost from the number of onboarded features, a critical factor for growing SaaS products. The use of OpenTelemetry for end-to-end visibility ensures that the monitoring system itself can be monitored, preventing silent failures that erode trust in dashboards. However, the fluctuating recall rates remind builders that faster data does not automatically mean better detection logic.

The architectural choice to use idempotent sinks and dedicated Kafka topics addresses common pain points in distributed systems: data duplication and resource contention. By treating operator UIDs as a stable API and tuning autoscaler boundaries, the team eliminated frequent restarts that previously disrupted state. This approach offers a blueprint for teams struggling with noisy, expensive, and slow legacy monitoring stacks that rely on shared queues and in-memory caches.

What you can do

  • Evaluate your current event pipeline for shared queue dependencies that may cause noisy neighbor issues during peak load.
  • Consider using HyperLogLog or similar probabilistic data structures if you need distinct counts across distributed systems.
  • Implement server-side filtering at the message bus level to reduce downstream processing volume and cost.
  • Design your stream processing jobs with idempotent sinks to allow safe replays without double-counting impacts.
  • Add synthetic end-to-end tests that inject fake alerts to verify your detection pipeline works even when real incidents are rare.
  • Separate operators in your stream graph based on resource profiles, splitting CPU-bound transformations from I/O-bound enrichment.

Tools from the Bytechap store

Keep reading

All stories