Skip to content
All articles

Sharding · System Design

A good hash can't fix a hot tenant

Hash partitioning on tenant_id spreads tenants evenly, but not load. When one merchant is 40% of writes, one shard sets your capacity. Three repairs, and what each one costs.

· 4 min read

Bar chart of write share across 8 shards partitioned by a hash of the tenant ID: shard 2 carries 49% of all writes, 3.9 times the average, against an even share of 12.5%; the other shards carry between 2% and 11.6%. One large merchant is 40% of all writes, so 8 machines deliver roughly the throughput of 2.

You avoided the monotonic key. You hashed tenant_id, as the documentation recommends for a multi-tenant system, and that choice is right about a lot: every tenant's rows live together, so a per-tenant query hits exactly one shard, and tenants spread evenly across the cluster.

Then the marketplace signs one large merchant, and that merchant is 40% of all writes.

What the shards see

The course's lab runs 100,000 writes from 149 tenants across 8 shards partitioned by hash(tenant_id):

Shard01234567
Share of writes2.0%11.2%49.0%7.2%7.6%7.4%11.6%4.0%
Load vs. average0.16×0.90×3.92×0.58×0.61×0.59×0.93×0.32×

Shard 2 carries almost four times the average load while shard 0 sits nearly idle. Capacity is set by the hottest shard, so eight machines deliver roughly the throughput of two.

Hashing distributes keys uniformly. It cannot distribute load uniformly, because load per key is not uniform. The skew is in the traffic, and no better hash function will remove it.

Three repairs, not interchangeable

Salting writes the large tenant's rows under tenant#0 to tenant#15, chosen at random for each write. Sixteen keys spread the load over up to sixteen shards, which here means all eight. The cost lands on reads: every read of that tenant becomes sixteen reads.

A composite key spreads on something the query already knows. Partitioning the large tenant by tenant_id#order_seq evens the load out to 1.01× in the lab, and every query for that tenant now fans out to all eight shards. It is the better choice when queries can supply the extra part of the key.

An explicit placement override gives the large tenant its own shard and leaves the other 148 alone. That shard still carries 40% of all writes by itself, so it needs a bigger machine or a further split. It is the only repair that needs no change to code on the read path.

Think Like an Engineer

The large tenant is salted across sixteen keys, and a report needs that tenant's orders for one day. How many shards does it ask, and what sets its latency? All eight, and the client merges the results, so the report is as slow as the slowest shard's reply. Salting bought an even write load with read fan-out.

You have probably met this problem under other names: "one customer's reports time out and nobody else's do", or "we had to give the enterprise account its own instance". The second one is a hot partition solved by hand.

Get one diagram a week

A short article built around one engineering diagram, from the same library as these courses.

One diagram-led article a week on AI and systems engineering. We email you once to confirm, and every newsletter has an unsubscribe link. Privacy policy