We replaced a Node.js aggregation service that consumed a shared cloud queue with a dedicated Kafka‑Flink pipeline that streams client‑side telemetry through OpenTelemetry on Kubernetes. The new design cuts event‑to‑metric latency from more than 40 seconds to under 10 seconds and isolates the detection path from noisy neighbours.
OpenTelemetry incident detection architecture
The rebuilt system follows a straight‑through flow: product services emit operational events to an analytics gateway, which forwards them to an Apache Kafka bus. A server‑side subscription filter – a generated ~770‑line YAML allow‑list – selects only the relevant products and experiences, directing the filtered stream to a dedicated Kafka topic with a seven‑day retention window. This topic serves as the replay buffer for the downstream pipeline.
A single Apache Flink 1.20 job, deployed via the Flink Kubernetes Operator, consumes the topic. The job enriches events with tenant‑context data and configuration from side‑input services, then computes per‑minute error‑rate and volume‑drop metrics. Results are exported through OpenTelemetry to a Prometheus‑compatible time‑series store and a commercial metrics service where detector logic runs.
Impact aggregates are written to a multi‑region key‑value store and to Apache Parquet files in object storage. A lightweight Go GraphQL service (the Impact API) reads these stores to answer “who is impacted and how many”. The AutoHOT engine – a Go service running active‑active in two regions – consumes detector alerts, queries the Impact API, applies a severity matrix, suppresses transient blips, creates incident tickets, pages responders, and continuously re‑evaluates impact until it clears.
Implementation choices and trade‑offs
- Filtering at the bus. By moving the allow‑list into Kafka subscription filters, the pipeline processes only ~55 % of total bus traffic, reducing compute load and storage cost for the seven‑day retention window.
- Single Flink job. Consolidating all detection logic into one Flink job simplifies deployment and operational monitoring, while still supporting per‑experience configuration via the side‑input repository.
- Idempotent writes. The Flink job emits idempotent updates, ensuring that restarts or replay of the seven‑day window do not double‑count impact.
- Isolation. The dedicated Kafka topic and separate Flink job eliminate dependence on the shared queue that previously caused two incidents due to tenant‑level lag.
- Cost model. Running costs now grow with event volume rather than the number of onboarded experiences; the earlier system’s CPU usage spiked to 100 % during routine changes, whereas the new pipeline scales horizontally on Kubernetes.
Operational impact
Latency targets are met: event‑to‑metric latency stays under 10 seconds, and detector decisions are made on the next minute boundary. Monthly measurements over eighteen months show in‑scope recall rising from ~60 % to a peak of 86 % (with a later dip to 64 % in a bad month), while precision remains an area for improvement.
Cost shifted from $120 K to $230 K per year as products were added, but the scaling is now predictable because it is tied to event volume. The architecture’s operability goal is satisfied by a single source of truth for filters, transforms, thresholds, and sinks – onboarding a new experience is a pull‑request change to the YAML filter, not a deployment of new services.
Related CloudNinjas coverage: hands-on guides.
What This Means For Practitioners
Teams building real‑time incident detection should consider moving filtering upstream into the event bus to shrink downstream workloads, and use a stateful stream processor like Flink for low‑latency metric derivation. Idempotent output and dedicated consumption paths improve reliability and isolate detection from unrelated traffic spikes. Exporting metrics via OpenTelemetry keeps the observability stack vendor‑agnostic while feeding both open‑source and commercial downstream detectors. Finally, aligning cost with event volume rather than feature count simplifies budgeting as services scale.


