Learn Labs
7. Sharding

7.6 Technology deep dives


6.1 Vitess (manual key-range sharding for MySQL)

Problem it solves. Let an existing, huge MySQL deployment scale horizontally without rewriting the application, and without giving up MySQL's transactional semantics inside a shard.

Why not just a bigger MySQL? Single-node write throughput and dataset size ceilings; and at YouTube/Slack scale the ceiling is real.

Why not automatic sharding? Vitess deliberately gives operators manual control over key-range boundaries — because resharding is expensive and the timing matters (§3.4's "human in the loop").

How it works internally. vtgate is the routing tier (§4 approach ②) speaking the MySQL wire protocol, so the application thinks it's talking to one MySQL. It parses the query, consults the VSchema (which declares the sharding key per table), and routes or scatters. vttablet sits in front of each MySQL instance. Resharding is done online: new shards are created, data is copied, changes are streamed via VReplication (built on the binlog — Ch 6 §1.5c), and traffic is cut over per-shard-per-direction with the ability to reverse.

Deployment. Topology service (etcd/ZooKeeper/Consul) holding the authoritative shard map; vtgate fleet behind a normal L4 LB; vttablet + MySQL pairs.

Monitoring. Scatter-query rate (queries hitting all shards — this is the metric that tells you your VSchema is wrong); per-shard QPS and row-count skew; VReplication lag during a reshard; vtgate query-plan cache hit rate; cross-shard transaction rate.

Scaling. Split a shard when it's too big or too hot; the key range makes the split a clean range operation.

Backup. Per-shard MySQL backups + binlog; plus the topology service, which must be backed up separately — losing the shard map is losing the database.

What actually breaks.

  • Scatter queries by accident. A query that omits the sharding key hits every shard. One such query on a hot path silently converts your 100-shard cluster into 100× the load. This is §5.1's problem, made operational.
  • Choosing a sharding key that doesn't match the access pattern — irreversible without a full reshard.
  • Reshard cutover with in-flight writes (§4 problem 3).
  • Cross-shard transactions — Vitess offers 2PC but it is slow and discouraged (Ch 8).
  • AUTO_INCREMENT across shards — needs a sequence table or a different ID scheme.

6.2 Cassandra / ScyllaDB token-range sharding

Problem it solves. Sharding that rebalances with minimal data movement, no coordinator, and no fixed shard count.

Why not mod N? §3.2a — most keys move on every topology change.

Why not a fixed shard count? §3.2b — you must guess the count correctly at creation, can't exceed it with nodes, and resharding is a downtime event.

How it works internally. Murmur3 hashes the partition key into a token; the token space is divided into ranges with random boundaries, with 16 vnodes/node in Cassandra, 256 in ScyllaDB by default. Adding a node takes slices from several existing nodes rather than splitting one, which spreads the streaming load. Within a partition, rows are clustered (sorted) by the clustering columns — which is exactly the "partition key first, range query on later columns" trick of §3.2c. ScyllaDB additionally shards per CPU core within a node (the §1 NUMA/thread-per-core idea).

Deployment. Multi-DC with NetworkTopologyStrategy; token allocation algorithm configured to reduce imbalance; clients use a token-aware driver — §4 approach ③, so the client computes the token itself and connects directly to a replica.

Monitoring. Per-node ownership percentage (should be near-uniform; if not, your vnode count or token allocation is wrong); partition size distribution — the p99 and max are what matter, because one giant partition is the classic Cassandra outage; tombstones per read; hot-partition detection via per-table read/write latency outliers; streaming progress during topology change.

Scaling. Add nodes; ranges adjust automatically. You cannot easily change the partition key — it determines everything.

What actually breaks.

  • The unbounded partition. Modeling (user_id) as the whole partition key for a table that grows forever produces a multi-GB partition on one node. Reads time out, compactions take hours, and the node is a permanent hot spot. The fix is bucketing — (user_id, month) — which is the §3.1 timestamp-prefix lesson in another dress.
  • Hot key from a celebrity (§3.3) — token-uniformity does not imply load-uniformity.
  • Local secondary indexes used as if they were global. Cassandra's CREATE INDEX is a local index; a query on it scatters to every node. It is the single most misused feature in Cassandra.
  • Adding many nodes at once, saturating the network with streaming.
  • Gossip-based topology disagreement during a partition — weaker than consensus, by design.

6.3 DynamoDB (hash-range sharding + adaptive capacity)

Problem it solves. Fully managed sharding where the operator never thinks about partitions, with automatic splitting on both size and throughput.

Why not manual? The whole product promise is that shard management is invisible.

How it works internally. Partition key hashed to place the item; sort key orders items within the partition (again §3.2c). Partitions split when they exceed 10 GB or their throughput ceiling. Adaptive capacity / heat management (§3.3) reallocates throughput toward hot partitions, and isolates a single hot key onto its own partition when needed. Global secondary indexes are asynchronously maintained — hence eventually consistent (§5.2).

Monitoring. ThrottledRequests and ConsumedCapacity per table and per GSI (a throttled GSI throttles the base-table write); SuccessfulRequestLatency; hot-partition indicators via CloudWatch Contributor Insights — the only way to actually see key-level skew.

Scaling. On-demand or provisioned with autoscaling. Splits are automatic and one-way — partitions never merge, so a table that once spiked keeps its partition count, which permanently divides your provisioned throughput. This is a genuinely surprising operational fact.

What actually breaks.

  • A low-cardinality partition key. status = "PENDING" as a partition key puts everything on one partition, and no amount of provisioned capacity helps.
  • Throttling on a GSI silently throttling base-table writes.
  • Stale GSI reads treated as strongly consistent (§5.2).
  • Scan operations in production — the anti-pattern equivalent of Vitess scatter queries.
  • Partition count inflation after a traffic spike, permanently diluting per-partition capacity.

6.4 Elasticsearch (fixed shard count, local secondary indexes)

Problem it solves. Distributed full-text search where each shard is a self-contained Lucene index.

Why local indexes? A search index is a secondary index; term-partitioning it globally would require cross-shard postings intersection on every multi-term query (§5.2). Elasticsearch chooses document-partitioned, accepting scatter/gather.

How it works internally. shard = hash(routing) % number_of_primary_shards — the fixed shard count of §3.2b, chosen at index creation and immutable without a reindex. A search scatters to one copy of every shard, each computes local top-K, and the coordinating node gathers and merges. Scoring requires a second phase (query-then-fetch) because term statistics are per-shard.

Monitoring. Shard count per node; search latency broken into query vs fetch phase; the slowest shard per query (tail latency amplification made visible); rejected search-thread-pool tasks; skew in doc count per shard.

Scaling. More nodes redistributes shards but does not increase query throughput for a query that must touch every shard — §5.1's exact warning. Increasing parallelism requires more indices (e.g. time-based indices) plus routing so queries hit fewer of them.

What actually breaks.

  • Oversharding — the fixed-count guess made too high, so every query fans out to hundreds of shards and tail latency dominates.
  • Undersharding — the guess made too low, so shards grow past ~50 GB and you need a full reindex to change it.
  • Custom routing forgotten on the query side, so a routed write is searched with a full scatter.
  • Deep pagination across shards — every shard must produce from+size hits.
  • One slow node making every query slow, because every query touches every shard.

On this page