Concepts

Epoch sharding

SlateDB allows one writer per database. Plural Telemetry gets many writers by routing each record to a storage shard, using a hash of its identity and the routing epoch in effect at its timestamp.

Fig. — Epoch-based sharding, live
now 11:18
shard 0shard 1epoch 0 · 2 shardsnow11:0012:0013:0014:0015:000x0…0x4…0x8…0xc…0xf…BLAKE3-128 hash spacerecord time →
Write flow · traced record

Waiting for the first traced record…

Shard ownership
shard 0writer-0
shard 1writer-1
Scaling up schedules a new epoch at the next aligned hour that is at least 2 minutes away, so every writer sees it before any record routes by it. Repeated requests before the cutover replace the pending epoch. Old data never moves, and late records route by their own timestamp.

Press Scale writers to schedule a new epoch, and watch the cutover land on the next hour. Send late record emits an entry with an old timestamp: it routes by the epoch that was in effect then, into a shard that already holds that time range. Switch to Read path to watch a 90-minute range query fan out across every shard, including the ones a scale-up added mid-range, and merge back on the reader.

Routing

Every product hashes a canonical key with the first 128 bits of BLAKE3:

ProductRouting keyRouted by
Metricsnamespace + canonical labelseach sample's timestamp. A series straddling a cutover is split into one write per epoch
Logsnamespace + canonical stream labelseach entry's timestamp
Tracesnamespace + trace IDthe trace's earliest span start

The hash selects a range in the epoch's hash-range map. Ranges are aligned to 4,096 routing slots (the high 12 bits), which caps a deployment at 4,096 storage shards. The slot is only a routing granularity and is never stored in keys.

text
record ─► (namespace + labels | trace ID, timestamp)
       ─► epoch = last epoch with effective_from <= timestamp
       ─► shard = epoch.range_for(blake3_128(key))
       ─► owner of shard (Lease holder)
            ├── this writer ──────────────► SlateDB put
            └── another writer ─► gRPC ──► SlateDB put   (retry on stale owner)

Stale ownership responses refresh the assignment and are retried up to write.remoteRetries times. Forwarding fan-out is bounded by write.remoteConcurrency.

Routing epochs

A ShardMap holds an ordered, append-only list of epochs, each with an effective_from_ns timestamp and a shard count. Because epochs only grow, every shard referenced by an older epoch still exists and stays readable.

Scaling writers from 2 to 4:

  1. The operator sees the new writer.replicas and creates pods writer-2 and writer-3.
  2. The Lease-elected coordinator appends epoch 1 to the ShardMap using a resourceVersion compare-and-swap. The epoch is effective at the next aligned boundary that's at least the lead time away: one-hour alignment and two minutes of lead by default.
  3. Writers and readers watch the ShardMap, apply only increasing generations, and acquire Leases for their assigned shards.
  4. At the cutover, new records route by epoch 1. Shards 2 and 3 start as new, empty SlateDB databases. No data is copied, cloned or drained.

Repeated scale requests before a cutover replace the pending epoch instead of stacking epochs. The coordinator then rebalances ownership across the StatefulSet, which moves Leases but not data.

Scale-down is blocked

Shards referenced by any epoch must stay writable until their data expires, so the operator refuses to reduce writer replicas and reports this in the product status.

Reads

Readers open every storage shard [0, shard_count) and merge:

  • Metrics deduplicates series by fingerprint across shards.
  • Logs merges streams with equal labels.
  • Traces merges partial traces from every shard a trace ID may have been routed to.

Every shard opened by a process shares one block cache and one meta cache, so reader cache memory stays flat as the shard count grows.

Why Kubernetes

PrimitiveUsed for
StatefulSetStable writer identities (writer-0, writer-1, …) across restarts and scaling
ShardMap CRThe single versioned source of truth for epochs and assignments
LeaseExclusive ownership of each physical shard, plus coordinator election

Kubernetes already gives you a CP datastore with watches and compare-and-swap. Using it means Plural Telemetry needs no ZooKeeper, etcd or gossip ring of its own, and adds no new dependency to clusters that already run it.

Tuning

  • Keep sharding.renewIntervalSeconds well below sharding.leaseDurationSeconds (defaults 5 and 15).
  • sharding.ioConcurrencyLimit (default 128) is a per-pod storage I/O budget, independent of shard count.
  • Every forwarding participant must share the internal auth token when auth.internalTokenSecretRef is set.