Running a distributed system sounds great, until one node starts

Arpit Bhayani

Arpit Bhayani

Aug 06, 2025 • 2 min read


Running a distributed system sounds great, until one node starts doing all the heavy lifting. That’s what hot shards are. But why do they occur?

Hot shards occur when a specific node of a distributed database becomes overloaded and has to work significantly harder than other parts to handle the workload, leading to slowdowns and bottlenecks. Here are some common reasons for it

When I used to manage a production ElasticSearch cluster on raw EC2 virtual machines, hot shards were a big pain, and here are the reasons that were the root cause 80% of the time

  • we did not follow the best practices
  • we wrote highly inefficient queries
  • we chose inefficient routing (partitioning) key
  • unexpected surges happened due to performance marketing
  • massive stop-the-world GC pauses, accumulating the backlog

Key things we did to address this were to have high observability in place, along with spending time knowing the internals and best practices. This looked counter-productive at first, but it turned out to be the best decision in the long run.

Note: although I mentioned these points w.r.t ElasticSearch, they hold true for almost all distributed databases that exist because the underlying limitations of hardware and software remain the same.

btw, enrollments are open for my august sys design cohort (starts this weekend 9th August, ~3 seats left), filled with no-fluff and highly practical engineering discussions aimed at making you a better engineer - arpitbhayani.me/course

Arpit Bhayani

Principal Engineer II at Razorpay - building Agent Studio, Ex-staff engg at GCP Memorystore & Dataproc, Creator of DiceDB, ex-Amazon Fast Data, ex-Director of Engg. SRE and Data Engineering at Unacademy. I spark engineering curiosity through my no-fluff engineering videos on YouTube and my courses