Skip to main content
Create your own
Lesson illustration

Sharding Relational Databases: Implementation & Complexities

Welcome to our next session on scaling databases. In the previous lesson, we successfully addressed the challenge of read-heavy workloads by implementing read replicas. This pattern allows us to scale read capacity horizontally, but it leaves a critical bottleneck untouched: all write operations still go to a single primary database. Eventually, this primary server will hit its limits in terms of storage, I/O, or CPU, no matter how much you scale it vertically.

Today, we tackle this fundamental limitation head-on. We will explore horizontal partitioning, more commonly known as sharding. This lesson will equip you to apply this powerful technique to a relational database, breaking a large dataset into smaller, more manageable pieces called shards. We'll analyze the critical decisions involved, particularly the choice of a sharding strategy, and dissect the significant complexities that arise, such as hotspots, cross-shard queries, and maintaining data consistency. Mastering sharding is a crucial step in your journey from managing small-scale systems to architecting for massive scale.

1. The Need for Sharding: Beyond Vertical Scaling

When a single database can no longer handle the write load or store the sheer volume of data, the first instinct is often vertical scaling—upgrading to a more powerful server. However, this approach has its limits. There's a finite amount of CPU, RAM, and storage you can pack into a single machine, and the cost increases exponentially. When you hit this wall, the only way forward is to scale horizontally.

The following video provides an excellent explanation of why and when sharding becomes necessary.

Sharding in System Design Interviews w/ Meta Staff Engineer

This video from Hello Interview, featuring a Meta Staff Engineer, provides a clear, practical motivation for sharding.

Watch the segment from the beginning, which walks through the limitations of a single large database instance and introduces sharding as the solution for scaling storage, read, and write throughput by distributing data across multiple machines.

As the video explains, sharding is the process of splitting your data across multiple database instances. Each instance, or shard, holds a subset of the data and operates as a standalone database with its own resources. Together, the shards form a single logical database.

This diagram provides a great roadmap for the concepts we'll discuss today.

This diagram outlines the key components of database sharding, including its definition, different types (Range-based, Directory-based, Key-based), criteria for selecting a shard key, and patterns for routing requests.

2. The Core Decisions: Shard Key and Distribution Strategy

Successfully implementing sharding hinges on two fundamental decisions:

  1. What to shard by? (The Shard Key)
  2. How to distribute data based on that key? (The Distribution Strategy)

Let's break these down.

The Shard Key: The Most Critical Choice

The shard key is the column or field you use to determine which shard a piece of data belongs to. For example, in a users table, you might use user_id as the shard key. All data related to a specific user would then reside on the same shard. Choosing the right shard key is arguably the most important decision in a sharding strategy.

A good shard key has three essential properties:

  • High Cardinality: It should have many unique values to allow for wide distribution. A boolean flag like is_premium would be a terrible shard key as it only has two values.
  • Even Distribution: The key should naturally spread data and load evenly across shards, avoiding "hotspots" where one shard receives a disproportionate amount of traffic.
  • Query Alignment: It should align with your application's primary query patterns. If you frequently query data by user_id, sharding by user_id ensures most queries can be served by a single shard.

The video below explains these properties with clear examples of good and bad shard keys.

Sharding in System Design Interviews w/ Meta Staff Engineer

The same video from Hello Interview continues with a detailed analysis of what makes a good shard key. This is a common topic in system design interviews.

Watch the section discussing shard keys. Pay close attention to the three properties of a good shard key and the examples provided, such as why user_id is often a good choice, while is_premium or creation_date can lead to problems like low cardinality or hotspots.

Distribution Strategies

Once you have a shard key, you need a strategy to map its values to specific shards. There are three common approaches:

  1. Range-Based Sharding: Data is partitioned based on ranges of the shard key. For example, users with IDs 1-1,000,000 go to Shard 1, 1,000,001-2,000,000 go to Shard 2, and so on. This is simple to implement but can easily lead to hotspots, especially if the shard key is time-based or monotonically increasing (e.g., all new user sign-ups hitting the last shard).

  2. Hash-Based Sharding (or Key-Based Sharding): The shard key is fed into a hash function. The output of the hash determines which shard the data goes to (e.g., shard_index = hash(shard_key) % number_of_shards). This generally leads to a very even, random distribution of data, mitigating hotspots. A significant challenge here is rebalancing: when you add a new shard, the number of shards changes, and nearly all data must be reshuffled. Consistent Hashing is an advanced form of hash-based sharding that minimizes this reshuffling problem.

  3. Directory-Based Sharding: A lookup table (the "directory") explicitly maps each shard key value to a shard. This offers maximum flexibility—you can move data around by simply updating the lookup table. However, it introduces a single point of failure (the directory) and adds an extra network hop for every query to first consult the directory.

Sharding in System Design Interviews w/ Meta Staff Engineer

The video continues to explain these three distribution strategies, outlining their pros and cons.

Watch the next segment on distribution strategies. Understand the mechanics, advantages, and disadvantages of Range-based, Hash-based (including a mention of consistent hashing), and Directory-based sharding. Note the advice that hash-based sharding with consistent hashing is the industry default and often the expected answer in interviews.

3. Applying Sharding in Practice

Let's move from theory to how a sharded system is implemented. In most sharded architectures, your application doesn't talk to the shards directly. Instead, it communicates with a proxy or routing tier that understands the sharding scheme.

THE way to scale a database (sharding MySQL and Postgres)

This video from Ben Dicken provides a clear architectural diagram of a sharded setup.

Watch the segment from the start, which illustrates how a proxy (like Vitess for MySQL) sits between the application and the database shards. The proxy inspects the query, uses the shard key to determine the correct shard, and routes the request accordingly. This makes the sharded cluster appear as a single database to the application.

This routing logic can be implemented in various ways, from a dedicated proxy like Vitess or Citus to logic built directly into your application's data access layer.

Since you're proficient in Go, let's look at a concrete example. The following article demonstrates sharding an e-commerce order system.

Database Sharding in Go: A Practical E-commerce Case Study | GoFrame - A powerful framework for faster, easier, and more efficient project development

This article from the GoFrame project provides a practical case study that combines multiple sharding strategies, with Go code examples.

First, read the section Sharding Strategy Design. Notice how it uses a hybrid approach: Database sharding by user_id modulo (a hash-based strategy) to spread users across databases. Table sharding by order creation time (a range-based strategy) to organize data within each database. Next, review the Go code in the section Custom Sharding Rules. You don't need to understand the GoFrame framework details, but observe how the SchemaName function implements the hash-based sharding for databases (userId % schemaCount) and the TableName function implements range-based sharding for tables by month (createTime.Year(), createTime.Month()). This demonstrates how the routing logic can be encapsulated in code.

For PostgreSQL, a popular tool for sharding is the Citus extension. It turns a cluster of PostgreSQL servers into a single logical distributed database.

Sharding PostgreSQL with Citus and Golang

This article shows how simple the command to apply sharding can be when using a tool like Citus.

Skim the "Step 5" section, focusing on the SQL command: distributing data. The create_distributed_table('logs', 'id') command tells Citus to shard the logs table across the cluster, using the id column as the shard key. This single command abstracts away the complex work of data distribution.

4. Analyzing the Complexities

Sharding solves the write-scaling problem but introduces significant new complexities. Acing a system design interview requires not just proposing sharding, but also discussing its challenges and trade-offs.

Hotspots and Load Imbalance

Even with a good sharding strategy, you can get hotspots. The classic example is the "celebrity problem": if you shard a social media app by user_id, a celebrity's profile will generate massive traffic, overwhelming its designated shard.

Solutions include:

  • Compound Shard Keys: Instead of just user_id, you could shard by a combination like (user_id, post_id) to spread a single user's data across more shards.
  • Dedicated Shards: Move high-traffic tenants (like celebrities) to their own dedicated, possibly more powerful, shards.

Cross-Shard Queries

If a query doesn't include the shard key (e.g., searching for users by username when the shard key is user_id), the routing proxy must "scatter-gather"—query all shards and aggregate the results. These queries are slow, expensive, and can overwhelm the system.

Solutions include:

  • Good Shard Key Choice: The best defense is a shard key that aligns with your most common queries.
  • Caching: For frequent, expensive global queries (e.g., "trending posts"), cache the results.
  • Denormalization: Duplicate data where necessary to ensure that queries can be served from a single shard.

Distributed Transactions and Consistency

ACID transactions are straightforward on a single database. Across multiple shards, they become extremely difficult. If you need to deduct money from Bob's account on Shard 1 and add it to Alice's on Shard 2, what happens if the first operation succeeds but the second fails? This breaks atomicity.

While we'll explore this topic in-depth later, the two main approaches are:

  • Two-Phase Commit (2PC): A protocol that coordinates the transaction across all participating shards. It ensures consistency but is slow and complex.
  • Saga Pattern: An application-level pattern that breaks a transaction into a series of smaller, local transactions, each with a compensating action to undo it in case of failure.

Sharding in System Design Interviews w/ Meta Staff Engineer

This is the most critical part of the video for interview preparation. It covers the common follow-up questions you'll get after proposing a sharding solution.

Watch the final segments covering the main challenges: Hotspots and Cross-Shard Operations: This section details the celebrity problem and expensive scatter-gather queries, along with their solutions. Consistency: This part explains the difficulty of distributed transactions and introduces Two-Phase Commit and the Saga pattern.

Rebalancing

What happens when your shards fill up and you need to add more? This process is called rebalancing or resharding. You need to add new database servers and redistribute data from the existing shards to the new ones, all while the system is live. This is a complex and risky operation.

This diagram shows the process of rebalancing. When Node 4 is added, shards (like S4, S8, S12) are moved from existing Nodes 1, 2, and 3 to distribute the load more evenly.

Consistent hashing, as mentioned earlier, is a key algorithm designed to make this rebalancing process far more efficient by minimizing the amount of data that needs to move.

5. Discussing Sharding in an Interview

Proposing sharding in an interview demonstrates you're thinking about scale. Justifying it and articulating the trade-offs demonstrates seniority.

Sharding in System Design Interviews w/ Meta Staff Engineer

The final part of the Hello Interview video provides a superb framework for discussing sharding in an interview context.

Watch the closing segment from how sharding comes up in interviews. The four-step process—propose a shard key, choose a distribution strategy, call out trade-offs, and address growth—is a perfect template to structure your answer.

Conclusion

Today, we've dissected horizontal partitioning, the primary method for scaling a database's write capacity and storage. You've learned that while it's a powerful solution, it's not a silver bullet and introduces its own set of deep architectural challenges.

Key Takeaways:

  • Sharding is the process of horizontally partitioning data across multiple database instances to scale beyond the limits of a single server.
  • The choice of a shard key (high cardinality, even distribution, query alignment) and a distribution strategy (range, hash, directory) are the most critical design decisions.
  • A proxy/routing tier is typically used to direct application queries to the correct shard, abstracting the complexity from the application.
  • Sharding introduces significant complexities: hotspots, cross-shard queries, maintaining transactional consistency, and the operational burden of rebalancing. Being able to discuss these trade-offs is crucial.

In our previous lesson, you saw how to implement read replicas in PostgreSQL and MySQL. Today, you've seen the principles of sharding a relational database. In our next lesson, we will apply these concepts in a different context: implementing sharding in a MongoDB cluster, where you will see how a NoSQL database handles these challenges natively.

Can't find a good explanation? Sign up and we'll make it for you

Sign up