10 min read

Beyond the Speed of Light: Engineering 100 Million Writes Per Second with Hierarchical Consensus

100M Writes Per Second via Hierarchical Consensus

Imagine this: It’s 3:00 AM, and you’re looking at a Grafana dashboard. A global event—perhaps a massive product drop or a sudden viral surge—is hammering your metadata layer. In a fraction of a second, tens of millions of concurrent users are vying for the same state. Your traditional Raft-based consensus cluster is choking. The leader is pinned at 100% CPU, heartbeats are dropping, and the tail latency for a simple "write" has ballooned from 5ms to 5 seconds.

This is the Global Metadata Wall.

In the world of distributed systems, we’ve long accepted the CAP theorem as a static law of gravity. You want consistency and availability? Prepare to pay the latency tax. You want global scale? Prepare to sacrifice strong consistency. But what if you need to sustain 100 million writes per second across the planet without turning your system into an eventual-consistency nightmare?

At this scale, the "standard" ways of doing things—single-leader Paxos, vanilla gossip protocols, and synchronous replication—don't just slow down; they disintegrate. To solve this, we had to rethink the very geometry of how machines agree on reality.

Today, we’re diving deep into the architecture of Hierarchical Consensus and Sharded Gossip Protocols: the twin pillars that allow us to shatter the throughput ceilings of the past.


The Death of the Global Leader

To understand where we’re going, we have to look at why we’re stuck. Traditional consensus protocols like Raft or Multi-Paxos rely on a centralized leader. Every write must pass through this bottleneck to be sequenced.

When you aim for 100 million writes per second, the math becomes grim:

  1. Network Saturation: Even with a 100GbE NIC, the leader cannot physically ingest and broadcast the metadata packets for 100M operations.
  2. The Heartbeat Problem: As the cluster grows to hundreds of nodes to handle read volume, the leader spends more time managing heartbeats and peer lists than actually committing logs.
  3. Speed of Light: If your leader is in US-East and a write comes from Tokyo, you are physically tethered to a ~200ms round-trip time (RTT).

The industry hype around "Serverless Databases" and "Global State Stores" often glosses over these physics. They offer "Global Tables," but under the hood, they are either using asynchronous replication (losing consistency) or forcing you into a single-region leader (losing performance).

We needed a paradigm shift. We needed to stop thinking about a single "Global Log" and start thinking about Federated Sovereignty.


Architecture: The Hierarchical Consensus Model

Hierarchical Consensus is the art of breaking a global problem into localized certainties. Instead of one giant Raft group of 1,000 nodes, we organize nodes into a recursive tree of consensus cells.

1. Local Consensus Cells (The Leaves)

At the edge—closest to the user—we deploy small, hyper-fast consensus groups (usually 3 or 5 nodes). These cells use a highly optimized version of Raft implemented in Rust, leveraging zero-copy memory mapping and io_uring for disk I/O.

These local cells possess "lease-based sovereignty" over a specific shard of the metadata keyspace. If a user in London updates a key, the London cell reaches consensus internally in sub-millisecond time. The write is "locally committed."

2. Regional Aggregators (The Branches)

Local cells don't talk to the global core individually. That would create a $O(N^2)$ communication storm. Instead, they stream their committed logs to Regional Aggregators.

The Aggregator’s job is Log Compaction and Batching. It takes 10,000 individual metadata updates and collapses them into a single "Merkle Proof" or a highly compressed state delta.

3. The Global Root (The Trunk)

The Global Root only sees the aggregated proofs from the Regional Branches. It doesn't care about individual key-value pairs; it only cares about the consistent ordering of regional batches.

By the time a write reaches the global level, it has been "pre-verified" by the hierarchy. This allows the root to process millions of logical writes by only handling a few thousand high-level proofs per second.


Sharded Gossip: Silencing the Noise

While the Hierarchical Consensus handles the "Ordering" of writes, we still have the problem of Propagation. How do you tell 10,000 nodes across the world about these changes without melting the backbone fiber?

Standard Gossip protocols (like SWIM) work by having nodes randomly pick peers to swap data with. At 100M writes/sec, the amount of redundant data being gossiped would exceed the total bandwidth of the internet.

Enter Sharded Gossip Protocols.

The "Interest-Based" Topology

In a Sharded Gossip model, we don't treat the network as a flat pool. We shard the gossip traffic based on the metadata keyspace.

  • Topic Awareness: Nodes subscribe only to the shards of metadata they are currently serving or caching.
  • Partial View Clusters: Instead of knowing every node in the global fleet, a node only maintains a "partial view" of its immediate neighbors and a few long-distance "express links" to other regions.

Adaptive Push-Pull

To reach 100M writes/sec, we utilize an Adaptive Push-Pull mechanism. When the network is quiet, nodes "push" updates aggressively. When congestion is detected (via eBPF-monitored socket buffers), the protocol automatically switches to a "pull-on-demand" model.

// A simplified look at an Adaptive Gossip Heartbeat in Rust
struct GossipNode {
    id: NodeId,
    interest_shards: Vec<ShardId>,
    transport: QuicTransport,
}

impl GossipNode {
    async fn propagate_update(&self, update: MetadataUpdate) {
        let targets = self.topology.get_peers_for_shard(update.shard_id);
        
        for peer in targets {
            if self.network_pressure() > THRESHOLD {
                // Switch to Pull: Just send the Update ID (the hash)
                self.transport.send_header(peer, update.hash).await;
            } else {
                // Push: Send the full metadata payload
                self.transport.send_full_payload(peer, update).await;
            }
        }
    }
}

Engineering the Data Plane: 100M Writes Under the Hood

You cannot hit these numbers using the standard Linux networking stack. The overhead of context switching between kernel space and user space is a death sentence. To achieve 100 million writes, we move almost everything into User-Space Networking.

DPDK and eBPF: The Secret Weapons

By using the Data Plane Development Kit (DPDK), our consensus engine bypasses the Linux kernel entirely. It takes direct control of the NIC (Network Interface Card), pulling packets directly into the application memory.

Furthermore, we use eBPF (Extended Berkeley Packet Filter) to perform "Early Rejection" and "XDP Load Balancing." If a consensus vote arrives that is stale or malformed, the eBPF program at the NIC level drops it before it even touches the CPU. This preserves precious cycles for the state machine.

The Storage Engine: LSM-Trees on Steroids

A write isn't "done" until it's durable. Writing 100 million entries per second to disk is impossible for standard relational databases. We utilize a custom-built Log-Structured Merge-Tree (LSM-Tree) optimized for NVMe drives.

  • Immutable Memtables: We use lock-free skip lists in memory to handle the incoming 100M writes.
  • Zero-Copy Serialization: We use FlatBuffers or Cap’n Proto instead of JSON or Protobuf. This allows us to map the data directly from the network packet to the disk buffer without an intermediate "deserialization" step.

Conflict Resolution: CRDTs and Semantic Consistency

When you shard consensus and use gossip, you occasionally run into the "Split-Brain" boogeyman. What if two local cells accept a write for the same key at the exact same microsecond?

Traditional systems use "Last Writer Wins" (LWW), which is dangerous. We utilize Conflict-free Replicated Data Types (CRDTs).

By encoding the metadata as G-Counters (Grow-only counters) or OR-Sets (Observed-Remove sets), we ensure that the state is mathematically guaranteed to converge, regardless of the order in which updates arrive at different nodes. The Hierarchical Consensus provides the "strong ordering" for critical operations (like account balances), while the Sharded Gossip handles "semantic convergence" for less critical metadata (like presence indicators or session tags).


The Infrastructure Scale: The Raw Math

Let's do the "Back of the Envelope" calculation for 100,000,000 writes per second.

  • Nodes: 1,000 high-performance servers.
  • Per Node Load: 100,000 writes per second per node.
  • Packet Size: 128 bytes per metadata update.
  • Bandwidth per Node: $100,000 \times 128 \text{ bytes} \approx 12.8 \text{ MB/s}$ of raw data.

On the surface, 12.8 MB/s sounds easy. But in a distributed system, that 12.8 MB/s is multiplied by:

  1. Replication Factor: (e.g., 3x) = 38.4 MB/s.
  2. Gossip Overhead: (e.g., 5x) = 192 MB/s.
  3. Consensus Chatter: (Raft AppendEntries, Heartbeats, etc.) = ~400 MB/s.

Suddenly, each node is handling nearly 4Gbps of pure metadata traffic. When you factor in the CPU cost of cryptographic signing, Merkle tree updates, and state machine transitions, you realize why this requires specialized engineering.

We aren't just building an application; we are building a Metadata Switch.


Real-World Context: Why This Hype Matters

Over the last 18 months, we’ve seen a massive surge in "Real-Time AI" and "Edge Computing." Everyone from OpenAI to Vercel is trying to solve the problem of Global Context.

If an AI agent is interacting with you, it needs to know your state now, regardless of whether you’re in New York or Singapore. The "Hype" around Vector Databases and Edge Functions has hit a wall because the underlying consistency layers were built for the 2010s—an era where 10,000 writes per second was considered "big."

The industry is realizing that Global Consistency is the new bottleneck for AI. If your metadata layer can't keep up with the inference speed, your "Real-Time AI" is just "Slow-Time AI" with a better UI. This is why Hierarchical Consensus is moving from a theoretical research paper to the core of the modern stack.


Observability: Debugging a 100M Ops System

How do you even know if it’s working? At 100M writes/sec, traditional logging is a denial-of-service attack on yourself. You cannot log.info() 100 million times a second.

1. Statistical Sampling and Distributed Tracing

Instead of logging every write, we use Probabilistic Tracing. We attach a "Trace Bit" to 0.001% of requests. This gives us a statistically significant view of the system's health without the overhead.

2. The "State Delta" Dashboard

Instead of monitoring "Writes Per Second," we monitor "Global Skew." This is the time delta between the first commit in a local cell and the final convergence in the global root. If the skew exceeds 500ms, the system automatically triggers backpressure.

3. TLA+ Formal Verification

When you’re dealing with hierarchical consensus, "testing" isn't enough. You can't simulate 100 million concurrent writes in a staging environment easily. We use TLA+ (Temporal Logic of Actions) to mathematically prove that our protocols are free from deadlocks and livelocks before we write a single line of Rust code.


The Path to One Billion

100 million writes per second is a milestone, but it isn't the finish line. As we move toward a world of trillions of IoT devices and ubiquitous AI agents, we are already looking at the Billion Write horizon.

The lessons are clear:

  • Centralization is the enemy of scale.
  • The kernel is too slow for the modern data plane.
  • Consistency doesn't have to be a global lock; it can be a hierarchical agreement.

Engineering at this scale requires a certain level of fearlessness. You have to be willing to throw away the libraries you love (like standard std::net) and dive into the depths of assembly, eBPF, and custom consensus logic.

But when that Grafana dashboard finally levels out—when you see 100,000,000 writes per second flowing through your system with sub-100ms global convergence—the "magic" of distributed systems feels very, very real.

We are no longer just managing data. We are orchestrating the global heartbeat of the internet.


Engineering Checklist for High-Throughput Consensus:

  • Is your consensus protocol leaderless or hierarchical? (If not, you'll hit a CPU bottleneck).
  • Are you using a zero-copy data format? (FlatBuffers > Protobuf for 10M+ ops).
  • Have you bypassed the Linux Kernel? (DPDK/XDP is mandatory at this scale).
  • Is your gossip protocol "Interest-Aware"? (Avoid the $O(n^2)$ broadcast storm).
  • Are your data structures CRDT-compliant? (Ensure mathematical convergence).

If you’re building the next generation of global infrastructure, these aren't just "nice-to-haves." They are the survival kit for the 100M write-per-second world. Let's keep pushing the boundaries of what's possible. The speed of light is the only limit we have left.


More to explore

Keep diving in