Sharding is a horizontal partitioning technique for distributed databases and blockchain networks in which a dataset or workload is divided into disjoint subsets called shards, each maintained by a distinct subset of nodes, so that the total system throughput scales with the number of shards rather than being bounded by the capacity of a single node. In databases, sharding routes queries to the appropriate shard by a sharding key. In blockchain, each shard processes its own subset of transactions and stores its own portion of the state, with cross-shard communication handled by a coordination layer. Sharding dramatically increases transaction throughput and reduces storage requirements per node at the cost of increased architectural complexity and cross-shard coordination overhead.

Content

  • Sharding as a database concept dates to large-scale web systems in the late 1990s and 2000s, when companies such as eBay, Google, and Facebook horizontally partitioned relational databases across multiple servers to handle growth beyond what a single machine could sustain. The sharding key — a value such as user ID or geographic region — determines which shard stores and processes each record, ensuring that related data co-locates for efficient query processing. MongoDB, Cassandra, and Vitess popularised database sharding in the 2010s.
  • In blockchain, sharding attempts to solve the scalability trilemma: the observation that it is difficult to simultaneously achieve decentralisation, security, and high throughput. A traditional blockchain requires every full node to process every transaction, bounding throughput to what a single node can handle (approximately 15–30 TPS for Ethereum pre-Merge). In a sharded blockchain, the validator set is divided into committees, each responsible for one shard. Transactions are routed to the appropriate shard based on the sender’s or contract’s address. The beacon chain or coordination layer maintains cross-shard state roots and handles finality. Cross-shard transactions require a receipt-based protocol: the originating shard creates a receipt that the destination shard can redeem in a subsequent block.
  • Ethereum’s roadmap has included database sharding (originally “state sharding”) since 2016. The execution sharding plan was substantially revised after the rise of Layer 2 rollups: rather than sharding execution on Layer 1, Ethereum now pursues “danksharding” — providing large-scale blob data availability to L2 rollups via a data sharding layer, without sharding execution itself. Proto-danksharding (EIP-4844, activated March 2024) introduced “blob” data types providing L2s with cheap temporary data storage, a stepping stone toward full danksharding. Near Protocol, Zilliqa, and Harmony implement transaction-level sharding. The Ethereum Beacon Chain uses committee-based sharding for attestations.
  • In 2024–2025, EIP-4844 has substantially reduced the cost of posting rollup proofs to Ethereum L1, validating the data sharding approach. Full danksharding with 64 data shards is targeted for a future Ethereum upgrade. Academic research continues on the security of sharded systems — particularly the one-percent attack and cross-shard MEV — and on ZK-rollup-based sharding where validity proofs eliminate the need for fraud-proof windows. In databases, distributed OLTP systems such as CockroachDB and Google Spanner provide automatic sharding with global transactions, influencing the design of Web3 storage systems.