Valkey Cluster Architecture and Why the Gossip Bus Is the Ceiling
Valkey cluster is the horizontally scalable mode. Every shard owns a slot range, a primary serves writes, and replicas handle reads. Adding shards scales both write and read capacity. The cluster bus is a full mesh: every node talks to every other node. At around 1,000 nodes, gossip consumes roughly 1 to 2% CPU per node. There are exactly 16,384 slots total, so clusters approaching 2,000 nodes start seeing uneven slot heat. These two limits define the practical ceiling.
Client Topology and Connection Management at Scale
Valkey clients connect directly to one node, fetch the full cluster topology, then route requests to the correct shard by IP. Skipping DNS cuts latency and removes CoreDNS as a dependency. The topology command is expensive enough that Valkey caches the response internally. On Kubernetes, thousands of pods each opening their own connections can take down individual cache nodes. Connection pooling at either the client or Envoy layer prevents this. The Valkey Glide client consolidates topology refreshes across connections rather than running one refresh per connection.
Observability at the Slot Level with cluster slot stats
Most large Valkey clusters at Amazon are memory-bound, not CPU-bound. The cluster slot stats command, shipped in Valkey 8, surfaces per-slot key counts, network bytes in, network bytes out, and CPU. When a node runs hot and resharding is needed, the right move is to migrate the second-hottest slot off that node first, not the hottest. Moving the hottest slot is the most expensive operation. This command gives operators the data needed to make that call rather than guessing from pod-level CPU metrics alone.
What Broke at 2,000 Nodes: Three Bugs Found in Testing
AWS ran months of tests on 2,000-node clusters, killing up to 499 primary nodes at once. Three problems surfaced. First, surviving nodes retried connections to dead nodes every 100 milliseconds, burning CPU on connections that could never succeed. A throttling mechanism now spreads reconnects evenly across the 15-second cluster node timeout. Second, failure-report expiry used an O(n) list. Replacing it with a radix tree indexed by timestamp dropped CPU from 100% to 28-30% immediately after a mass failure. Third, a split-vote problem blocked replica promotion entirely.
The Split-Vote Fix: Lexicographic Ordering of Failover Requests
When two primaries fail simultaneously, their replicas both request votes from surviving primaries at the same moment. Because each primary casts only one vote per epoch, votes split and neither replica reaches quorum. Valkey maintainer Bin solved this by ordering failover requests lexicographically by shard ID. Shard one requests votes at T1, gets them, and promotes. Then shard three requests votes independently. A small jitter between requests eliminates the overlap. After this fix, a 2,000-node cluster recovers from 499 simultaneous primary failures in under one minute.
Pubsub Header Reduction: 2 KB Down to 16 Bytes
Valkey pubsub messages travel over the cluster bus. A 100-byte pubsub payload previously carried a 2 KB header containing full slot-assignment information. Pubsub messages have no use for slot data. A new lightweight message header keeps only the fields pubsub actually needs and cuts the overhead to 16 bytes per message. At cluster bus scale, that change reduces per-message overhead by more than 99%. Combined with the gossip and failover fixes, 1 billion RPS across 1,000 primaries is now verified.
Q&A
Was the 1 billion RPS test using connection pooling? Yes, connection pooling was used, and the benchmark scaled as many nodes as possible. ▶ 27:43
What profiling tool did you use during testing? The standard Linux perf tool, used to generate flame graphs showing per-function compute spend. ▶ 28:17
How do you ensure cluster initialization assigns primaries and replicas correctly rather than all pods becoming primaries? Today the operator is the proper solution. A workaround is replica migration config, which attaches idle nodes to primaries that have no replica. ▶ 29:35
Notable Quotes
in a 2000 node cluster where there are thousand primaries and we are sending a bunch of right workload we are able to reach up to 1 billion RPS Yes. So that was pretty cool. Sarthak Aggarwal · ▶ 18:47
a typical request inside Valky only takes on the order of like one microcond but connection establishment takes hundreds of microsconds. So like if you’re doing a connection for every request like you would need a lot of you’re wasting a lot of CPU. Madelyn Olson · ▶ 27:55
we introduced a lightweight message header which uh basically reduced the uh 2k uh message header to 16 bytes just kept the information that we really need Sarthak Aggarwal · ▶ 25:23
we practically saw from 100% the we were around 28 to 30% uh of compute just after failing. Sarthak Aggarwal · ▶ 21:29
Key Takeaways
- A 2,000-node Valkey cluster sustains 1 billion RPS with 1,000 active primaries.
- Lexicographic ordering of failover requests eliminates split votes during simultaneous primary failures.
- Pubsub cluster bus overhead dropped from 2 KB to 16 bytes per message with a lightweight header.
About the Speaker(s)
Madelyn Olson is a co-creator and maintainer of Valkey and Principal Engineer at AWS. She focuses on building secure and highly reliable features, with a passion for working with open-source communities.
Sarthak Aggarwal is a software engineer at AWS contributing to Valkey and OpenSearch. He is passionate about large-scale distributed databases and the performance challenges that come with them.