All resources
Daily Learning · Distributed Systems: Principles and Paradigms

When Your Failover Plan Becomes the Problem

Thursday, 23 July 2026
When Your Failover Plan Becomes the Problem

🎯 After reading this you will understand exactly why failover scripts in sharded databases can themselves cause outages and how fencing tokens plus time-bucketed keys prevent serving stale data.

Step 1: The Starting Setup

Twitter stored timelines in many MySQL shards so write latency stayed low even as the firehose of tweets grew. Each shard held only a slice of user IDs. For the obvious failure case they added a second replica set in another rack.

Primary Shard (user slice A)
        │  replication stream
        ▼
Replica Set (user slice A)

What is a shard?

A shard is one small piece of a big database. Instead of one giant server holding every user, you split the users across many smaller servers.

What is a replica set?

A replica set is a primary database plus one or more copies (replicas) that stay in sync. If the primary dies, a replica can take over.

Step 2: What Happens When a Rack Dies

One day a whole rack lost power. The failover script tried to promote the replica in the other rack to become the new primary.

The script assumed the old primary would cleanly release its locks. In reality the replication stream simply stopped moving.

What is a replication stream?

The replication stream is the continuous flow of changes (new tweets, likes, etc.) sent from primary to replica so the replica stays up to date.

Because the stream backed up, the new primary had not yet received the latest writes. It started serving old data to clients.

Old Primary (dead)
        X  (no more changes)
Replica (now promoted)
        │
        ▼
Clients ← stale timeline data
The failover path itself had become the single point of failure.

Step 3: Adding Fencing Tokens

Engineers introduced fencing tokens. A fencing token is a unique increasing number given to a new primary. Any old primary that later wakes up sees a higher token and refuses to accept writes.

Think of it like a restaurant kitchen: only the cook holding the current “order number token” is allowed to send plates out. If an old cook shows up with an earlier token, the waiters ignore them.

Step 4: Changing the Partition Key

They also changed how data was split across shards. The new primary key became:

PRIMARY KEY ((user_id, bucket), tweet_id)

What is a partition key?

The partition key decides which shard a row lives on. Previously it was just user_id. Now it is the pair (user_id, bucket).

A bucket is a time window (for example, one hour). Every tweet for a user is placed into the bucket that matches its timestamp.

Because a single shard now only covers one user’s tweets inside one time window, losing that shard cannot wipe out an entire timeline.

Step 5: The New Trade-offs

Every client request must now include the correct bucket value. When reading a full timeline the system must query several shards (one per bucket) instead of one. This increase in parallel requests is called query fanout.

To keep the 99th-percentile latency acceptable they added an extra caching layer in front of the shards.

Client request
    │
    ▼
Cache layer
    │
    ├── Shard (user, bucket-1)
    ├── Shard (user, bucket-2)
    └── Shard (user, bucket-3)

They accepted the extra moving parts because losing an entire timeline during a major sports event was no longer acceptable.


The core lessons

  • Failover scripts are code too; they can fail exactly when you need them most.
  • A replication stream that stops moving silently turns a replica into a source of stale data.
  • Fencing tokens give every new primary an unmistakable “I am current” signal that old primaries cannot fake.
  • Adding a time bucket to the partition key limits the blast radius of any single shard failure.
  • Extra query fanout and caching are the price of that safety; measure them before you need them.

Want this kind of thinking applied to your business?

Book a free 30-minute discovery call — we’ll show you your highest-value first automation, no jargon, no obligation.

Book a discovery call