Sharding & Rebalancing
Sharding & Rebalancing — Scale to Petabytes Without Downtime
UVP
When a single Raft group can’t hold your dataset, you shard. Most databases force you to pick a fixed sharding scheme up front and live with it forever. HeliosDB ships elastic sharding: hot shards split themselves, cold shards merge themselves, and the rebalancer moves data between nodes online — all while serving queries. Pick hash, range, or composite. Pick virtual-node count for distribution. Tune split/merge thresholds. The system observes load and reshapes the cluster while you sleep — with <10% data movement guaranteed during any rebalance.
Prerequisites
- A multi-node HeliosDB Full cluster (see raft-setup.md).
- A workload that doesn’t fit comfortably on a single shard (or a forecast that says it won’t soon).
- About 25 minutes.
1. The Three Strategies
| Strategy | Picks shards by | Best for |
|---|---|---|
| Hash (default) | hash(key) % N with virtual nodes | Uniform distribution, point lookups |
| Range | (low_key, high_key) per shard | Range scans, time-series |
| Composite | Hash bucket + range within bucket | Hot keys spread across buckets, but range scans within |
Hash is the safest default. Range becomes attractive when most queries are WHERE created_at BETWEEN.... Composite is a tuning knob, not a starting choice.
2. Configure the Elastic Shard Manager
Two important defaults:
virtual_nodes_per_shard: 150— consistent hashing uses virtual nodes so that every physical shard owns roughly equal slices of the ring. 100-200 is recommended for production. Lower → uneven distribution. Higher → more memory in the ring map.split_threshold: 0.8/merge_threshold: 0.3— fractions of CPU/memory utilization that trigger automatic split (hot) or merge (cold).
3. How a Shard Split Works
What happens:
- Identify hot shard (high CPU/memory/storage).
- Find optimal split point (median key or hash midpoint).
- Create two child shards.
- Copy data incrementally to the children.
- Atomic ring update.
- Redirect traffic.
- Delete parent shard.
The whole sequence is online. Reads and writes against shard-1 succeed throughout — the planner just reroutes them at step 5.
Indicative benchmark (representative fixtures; reproduce on your own hardware): ~100ms for 10K keys.
4. Merge Cold Shards
Same shape as split, in reverse. Constraint: shards must be adjacent in key space. The merger refuses non-adjacent merges to keep range queries fast.
Indicative benchmark: ~80ms for 10K keys.
5. Configure Auto-Rebalancing
The rebalancer runs in the background when auto_rebalance: true. Its config:
| Knob | Default | What it does |
|---|---|---|
target_cpu / target_memory | 0.6 | Per-shard utilization the rebalancer aims for. Lower = more headroom. |
max_variance | 0.3 | Acceptable load spread across shards. Lower = more aggressive rebalances. |
min_improvement | 0.1 | Plan must improve load by ≥10% to be executed. Higher = fewer no-op rebalances. |
check_interval_secs | 300 | How often to evaluate. Lower = faster reaction; higher = more stability. |
6. Rebalance Guarantees
Key invariant: “Guarantees <10% data movement relative to total cluster size”. The planner refuses any plan that would shuffle more.
If a step fails mid-execution, the rebalancer rolls back automatically.
7. Schema-Based Sharding (Multi-Tenant Twist)
The Full edition also supports schema-based sharding, where each tenant gets its own schema and the entire schema lives on one shard:
CREATE SCHEMA tenant_1234 DISTRIBUTED;
CREATE TABLE tenant_1234.users ( user_id SERIAL PRIMARY KEY, name TEXT, email TEXT);
CREATE TABLE tenant_1234.orders ( order_id SERIAL PRIMARY KEY, user_id INT REFERENCES tenant_1234.users(user_id), amount DECIMAL);No tenant_id column needed in every table. Foreign keys work natively because all of tenant_1234.* lives on the same physical shard.
8. Migrate a Shard to a Different Node
Moving shard-4 from one node to another is a four-phase online operation:
- Bulk Transfer — copy existing data to the target.
- Incremental Sync — stream changes that arrived during phase 1.
- Switchover — atomically redirect traffic.
- Verification — optional consistency check.
- Cleanup — drop the data on the source.
No downtime. No client retries.
9. Watch Shard Statistics
Per-shard CPU, memory, storage and requests/sec feed Prometheus via the standard observability surface.
Where Next
- raft-setup.md — each shard is its own Raft group.