14 min read

The Edge is Burning: Designing Hierarchical LSM-Trees for Write-Intensive Workloads in Geographically Distributed Key-Value Stores

Hierarchical LSM-Trees for Write-Intensive Geo-Distributed Storage

Imagine this: You’re sitting in San Francisco, running a globally distributed database. A user in Tokyo just swiped their credit card, a sensor in Frankfurt reported a temperature spike, and a logistics truck in São Paulo updated its GPS coordinates. All of these events are writing to your system simultaneously.

In a single-datacenter world, your storage engine handles this with the grace of a well-oiled machine. You have a Log-Structured Merge-Tree (LSM-Tree) that buffers writes in memory, flushes them to disk sequentially, and compacts them in the background. Life is good.

But the moment you distribute that database across the globe—spanning continents connected by flaky subsea cables—that elegant LSM-Tree becomes a bureaucratic nightmare. You introduce latency that dwarfs disk I/O, you face the specter of conflicting writes, and you realize that the bottleneck isn't your SSD; it's the speed of light.

Welcome to the wild world of designing Hierarchical LSM-Trees for Write-Intensive Workloads in Geographically Distributed Key-Value Stores. If you’ve ever wondered how systems like CockroachDB, TiKV, or YugabyteDB manage to keep their sanity while replicating petabytes of data across oceans, you’re in the right place.

Today, we’re going to tear down the architecture, look at the messy reality of consensus protocols, and build a mental model for how to structure storage engines that don’t collapse under the weight of global write traffic.


The Physics Problem: Why Geo-Distribution Breaks Traditional LSMs

First, let’s establish the baseline. The LSM-Tree is the darling of write-intensive workloads. It turns random writes into sequential writes. It uses a MemTable (in-memory) and a SSTable (Sorted String Table on disk). It’s fantastic.

But traditional LSM-Trees were designed for a single node or a single datacenter. When you stretch them across the globe, three fundamental constraints rear their ugly heads:

  1. The Speed of Light (Latency): A round trip from Virginia to Tokyo takes roughly 150ms. If your write path requires synchronous replication to a quorum of nodes across these regions, every write takes at least 150ms. For a write-intensive workload, that is an eternity.
  2. The CAP Theorem (Consistency vs. Availability): You can’t have perfect consistency and perfect availability during a network partition. In a geo-distributed setup, partitions are the norm, not the exception.
  3. Conflict Resolution: If two users write to the same key in different regions simultaneously, who wins? Last-write-wins is a recipe for data loss. You need a mechanism to order these writes deterministically.

The naive approach is to simply replicate the entire LSM-Tree state to every region. This is what we call Full Replication. It’s simple, but it’s a disaster for write-heavy workloads. Why? Because every write must be applied everywhere. If you have 10 regions, your write amplification just multiplied by 10. Your compaction overhead explodes. Your network bandwidth costs skyrocket.

We need a Hierarchical approach. We need to think about data locality.


The Architectural Shift: Hierarchical LSM-Trees

The core idea of a Hierarchical LSM-Tree in a geo-distributed context is to decouple the durability of the write from the visibility of the write across the globe.

Instead of a monolithic tree replicated everywhere, we build a tree of trees. We organize the LSM structure into layers based on geographic proximity.

The Three Tiers of Storage

Let’s conceptualize the hierarchy:

  1. Tier 1: The Local Edge (The Speedy Buffer)

    • Located in the region where the write originates.
    • Contains a MemTable and a small, fast local SSTable.
    • Goal: Ingest writes with sub-millisecond latency.
    • Consistency: Local quorum (e.g., 3 replicas within the same datacenter).
  2. Tier 2: The Regional Aggregator (The Consolidator)

    • A central node or cluster within a geographic region (e.g., "US-East" aggregates US-East-1, US-East-2, etc.).
    • Receives asynchronous streams of data from Tier 1 nodes.
    • Goal: Reduce network traffic by batching and merging writes before sending them globally.
    • Consistency: Regional quorum.
  3. Tier 3: The Global Backbone (The Source of Truth)

    • The central repository, possibly spanning multiple continents.
    • Stores the "cold" data or the authoritative state.
    • Goal: Durability and global consistency for reads that traverse regions.

This hierarchy allows us to batch writes. Instead of sending every single INSERT to Tokyo, we send a compressed batch of updates every 100ms. This amortizes the latency cost.


The Write Path: A Journey Through the Hierarchy

Let’s trace a write from a user in London to the global backbone.

Step 1: The Local Ingest (MemTable)

The write hits a node in London. It goes straight into the MemTable (a concurrent skip list or B-tree in memory). This is fast. It’s local. The client gets an acknowledgment almost immediately.

But wait, we need durability. If the London node dies, we lose the write. So, we replicate the MemTable to a Local Quorum (e.g., two other nodes in the London datacenter). This is still fast—sub-millisecond.

The key insight: We do not wait for the write to reach Tokyo before acknowledging. We wait for the local quorum. This is the Hierarchical part of the LSM.

Step 2: The Flush and the Log

The MemTable fills up. It flushes to a Level-0 (L0) SSTable on local disk. Simultaneously, we write to a Write-Ahead Log (WAL). This WAL is not just for local recovery; it is the stream that will be sent to the next tier.

Here is where the magic happens: The WAL is not a random stream of bytes. It is a structured log of key-value operations. We can compress it, batch it, and ship it.

Step 3: The Regional Aggregation

The London node sends its WAL to the Regional Aggregator (let’s say, a node in Frankfurt). Frankfurt is responsible for the European region.

Frankfurt receives WALs from London, Paris, Berlin, and Madrid. It does something clever: it merges these streams.

Imagine London writes key=A, value=1 and Paris writes key=A, value=2. Frankfurt sees both. If we are using a Last-Write-Wins (LWW) strategy with timestamps, Frankfurt can resolve this conflict locally. It only needs to send the winner to the global backbone.

This reduces the write amplification. Instead of 4 writes traveling to Tokyo, only 1 travels to Tokyo.

Step 4: The Global Backbone (The Consensus Layer)

Now, Frankfurt needs to replicate this to the global backbone. This is where we hit the Consensus wall.

We cannot use a simple async replication to the backbone if we want strong consistency. We need a consensus protocol like Raft or Paxos. But running Raft across continents is expensive. Every write requires a round trip to a majority of the backbone nodes.

The Trick: We don't run consensus on every write. We run consensus on batches.

Frankfurt accumulates writes for 50ms. It then proposes a single batch to the global Raft group. The Raft group replicates this batch to a majority of backbone nodes (e.g., Tokyo, Virginia, Sydney). Once committed, the batch is applied to the Global LSM-Tree.


The Read Path: Navigating the Labyrinth

Reads are trickier. If a user in London reads key=A, they want the latest value. But the latest value might be in the global backbone, or it might be in Frankfurt, or it might be in the London L0 SSTable.

We need a Read Path that doesn't always go to the global backbone.

1. Local Read (Hot Data)

The read hits the London node. It checks the MemTable. If it finds the key, it returns. This is the fastest path.

2. Local SSTable Read (Warm Data)

If not in the MemTable, it checks the local L0 SSTable. If found, return.

3. Regional Read (Consolidated Data)

If not found locally, it queries the Regional Aggregator (Frankfurt). Frankfurt has a consolidated view of the region. It has merged all the local WALs. It can answer the read without going to the global backbone.

4. Global Read (Cold Data)

If the data is not in the region, we go to the global backbone. This is the slowest path, but it’s rare for write-intensive workloads where data is usually "hot" for a short period.

The Hierarchy of Reads:

  • MemTable: Nanoseconds.
  • L0 SSTable: Microseconds.
  • Regional Aggregator: Milliseconds.
  • Global Backbone: 100+ Milliseconds.

By structuring the LSM hierarchically, we ensure that 99% of reads are served by the first three tiers.


The Compaction Nightmare (and How to Tame It)

Compaction is the bane of every LSM-Tree engineer. In a geo-distributed system, compaction is a network problem, not just a disk problem.

If we naively compact the global LSM-Tree, we might rewrite data that hasn't changed. But in a hierarchical setup, compaction happens at each level.

Local Compaction

At the edge, we compact L0 SSTables into L1. This is local disk I/O. It’s fast.

Regional Compaction

The Regional Aggregator compacts the merged WALs into a regional SSTable. This is where we can drop data that has been superseded. If London wrote key=A, value=1 and Paris wrote key=A, value=2, and Paris's write is newer, we can drop London's write. This reduces the data size before sending it to the global backbone.

Global Compaction

The global backbone compacts the regional batches into the global LSM. This is the most expensive operation. We need to be smart about it.

Optimization: We use Tiered Compaction at the global level instead of Leveled Compaction. Tiered compaction is more write-efficient but has higher read amplification. Since global reads are rare, we trade read performance for write throughput.


Handling Conflicts: The CRDT Approach

In a geo-distributed system, conflicts are inevitable. Two regions might write to the same key at the same time.

We have a few options:

  1. Last-Write-Wins (LWW): Simple, but relies on synchronized clocks. In a geo-distributed system, clocks are never perfectly synchronized. You can use Hybrid Logical Clocks (HLC) to mitigate this, but it's still a heuristic.
  2. Multi-Value (MV): Store all conflicting values. The application resolves them. This is safe but pushes complexity to the app.
  3. CRDTs (Conflict-Free Replicated Data Types): The holy grail. If your data type is a CRDT (e.g., a counter, a set, a map), you can merge writes automatically without conflicts.

For a write-intensive key-value store, HLC + LWW is the pragmatic choice. But if you need strong consistency, you need to run consensus on the write path.

The Hybrid Approach: Use Raft for the global backbone (strong consistency) and CRDTs for the regional aggregation (eventual consistency). This gives you the best of both worlds: fast local writes and strong global consistency.


Code Snippet: The Hierarchical Write Path

Let's look at a simplified pseudocode for the write path.

class HierarchicalLSM:
    def __init__(self, region):
        self.memtable = MemTable()
        self.local_sstables = []  # L0
        self.regional_aggregator = RegionalAggregator(region)
        self.global_backbone = GlobalBackbone()

    def write(self, key, value):
        # Step 1: Write to local MemTable
        timestamp = HybridLogicalClock.now()
        self.memtable.put(key, value, timestamp)
        
        # Step 2: Replicate to local quorum (e.g., 3 nodes)
        local_quorum.replicate(key, value, timestamp)
        
        # Step 3: Check if MemTable is full
        if self.memtable.is_full():
            self.flush_to_local_sstable()
        
        # Step 4: Asynchronously send to Regional Aggregator
        self.regional_aggregator.send_async(key, value, timestamp)

    def flush_to_local_sstable(self):
        # Flush MemTable to L0 SSTable
        sstable = SSTable.from_memtable(self.memtable)
        self.local_sstables.append(sstable)
        self.memtable.clear()
        
        # Trigger local compaction if needed
        if len(self.local_sstables) > THRESHOLD:
            self.compact_local()

    def compact_local(self):
        # Merge L0 SSTables into L1
        merged = merge_sstables(self.local_sstables)
        self.local_sstables = [merged]

And here is the Regional Aggregator:

class RegionalAggregator:
    def __init__(self, region):
        self.region = region
        self.buffer = []
        self.global_backbone = GlobalBackbone()

    def receive(self, key, value, timestamp):
        self.buffer.append((key, value, timestamp))
        
        # Batch and send to global backbone every 50ms
        if time.now() - self.last_flush > 50ms:
            self.flush_to_global()

    def flush_to_global(self):
        # Merge buffer and resolve conflicts (LWW)
        merged = resolve_conflicts(self.buffer)
        
        # Send batch to global backbone via Raft
        self.global_backbone.propose_batch(merged)
        self.buffer.clear()
        self.last_flush = time.now()

This is a simplified view, but it captures the essence: local writes are fast, global writes are batched.


The Infrastructure Reality: Compute and Storage Scale

Let’s talk numbers. What does it take to run this at scale?

Imagine you have 100 million writes per second globally. That’s a lot of writes.

  • Edge Nodes: You need thousands of edge nodes. Each node handles a shard of the key space. They are optimized for CPU and memory (for the MemTable).
  • Regional Aggregators: You need dozens of these. They need high network bandwidth (to receive WALs) and fast CPUs (to merge and resolve conflicts).
  • Global Backbone: You need a small number of high-capacity nodes (e.g., 5-7 nodes for Raft). They need massive storage (NVMe SSDs) and high network bandwidth (to replicate across continents).

The Cost: The network egress costs for geo-replication are brutal. If you replicate 1PB of data across 3 continents, you’re paying for 3PB of egress. This is why compression and batching are not optional; they are survival mechanisms.

Optimization: Use erasure coding at the global backbone instead of full replication. This reduces storage overhead but adds CPU overhead for encoding/decoding. It’s a trade-off.


The Hype Cycle: Why This Matters Now

You might have noticed a lot of buzz around "Geo-Distributed SQL" and "Global Databases". Companies like CockroachDB, YugabyteDB, and Google Spanner have popularized the idea that you can have a single database that spans the globe.

The hype is real, but the technical substance is often misunderstood. People think it’s magic. It’s not. It’s a carefully engineered hierarchy of storage and consensus.

The recent interest in Edge Computing has also fueled this. As more compute moves to the edge (e.g., 5G, IoT), the need for a write-intensive, geo-distributed key-value store becomes critical. You can’t have an autonomous car waiting 200ms for a write to replicate to a central datacenter.

The Shift: We are moving from "Cloud-Native" to "Edge-Native". The hierarchical LSM-Tree is the storage engine that makes this possible.


Lessons Learned: The Hard Truths

After building and operating these systems, here are the hard truths:

  1. Latency is Physics: You cannot beat the speed of light. You must design your system to hide latency, not fight it.
  2. Consensus is Expensive: Running Raft across continents is slow. Do it sparingly. Batch your writes.
  3. Conflicts are Inevitable: Don’t pretend they don’t exist. Use HLCs, CRDTs, or a hybrid approach.
  4. Compaction is a Network Problem: In a geo-distributed system, compaction is not just about disk I/O. It’s about minimizing cross-region traffic.
  5. Observability is Key: You need to monitor the lag between regions, the size of the WAL, and the compaction backlog. If you can’t see it, you can’t fix it.

The Future: What’s Next?

The future of hierarchical LSM-Trees lies in adaptive intelligence.

Imagine a system that automatically adjusts the hierarchy based on workload patterns. If a region suddenly becomes write-heavy, it spins up more edge nodes. If a key becomes hot globally, it replicates it to all regions.

We are also seeing the rise of Disaggregated Storage. Instead of storing SSTables on local disk, we store them in a shared object store (like S3). This separates compute from storage, allowing for elastic scaling. But it introduces new challenges: how do you compact data in S3 without moving it? How do you cache hot data?

These are the questions the next generation of storage engineers will answer.


Wrapping Up

Designing hierarchical LSM-Trees for write-intensive, geo-distributed workloads is one of the most challenging problems in distributed systems today. It requires a deep understanding of storage engines, consensus protocols, and network physics.

But when you get it right, it’s beautiful. You get a system that is fast, resilient, and globally consistent. You get a system that can handle the firehose of writes from the edge of the network while maintaining a coherent view of the world.

So, the next time you swipe your credit card in Tokyo and see the transaction appear instantly on your phone in San Francisco, remember the hierarchy. Remember the MemTables, the WALs, and the Raft groups. It’s not magic. It’s engineering.

And it’s a damn good time to be a storage engineer.


Further Reading:

  • The Log-Structured Merge-Tree (LSM-Tree) by O’Neil et al.
  • Raft Consensus Algorithm by Diego Ongaro.
  • Spanner: Google’s Globally-Distributed Database by Corbett et al.
  • CRDTs for Geo-Distributed Systems by Shapiro et al.

Now go build something that scales. The edge is waiting.


More to explore

Keep diving in