7.8 Decision cheat sheet
Only if data volume or WRITE throughput exceeds one machine.
Should I shard at all? Only if data volume or WRITE throughput exceeds one machine. Read throughput → use replicas instead. A single modern machine handles far more than folklore suggests, and sharding is a heavyweight solution mostly relevant at large scale.
Key-range or hash?
- Prefix with a higher-cardinality field (
sensor_id,user_id), and accept that cross-entity range queries now need one query per entity. - Fixed shard count — many shards per node; pick a highly divisible number; move whole shards.
- Hash-range sharding — shards split on demand, so the count adapts to data volume.
Composite keys are usually the answer. (partition_key, clustering_key) gives you hash-distributed shards and efficient ranges within a partition. Almost every well-designed sharded schema looks like this.
How many shards per node? More than one, always — it makes rebalancing cheap and lets you weight powerful nodes. Cassandra defaults to 16 vnodes/node, ScyllaDB to 256.
Automatic or manual rebalancing? Manual or semi-automatic if: you're near max write throughput, your failure detection is timeout-based, or a known traffic surge is coming. Fully automatic if the system is well below capacity and the vendor owns the operational risk.
Which routing model? Shard-aware client for lowest latency (one hop) if you control the client. Routing tier if clients are heterogeneous or you want a wire-protocol-compatible façade. Any-node forwarding for simplicity at the cost of an extra hop. In all three cases, put the authoritative map in a consensus-backed coordination service.
Local or global secondary index? Local when writes dominate, or when queries usually include the partition key. Global when read throughput exceeds write throughput and postings lists are short — and be explicit about whether you accept staleness or pay for a distributed transaction.
Multitenant: shard per tenant? Yes if tenants are meaningfully sized and you want resource/permission isolation, per-tenant restore, per-tenant residency, and easy GDPR deletion. No if you have hundreds of thousands of tiny tenants (overhead), or one whale tenant that doesn't fit a node (you'll need to shard within it anyway).