DEV Community

Saqib Ameen Subhan
Saqib Ameen Subhan

Posted on Edited on

I take down a healthy primary on purpose. You should too.

In one of my previous companies, I built an automation that used to scale our MongoDB clusters automatically. When CPU crossed a threshold in CloudWatch, the automation upgraded the nodes one by one. First a secondary, then the second secondary, and the primary last. It took around 12 minutes, nobody had to log in, and the application was not impacted.

The order matters here. Restarting the primary last means the cluster goes through only one election for every upscale, and the automation takes care of it. But this only works if you trust elections, and you will trust elections once you have triggered them yourself and seen how your application behaves when it happens.

What happens during an election?

In simple terms, this is how it works.

  1. Every node in a replica set checks on the other nodes every 2 seconds (heartbeats).
  2. If the primary does not respond for 10 seconds (electionTimeoutMillis), the other nodes assume it is gone.
  3. A secondary asks the other nodes to vote for it. It can only win if its data is at least as up to date as the nodes voting for it (i.e its oplog is as recent as theirs).
  4. The node that gets the majority of votes becomes the new primary, and the other secondaries start copying data from it.

One thing that surprises a lot of people is the priority setting. Priority decides which nodes get preference to become primary, but it cannot make a node with old data win. A node with priority 10 that is behind will lose to a node with priority 1 that is fully up to date. Up to date data matters more than priority. Once the higher priority node catches up with the others, it calls another election and takes over as primary, as the priority dictates.

In total, there can be up to 15 seconds with no primary, and during that time MongoDB cannot accept any writes.

The real problem is your application

The election itself only takes a few seconds. What really matters is what your application does during those seconds. I have seen three kinds of behaviour.

  • Users see errors. If the driver throws a NotWritablePrimary error and nobody has handled it, users start getting errors, and things like checkout start failing, which directly impacts the business.
  • The application retries blindly. Someone wrote their own retry loop, and sometimes the same write gets saved twice. This is actually worse than failing.
  • The driver handles it. With retryWrites=true and a sensible serverSelectionTimeoutMS, the driver holds the write, waits for the new primary, and retries it exactly once. The user just sees one slow request.

You don't want to find out which one you have during an incident.

The drill

I run this every quarter. First in staging, then in production during a low traffic window. Yes, production too. Staging tells you that the drill works. Production tells you how your real application behaves.

// 1. Check which node is primary right now
rs.status().members.map(m => ({ name: m.name, state: m.stateStr }))

// 2. Ask the primary to step down
rs.stepDown(60)   // it will not try to become primary again for 60 seconds

// 3. Check from a secondary which node became the new primary
rs.status()
Enter fullscreen mode Exit fullscreen mode

While the drill runs, I note down these numbers.

time taken to elect a new primary : ___ seconds
errors seen by the application    : ___
p99 latency during the election   : ___ ms
any write saved twice?            : ___ (check using a unique key)
Enter fullscreen mode Exit fullscreen mode

These numbers are the real output of the drill, so fill them in from your own run.

rs.stepDown() is the gentle way to do this. The primary finishes what it is doing and hands over cleanly. Once you are comfortable with that, try the harder version. Kill the primary's mongod process with kill -9, or cut its network. This is much closer to what happens in a real cloud outage, because the other nodes have to wait the full 10 seconds before they start an election.

Two things I learned doing this at scale

Always keep an odd number of voting members. With four voting members, the votes can split 2 and 2, and then nobody becomes primary. If you really cannot afford a third data node, you can use an arbiter, but remember that an arbiter does not help w:"majority" writes get acknowledged, because it does not hold any data.

Use priority to decide which region should be primary. In our multi region clusters, we gave higher priority to the nodes in the main region, so that once those nodes recovered, the primary moved back there automatically. But as I said above, priority only gives preference. The node with the most up to date data still wins.

Try it yourself

If you want to try this without touching a real cluster, mdbkit lab gives you a 3 node replica set on your laptop.

pip install mdbkit
mdbkit lab start          # 3 node replica set on 127.0.0.1:28110-28112
Enter fullscreen mode Exit fullscreen mode

Connect to the replica set and run a loop that saves an order every second, with the time printed next to it.

// mongosh "mongodb://127.0.0.1:28110,127.0.0.1:28111,127.0.0.1:28112/"
let i = 0
while (true) {
  const t = new Date().toISOString().substr(11, 8)
  try {
    db.drill.insertOne({ order: ++i })
    print(t + "  order " + i + " saved")
  } catch (e) {
    print(t + "  order failed, no primary available")
  }
  sleep(1000)
}
Enter fullscreen mode Exit fullscreen mode

In a second terminal, connect to the primary and run rs.stepDown(30). Then watch the first terminal and look at the timestamps. You will see exactly how long your orders were affected, and whether they failed or just waited.

When you are done:

mdbkit lab destroy --yes
Enter fullscreen mode Exit fullscreen mode

An election you have practised is just a normal operational event. An election you have never practised is an outage.

Top comments (0)