TOPIC #62Beginner 9 min read

Sharding vs Partitioning vs Replication

CSD
CompleteSystemDesign Editorial
Report an issue
Key takeawayCore Architecture Summary

Disentangle the 3 database scaling pillars: Vertical Partitioning (Column splitting), Table Partitioning (Single node), Sharding (Multi-node), and Replication (Duplication).

Key Glossary Concepts in this TopicAll Glossary Terms
Interactive Lab · 🏛️ Scaling Pillars AllocatorFull lab guide

Replication vs Partitioning vs Sharding Allocator

Compose the 3 scaling pillars to satisfy a workload demand — each pillar fixes a different bottleneck.

Writes OVER CAPACITY

need 12,000/s / have 10,000/s

only SHARDING scales this

Reads OVER CAPACITY

need 30,000/s / have 20,000/s

REPLICAS scale this linearly

Storage OVER CAPACITY

need 3 TB / have 2 TB

replication adds ZERO capacity

Physical servers

1

each shard = 1 primary + 0 replicas

Data per node

100%

100% = replication duplicate, 1/N = shard subset

Time-range scan

100% of table

partition pruning 1x faster

Native SQL JOINs

YES

complexity score 0/5

Notice the failure pattern: adding replicas never fixes the write or storage bar, and partitioning never raises the single-node CPU ceiling. Production architectures stack all three — a sharded cluster where each shard is partitioned by month and backed by two read replicas.

The 3 Database Scaling Pillars: Replication vs Table Partitioning vs Sharding 🏛️

Replication duplicates identical data across nodes; Table Partitioning splits large tables on a single node; Sharding distributes data subsets across independent physical servers.

The 3 Database Scaling Pillars: Replication vs Table Partitioning vs Sharding 🏛️
100%
Touchpad: Pinch to zoom • Drag to pan
Rendering visual architecture flowchart...

01.Disentangling the 3 Core Scaling Techniques

Engineers frequently confuse Replication, Partitioning, and Sharding in system design discussions. These three mechanisms operate at different physical layers and solve fundamentally different bottlenecks:

1. Database Replication (Data Duplication)

  • What it does: Copies the exact same complete dataset (100% of rows) across two or more physical server instances.
  • Hardware Topology: Multi-node.
  • Problem it Solves: High Availability (HA), disaster recovery failover, and scaling READ throughput.
  • What it cannot do: Does not scale write throughput (Primary is still single-node) and does not scale storage capacity (every node stores the full dataset).

2. Table Partitioning (Single-Node Logical Division)

  • What it does: Divides a single massive table (e.g. 500M rows) into smaller physical sub-table files (e.g. orders_2025, orders_2026) inside a SINGLE database server instance.
  • Hardware Topology: Single node.
  • Problem it Solves: Partition Pruning (the query engine skips scanning entire years of historical data) and instant operational maintenance (dropping an entire month of old logs with DROP TABLE logs_jan_2024 in 0ms without running slow row-by-row DELETE queries).

3. Database Sharding (Horizontal Partitioning Across Nodes)

  • What it does: Partitions a table and distributes distinct subsets of rows (e.g. 25% of rows per node) across multiple independent physical database servers.
  • Hardware Topology: Multi-node distributed cluster.
  • Problem it Solves: Scaling WRITE throughput and scaling total storage beyond single-server hardware limits.

02.Comparative Architectural Matrix

Dimensional AxisReplicationTable PartitioningSharding
Physical TopologyMultiple ServersSingle ServerMultiple Servers
Data on Each NodeFull 100% DuplicateSegmented Sub-tablesDisjoint Subset (~1/Nth)
Scales Read Throughput?Yes (Linear: 5 nodes = 5x reads)Slightly (via Partition Pruning)Yes
Scales Write Throughput?No (Primary is bottleneck)No (Bound by single server disk/CPU)Yes (Linear: 10 nodes = 10x writes)
Scales Total Storage Capacity?No (Every node stores all data)No (Bound by single server disk)Yes (Petabyte-scale aggregation)
ACID Multi-Table JOINs?Yes (Full relational support)Yes (Within local server)❌ No (Cross-shard joins prohibited)

Architectural Trade-offs & Production Realities

Architectural Advantages

  • Replication is simple to configure and provides instant failover redundancy.
  • Table partitioning speeds up single-node queries with zero distributed system complexity.
  • Sharding unlocks petabyte-scale storage and millions of writes/sec.

Trade-offs & Constraints

  • Table partitioning on a single node is still bound by single-node CPU/RAM limits.
  • Sharding introduces distributed routing and cross-shard query complexity.
Production Implementation in Big Tech
Uber & GitLab• Combining Partitioning, Sharding, and Replication

GitLab partitions heavy CI/CD build log tables locally by month using PostgreSQL declarative partitioning (enabling instant pruning and table truncation), runs asynchronous read replicas for web dashboards, and shards repository metadata across distributed database clusters.

Staff+ Engineering Takeaways

  • Replication = Duplicates full data across nodes (Read scale + High Availability).
  • Table Partitioning = Splits large tables into sub-files on 1 single node (Partition pruning).
  • Sharding = Distributes disjoint row subsets across multiple independent physical servers (Write scale).
  • Production systems combine all three in a hierarchical topology.

Topic Knowledge Check

Exercise 1 of 2 • Test your architectural comprehension.

Exercise 1 of 20 answered
1

If a PostgreSQL database table contains 200,000,000 rows and is partitioned by month using native Declarative Table Partitioning on a single server, what is the primary performance benefit?

Rate This Architecture ChapterFeedback & Rating

How clear and actionable was this distributed systems breakdown?

Interactive Engineering Workbenches: