Limited Offer

30% OFF Lifetime Access ($139) with code SYSTEM30

TOPIC #62Beginner 9 min read

Sharding vs Partitioning vs Replication

💡
Core 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

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%
Rendering visual architecture flowchart...

01.1. 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.2. 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 Assessment 🧠

Step through 2 scenario questions to test your staff-level grasp.

Question 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 Chapter4.9 / 5.0 (38 ratings)

How clear and staff-actionable was this system breakdown?