What is a ZooKeeper 'watch storm' in the context of Kafka, and why did it hurt at scale?
answer
- ZK watch = one-shot notification
- re-read + re-register after every fire
- broker death → many znodes change → flood
- ZK becomes metadata-change-rate bottleneck
- KRaft replaces watches with log fetch (pull)
basics
~20 sKafka watched many ZooKeeper znodes to learn about changes. One event (like a broker dying) could fire huge numbers of watch notifications at once, flooding ZooKeeper and the controller with work and slowing metadata propagation.
solid answer
~50 sZooKeeper notifies clients of changes via **watches**: a client registers a one-shot watch on a znode and gets a notification when it changes. Kafka relied heavily on watches to detect broker membership, ISR, and leadership changes. The trouble is watches are one-shot — after firing you must re-read and re-register — and a single event (e.g., a broker going down) could touch many znodes and wake many watchers simultaneously. At large partition/broker counts this produced bursts of notifications plus follow-up reads — a 'watch storm' — that loaded the ZK ensemble and the controller, increasing latency and sometimes causing notification backlogs or missed updates. It also made ZK a throughput bottleneck for metadata change rate. KRaft eliminates watches entirely: brokers tail an ordered metadata log via fetch, so change notification is just normal log consumption with no per-znode watch fan-out.
go deeper
Know Kafka used ZooKeeper 'watches' to learn about changes, and one event could trigger a flood that slowed things down.
Explain watches are one-shot, requiring re-read/re-register, and how a broker failure causes a burst of notifications.
Tie watch storms to ZK throughput limits and propagation-latency spikes during correlated failures; contrast with KRaft fetch.
Discuss push-vs-pull notification trade-offs, intermediate-state collapsing, and why pull-based log replication scales metadata change rate better.
## What a ZooKeeper watch is **ZooKeeper (ZK)** lets a client register a **watch** on a znode (a node in ZK's data tree). When that znode changes, ZK sends the client a **one-shot notification**. 'One-shot' is the key constraint: the watch fires **exactly once**, then is gone. To keep tracking the node you must, on each notification, **re-read the node and re-register a new watch**. Between the change and your re-read, more changes can happen. ## How Kafka used watches Legacy Kafka leaned on watches to stay current: - Watches on `/brokers/ids` to detect brokers joining/leaving. - Watches related to ISR / partition state / leadership. - The controller watched membership and reacted to changes by recomputing leaders and pushing updates. ## Anatomy of a 'watch storm' A **watch storm** is a sudden burst of many watch notifications (and the read traffic that follows them) triggered by a single underlying event. Examples: - A **broker dies**: every partition it led now needs a new leader. Many znodes change; many watchers wake at once; the controller must process a flood of changes and then fan out updates. - A **rolling restart / mass reassignment**: large numbers of metadata changes in a short window. Because watches are one-shot, each notification is followed by a **re-read + re-register**, multiplying the load. At small scale this is invisible. At **hundreds of thousands of partitions / many brokers**, the burst can: - Saturate the ZK ensemble's request throughput. - Backlog notifications, increasing the latency between a real change and brokers learning about it. - Risk **missing intermediate states** (you only see the latest value when you re-read), so rapid successive changes collapse — fine for some cases, lossy for others. The net effect: **ZK became a bottleneck on the rate of metadata change**, and propagation latency spiked exactly when the cluster was already stressed (a failure). ## Why this capped scalability The cost of reacting to events grew with cluster size, and the reaction happened through a coordination service not designed for high-churn, high-fan-out Kafka-metadata workloads. This is a structural limit, not a tuning problem. ## How KRaft removes it KRaft (KIP-500) has **no watches at all**. Metadata changes are **records appended to a replicated log**. Brokers learn about changes by **fetching the log** (ordinary, batched, pull-based replication) — the same mechanism Kafka already uses for data. Benefits: - No one-shot re-register churn; you just keep fetching from your last offset. - No fan-out 'wake everyone' storm — consumers pull at their own pace, in order, in batches. - You see **every** change in order (no collapsing of intermediate states beyond snapshotting), and propagation is bounded by fetch latency. ## Nuance / edge cases - Watch storms were worst during **correlated failures** (the moment you most need fast, calm metadata propagation). - The one-shot nature also created subtle correctness work in Kafka's ZK client code (re-registration races). - KRaft trades 'push notification' for 'pull/fetch', which is generally cheaper and more predictable under load.
- Why does the one-shot nature of watches make storms worse?After each notification you must re-read the znode and register a fresh watch. Under a burst, every notification spawns follow-up read/register traffic, multiplying load and creating re-registration races.
- How does KRaft notify brokers of metadata changes without watches?Brokers tail the replicated `__cluster_metadata` log by fetching from their last offset — pull-based, batched, ordered replication. No per-znode watch fan-out exists.
saying these in an interview costs you the question
- Saying ZooKeeper watches are persistent subscriptions — they're one-shot and must be re-registered.
- Claiming KRaft still uses watches internally — it uses log fetch, no watches.
- Describing a watch storm as merely 'high CPU' — the core issue is fan-out of one-shot notifications plus re-read traffic saturating ZK at scale, worst during failures.