Hard20 minDatabase Fundamentals
UpdatedAug 1, 2026
Edit

Sharding vs Partitioning

Question Variations

  • "What is the difference between horizontal partitioning and sharding?"
  • "When would you choose to shard a database instead of just vertically scaling the server?"
  • "What are the common challenges when performing joins across shards?"
  • "How do you handle 'hotspots' in a sharded database architecture?"

Why This Is Asked

This question reveals whether you can design databases for scale. Interviewers want to see that you understand the difference between splitting data within a single database (partitioning) and across multiple database servers (sharding), and that you can articulate the operational complexity each introduces.

Key Concepts

  • Partitioning splits a table into smaller pieces within the same database instance (horizontal or vertical)
  • Sharding distributes data across multiple independent database servers using a shard key
  • Sharding enables horizontal scaling beyond a single machine’s limits but introduces cross-shard query complexity
  • Cross-shard joins and distributed transactions are expensive and often require denormalization
  • Shard key selection is critical: poor keys create hotspots, good keys distribute load evenly

Question Variations

  • “What is the difference between horizontal partitioning and sharding?”
  • “When would you choose to shard a database instead of just vertically scaling the server?”
  • “What are the common challenges when performing joins across shards?”
  • “How do you handle ‘hotspots’ in a sharded database architecture?”

Answers by Technology

+ Add Variant
System DesignImprove this answer ✏️

Expected Answer

While both techniques split large datasets, they operate at different levels:

Sharding

Application Logic / Router

Shard 1 / Server A

Shard 2 / Server B

Shard 3 / Server C

Vertical/Horizontal Partitioning

Large Table

Partition A

Partition B

Partition C

Single Database Server

  • Partitioning: A database-level feature that splits a table into smaller pieces within a single database instance. It is usually transparent to the application.
  • Sharding: An architectural pattern that distributes data across multiple independent database servers. Each server (shard) holds a subset of the data.

Why It Matters

Scalability limits. Partitioning helps manage large tables and improves query performance (via partition pruning) but is still limited by a single server’s resources (CPU/RAM). Sharding allows for Horizontal Scaling, enabling a system to handle billions of rows and millions of requests by adding more machines. However, sharding introduces massive complexity: cross-shard joins become impossible, and transactions across shards require complex coordination.

Example Code

Partitioning (PostgreSQL)

-- Create a partitioned table by range
CREATE TABLE orders (
    id int,
    order_date date not null,
    total decimal
) PARTITION BY RANGE (order_date);

-- Create specific partitions
CREATE TABLE orders_2023_q1 PARTITION OF orders
    FOR VALUES FROM ('2023-01-01') TO ('2023-04-01');

Sharding (Application Logic)

// A simple shard routing function
function getShardConnection(userId) {
  const shardCount = 4;
  const shardId = userId % shardCount;
  return dbConnections[shardId];
}

const db = getShardConnection(user.id);
await db.query("INSERT INTO orders ...");

Common Mistakes

  • Premature Sharding: Sharding is an operational burden. Most apps should exhaust vertical scaling, read replicas, and caching before moving to a sharded architecture.
  • Choosing a bad Shard Key: A shard key that creates “hotspots” (e.g., sharding by country where 90% of users are in one country) defeats the purpose of sharding.
  • Confusing Partitioning with Sharding: Thinking that partitioning a table will allow it to scale beyond the storage limits of a single machine.

Follow-up Questions

  • What is a “Hot Shard”? (Answer: A shard that receives a disproportionate amount of traffic compared to others, usually due to an uneven distribution of the shard key).
  • How do you handle cross-shard queries? (Answer: Either by denormalizing data so it exists on all shards, or by using a query aggregator/proxy layer that executes queries on all shards and merges results).

References