Skip to content

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

StrategyPicks shards byBest for
Hash (default)hash(key) % N with virtual nodesUniform distribution, point lookups
Range(low_key, high_key) per shardRange scans, time-series
CompositeHash bucket + range within bucketHot 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:

  1. Identify hot shard (high CPU/memory/storage).
  2. Find optimal split point (median key or hash midpoint).
  3. Create two child shards.
  4. Copy data incrementally to the children.
  5. Atomic ring update.
  6. Redirect traffic.
  7. 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:

KnobDefaultWhat it does
target_cpu / target_memory0.6Per-shard utilization the rebalancer aims for. Lower = more headroom.
max_variance0.3Acceptable load spread across shards. Lower = more aggressive rebalances.
min_improvement0.1Plan must improve load by ≥10% to be executed. Higher = fewer no-op rebalances.
check_interval_secs300How 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:

  1. Bulk Transfer — copy existing data to the target.
  2. Incremental Sync — stream changes that arrived during phase 1.
  3. Switchover — atomically redirect traffic.
  4. Verification — optional consistency check.
  5. 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