One Shard Took Every New User — JavaScript Bug Hunt

Modelled on Foursquare, October 2010: the service was down for about eleven hours.

  • Language: JavaScript
  • Layer: Database
  • Difficulty: Hard
  • Concepts: Sharding, Scaling
  • Modelled on: Foursquare · 2010
  • Visible tests: shards are in range and stable; sequential users are spread evenly
  • Reward: 50 XP for a complete fix

Briefing

Modelled on Foursquare, October 2010: the service was down for about eleven hours. Its check-in data was sharded across MongoDB servers by user, the shards had filled unevenly, and one of them grew until its data no longer fit in memory — at which point it slowed to the speed of disk and took the site with it.

This reconstruction places users on shards by contiguous id ranges fixed when the cluster was built. User ids only grow, so every new user lands on the same shard.

Fix shardFor so placement is balanced for any set of user ids.

Bug report

BUG-4SQ-SHARD · Priority: Critical · Reported by: database operations

shardFor(userId, shardCount) must:

  • return an integer in [0, shardCount)
  • be stable: the same userId and shardCount always give the same shard
  • spread users evenly whatever their ids look like — sequential ids, ids far above the original ranges, and ids that share a common stride (our id service hands out ids in strides of the server count) must each put every shard within 10% of an equal share

Hash the id (hash.fnv1a32 over the decimal string is available); do not use ranges or the raw id modulo the shard count.

shardLoads(checkins, shardCount) counts check-ins per shard (a check-in follows its user) and fitsInRam(loads, capacity) must then hold for a cluster sized for the average load plus headroom.

Observed: every user created since launch is on the last shard.

Logs

[shards] shard3 working set above RAM, page faults climbing
[shards] shard0 18%, shard1 17%, shard2 18%, shard3 47% of all check-ins

The code as shipped

src/checkins/placement.js (editable)

var config = require("./config");

// Users were split into contiguous id ranges when the cluster was built.
exports.shardFor = function (userId, shardCount) {
  var shard = Math.floor(userId / config.RANGE_SIZE);
  return shard < shardCount ? shard : shardCount - 1;
};

exports.shardLoads = function (checkins, shardCount) {
  var loads = [];
  for (var s = 0; s < shardCount; s++) loads.push(0);
  for (var i = 0; i < checkins.length; i++) {
    loads[exports.shardFor(checkins[i].userId, shardCount)] += 1;
  }
  return loads;
};

exports.fitsInRam = function (loads, capacity) {
  for (var i = 0; i < loads.length; i++) if (loads[i] > capacity) return false;
  return true;
};

Read-only context: src/checkins/config.js, src/checkins/hash.js.

Open the hunt to edit the files, run the visible tests and submit against the hidden ones. More JavaScript bug hunts.