13 min read

The Impossible Triangle: How Uber Built a Globally Linearizable Database Without Sacrificing Latency

Uber: Building a Globally Linearizable Database With Low Latency

It’s 2:00 AM in San Francisco. A rider opens the app in Tokyo. Simultaneously, a driver in London accepts a trip, and a payment is processed in São Paulo. For Uber, the world never sleeps, and neither does its data. But here is the brutal truth that keeps distributed systems engineers up at night: You can have Consistency, Availability, and Partition Tolerance—but you can only pick two.

For years, Uber’s engineering teams grappled with this fundamental law of physics (and computer science). They had the scale—millions of writes per second—but they were drowning in the complexity of eventual consistency. When you’re moving money and people, "eventually" isn't good enough. If a rider cancels a trip in New York, a driver in New Jersey cannot see that trip as available three seconds later. That’s not a bug; that’s a lawsuit.

Enter Docstore, Uber’s new distributed database layer. It’s a system that claims to have solved one of the hardest problems in computer science: Global Linearizability at scale.

Today, we’re going to tear down the architecture of how Uber engineered their way out of the CAP theorem trap. We're talking atomic clocks, hybrid logical clocks, and a relentless pursuit of low latency.


The Hype vs. Reality: Why Everyone is Talking About This

If you follow cloud infrastructure news, you’ve likely heard rumblings about "NewSQL" and "Globally Distributed SQL." The hype cycle has been dominated by Google Spanner, CockroachDB, and Yugabyte. The promise? A database that looks like a single machine but spans the globe.

Why did this gain traction now? Because microservices shattered our data.

In the monolith era, you had one database. Life was simple. You had ACID transactions (Atomicity, Consistency, Isolation, Durability). When Uber moved to microservices, they sharded their data. To keep things fast, they adopted eventual consistency—write to a local node, and replicate asynchronously. It’s fast, but it breaks developer brains. You have to write complex application logic to handle conflicts, out-of-order events, and stale reads.

Uber’s Docstore isn't just another database; it’s an abstraction layer that sits on top of their storage engines. It provides the holy grail: Strict Serializability (the highest level of isolation) across geographically dispersed regions, without the penalty of waiting for the speed of light.

How did they do it? They didn't reinvent the storage engine. They reinvented Time.


The Core Problem: The Speed of Light is a Jerk

To understand the solution, we have to define the enemy.

In a multi-region setup, Region A (US-East) and Region B (Asia-Pacific) are separated by thousands of miles of fiber optic cable. A round-trip time (RTT) is roughly 150-200ms. This is physics. You cannot beat it.

If you want Linearizability (the illusion of a single, instantaneous copy of data), you need a total global order of events. If I write to the database in New York, and you read from the database in Tokyo, you must see that write.

Traditional consensus algorithms like Paxos or Raft require a majority of nodes to agree before a write is committed. If you have nodes in the US, Europe, and Asia, every write requires a cross-ocean conversation. This is the "Latency Tax."

For years, the industry standard was: If you want global consistency, pay 200ms per write. Uber said: No.


The Architecture: How Docstore Works

Docstore is built on a foundation of three key pillars:

  1. Hybrid Logical Clocks (HLC).
  2. A Layered Consensus Model.
  3. Regional Autonomy with Global Ordering.

Let’s break these down.

1. The Time Lords: Hybrid Logical Clocks (HLC)

The bedrock of Uber’s linearizability is time. But not the time on your wall clock. Physical clocks drift. NTP (Network Time Protocol) can be off by milliseconds, which is an eternity in a database.

Uber implemented Hybrid Logical Clocks. HLC combines the best of physical time (wall clocks) and logical time (Lamport timestamps).

  • Physical Time: Gives us a rough ordering.
  • Logical Time: Ensures causality even if physical clocks drift.

In Docstore, every transaction gets an HLC timestamp. This timestamp is a tuple: (physical_time, logical_counter).

Here is the magic: HLC allows Docstore to assign a globally unique, monotonically increasing timestamp to every transaction without requiring a round-trip to a central clock server. This means a node in Tokyo can assign a timestamp that is guaranteed to be consistent with a node in Virginia.

// Simplified HLC Structure in Docstore
message HybridLogicalClock {
  int64 physical_time = 1; // Nanoseconds since epoch
  int32 logical_counter = 2; // Counter for events within the same nanosecond
}

If two events happen at the exact same physical nanosecond, the logical counter breaks the tie. This gives us a total order.

2. The Layered Consensus: Separating Metadata from Data

This is where Uber’s architecture gets clever.

In a naive Spanner-like system, you replicate everything everywhere. That’s expensive and slow. Uber realized that their data has different access patterns. Some data is "hot" (needs to be fast), some is "cold" (archival).

Docstore separates the Metadata Plane from the Data Plane.

  • The Metadata Plane: This is where the global consensus happens. It’s a Paxos group that manages the mapping of keys to shards. It handles the "who owns what" and "what time is it" questions.
  • The Data Plane: This is where the actual bytes live. It’s sharded and replicated, but it doesn't need to talk to every region for every read.

The Workflow:

  1. A client wants to write data for user_id: 123.
  2. It asks the Metadata Plane: "Who owns this key, and what is the latest timestamp?"
  3. The Metadata Plane uses its Paxos group to ensure the client gets the latest timestamp (Read-Latest).
  4. The client then goes to the Data Plane (a specific shard in a specific region) and performs the write.

By offloading the heavy lifting of ordering to a lightweight metadata layer, the data layer can operate with high throughput.

3. Regional Autonomy: The "Local Read" Trick

The most impressive part of Docstore is how it handles reads.

If you are in Tokyo and you want to read a record, you don't want to wait for a round-trip to the US. Docstore uses Leaseholder-based Reads.

Each shard has a "Leader" in one region, but it also has "Followers" in other regions. The Leader holds a Lease—a promise from the Metadata plane that it won't change the data for X milliseconds.

If a client in Tokyo wants a strictly consistent read, it can read from the local Follower if the Follower has caught up to the Leader's lease timestamp.

Here’s the code logic for how Docstore ensures linearizability on a local read:

def read_linearizable(key):
    # 1. Get the latest timestamp from the metadata plane
    # This is the "Read-Latest" operation. It's fast because it's a metadata query.
    latest_ts = metadata_plane.get_latest_timestamp(key)
    
    # 2. Find the local replica for the key
    replica = data_plane.get_local_replica(key)
    
    # 3. Wait until the local replica has caught up to the latest_ts
    # This is usually sub-millisecond if replication is healthy
    replica.wait_until_synced(latest_ts)
    
    # 4. Perform the read
    return replica.read(key, latest_ts)

This is the secret sauce. The wait in step 3 is usually negligible because Uber’s network is optimized. But by doing this, they guarantee that the read is linearizable. You will never read stale data. You will never see a "time travel" anomaly.


The Secret Weapon: TrueTime vs. HLC

You might be asking: "Wait, isn't this just Google Spanner?"

Spanner uses TrueTime, which relies on GPS antennas and atomic clocks in every datacenter. TrueTime gives Spanner a bounded uncertainty interval (e.g., "it is currently 10:00:00 ± 5ms"). Spanner then waits out the uncertainty (the "commit-wait" phase) to ensure consistency.

Uber took a different path. They couldn't rely on atomic clocks in every edge location. So, they built Docstore using Hybrid Logical Clocks and a Centralized Timestamp Oracle (TSO) .

Wait, a centralized TSO? Isn't that a single point of failure?

No. Uber’s TSO is sharded. It’s a Paxos group itself. It doesn't get queried for every write; it gets queried for timestamp ranges.

Here’s the flow:

  1. A Region (e.g., Asia) requests a batch of timestamps from the TSO.
  2. The TSO grants a range: [1000, 2000].
  3. The Asia region can now assign timestamps 1001, 1002... locally without talking to the US.
  4. This is called Timestamp Batching.

This drastically reduces the latency. The region only talks to the global TSO when it runs out of timestamps. This allows Uber to achieve single-digit millisecond writes across continents, something that was previously thought impossible for strict serializability.


The Engineering Curiosity: The "Uncertainty" Problem

There is a catch. HLC relies on physical clocks. If a clock in Tokyo drifts backwards, it could cause a consistency violation.

Uber implemented a Clock Skew Detection mechanism.

If a node detects that its clock is drifting too far from the TSO, it "fences" itself. It stops serving reads and writes until it can resync. This is a safety mechanism. In a distributed system, a node with a bad clock is worse than a dead node.

They also use a technique called Epoch-based Replication. The data is versioned by "Epochs." If a node is fenced, it can't just rejoin with old data. It must catch up to the latest epoch.

Let's look at the pseudo-code for the Fencing Logic:

func (n *Node) checkClockSkew() error {
    // Get the current time from the central TSO
    tsoTime := tso.GetTime()
    
    // Calculate the drift
    drift := tsoTime - n.localTime
    
    // If drift exceeds the threshold (e.g., 10ms), we are unsafe
    if drift > maxDrift {
        n.setState(FENCED)
        return errors.New("clock skew too high, node fenced")
    }
    
    return nil
}

This is the kind of gritty engineering that makes distributed systems work. It’s not just about the happy path; it’s about handling the chaos of reality.


The Scale: What Are We Actually Talking About?

Uber’s scale is mind-boggling. To understand why Docstore is a marvel, let’s look at the numbers (based on public engineering disclosures and typical Uber scale):

  • Millions of writes per second: Uber processes a massive amount of location data, trip states, and payment info.
  • Petabytes of data: The dataset is massive.
  • Global Footprint: Datacenters in multiple regions (US, EU, APAC).
  • Latency SLA: P99 latency for writes is often under 100ms, and reads are often under 20ms.

If you tried to run this on a traditional RDBMS with synchronous replication, the throughput would tank. If you ran it on Cassandra with eventual consistency, you’d have to write complex conflict resolution logic in the application layer, leading to bugs and data inconsistencies.

Docstore provides a single API that feels like a local database but is globally consistent.


The "Why": The Business Case for Linearizability

Okay, so it’s technically cool. But why did Uber spend millions of dollars and years of engineering effort on this?

Because "Eventual Consistency" is expensive.

Here is the hidden cost of eventual consistency:

  1. Developer Complexity: Every team has to write code to handle conflicts. "What if the driver status is 'Offline' in one region but 'Online' in another?"
  2. Data Corruption: Without a global order, it’s easy to lose writes. If two updates happen simultaneously, one might overwrite the other.
  3. Customer Trust: If a rider sees a surge price of $50 in one region and $20 in another, they will lose trust.

Linearizability simplifies the application layer. Developers can write code as if they are writing to a single machine. They don't have to worry about "read-your-writes" consistency (the guarantee that you see your own writes). It just works.

This allows Uber to move faster. It allows them to launch new features (like Uber Reserve or Uber Eats) without worrying about the underlying data infrastructure.


The Trade-offs: What Did They Sacrifice?

No architecture is free.

To achieve global linearizability, Docstore sacrifices availability during network partitions.

If the network connection between the US and EU breaks, and you need a quorum to write, the system will block writes to the minority partition. This is the "C" in CAP theorem. Uber chose Consistency over Availability during a partition.

However, they mitigated this by:

  1. Regional Autonomy: If the partition is within a region, the region can continue to serve traffic as long as it has a majority of nodes.
  2. Read-Only Mode: During a partition, the minority partition can serve reads if the data is not stale, but it cannot accept writes.

This is a conscious trade-off. For Uber’s core use cases (payments, trip state), consistency is more important than availability. If a trip status can't be updated, the driver can wait 200ms. But if the trip status is wrong, the user is charged incorrectly.


The Implementation Details: A Peek Under the Hood

Let’s get into the weeds of how Docstore implements the Two-Phase Commit (2PC) protocol with HLC.

When a transaction spans multiple shards (e.g., updating a user's balance and a trip's status), Docstore uses a variant of 2PC.

  1. Prepare Phase: The coordinator asks all shards to prepare the transaction. Each shard locks the resources and returns the HLC timestamp it can commit.
  2. Commit Phase: The coordinator collects all timestamps, picks the maximum one, and tells all shards to commit with that specific timestamp.

By using the maximum timestamp, Docstore ensures that the transaction is atomic and globally ordered.

Here is a simplified code snippet of the Commit Phase:

// Coordinator logic
List<PrepareResponse> responses = shards.prepare(txn);

// Find the max timestamp from all shards
long maxTimestamp = responses.stream()
    .mapToLong(PrepareResponse::getTimestamp)
    .max()
    .orElseThrow();

// Commit with the max timestamp to ensure global order
for (Shard shard : shards) {
    shard.commit(txn, maxTimestamp);
}

This ensures that even if shard A is in the US and shard B is in Asia, the transaction is committed at the same logical time in both places.


The Future: What’s Next for Docstore?

Uber is not stopping here. The engineering team is working on:

  • Geo-Partitioning: Allowing data to be pinned to specific regions for data sovereignty (e.g., GDPR). This means a European user's data stays in Europe, but they still get the benefits of global consistency.
  • Serverless Integration: Making Docstore the backend for Uber’s serverless functions, so developers can write stateless code that reads/writes to a globally consistent store.
  • AI/ML Workloads: Using the consistent snapshot capability of Docstore to train ML models on a globally consistent view of data.

The Takeaway: Time is the Final Frontier

The story of Uber’s Docstore is a story of mastering Time.

For decades, distributed systems engineers tried to solve consistency by throwing hardware at the problem—more memory, faster disks, bigger networks. Uber solved it by realizing that the problem wasn't the data; it was the order of the data.

By implementing Hybrid Logical Clocks, Timestamp Batching, and Regional Autonomy, they built a system that provides the illusion of a single machine across a planet.

It’s a testament to the fact that in the world of distributed systems, you don't always have to choose between speed and correctness. Sometimes, you just have to be clever about how you define "now."

If you’re building systems at scale, the lesson from Uber is clear: Don’t ignore the clock. Master it.


This deep dive was inspired by the brilliant work of the Uber Engineering team. If you want to geek out on the nitty-gritty details, I highly recommend checking out their official engineering blog and the papers on Hybrid Logical Clocks.


More to explore

Keep diving in