Skip to content

The 131 GB partition: per-tenant time buckets in ScyllaDB

Published 6 October 2026

A problem faced at Mozilor.

Every time a visitor accepts or rejects cookies on a customer's website, a consent log is written: a small record that is never updated and has to be kept for compliance. Customers read the logs back in two ways:

  1. Dashboard — recent activity for a website, plus search by consent reference id.
  2. Export — every log in a date range, when they need to prove what happened.

The logs live on a self-managed ScyllaDB cluster. Volume kept growing, and it didn't grow evenly: most websites produce a steady trickle, a few produce torrents, and the torrents kept getting bigger. Writes slowed, reads got unpredictable, deletes started timing out. The cluster metadata showed why — one partition at 131.9 GB, another over 120 GB, a third at 47.8 GB, and a long tail heading the same way.

The original key

PRIMARY KEY ((tenant_id), event_time, status, entry_id)

One partition per website, clustered by time. Both read patterns are single-partition queries, which is exactly why it's the obvious design. The flaw: every log a website ever produces lands in that one partition, and it only grows. A thousand events a day adds a few MB a month; millions a day adds GBs.

Why one huge partition hurts everything

In ScyllaDB (and Cassandra) a partition is the unit of everything:

  • It never splits across nodes. All of it lives on the same replica set, so the busiest websites pin the same few replicas.
  • Compaction and repair move the whole thing. A 131 GB partition is a 131 GB unit of work.
  • Reads pay for the partition, not the result. Indexes and caches work at partition granularity, so even a narrow slice of a huge partition is expensive.
  • Deletes are worst. A delete writes a tombstone that compaction clears later. Deleting a website's whole history creates tombstones faster than compaction can clear them, until reads wade through tombstones and deletes time out.

The usual guidance is to keep partitions under ~100 MB. The worst ones here were over a thousand times that.

Measure first

Normalise each website's volume by how long it has been a customer, take the monthly average, and compare it to a 100 MB target partition size. The result:

Bucket width needed Websites
monthly the clear majority — and will be for years
weekly most of the few thousand that don't fit monthly
daily several hundred
4 hours a couple of dozen
2 hours 2
finer none

That distribution is the whole story. No single width is right for everyone.

Three fixes that don't work

Fix Why it fails
Bucket everyone by month The 131 GB site writes ~7 GB a month — 70× the target. It doesn't help the sites that have the problem.
Bucket everyone by day Fine for the big sites, but every small site gets 365 tiny partitions a year, and a one-year export becomes 365 queries instead of one.
Add a hash to the key, e.g. hash(consent_ref) % N Fixes sizing, but scatters each day across N partitions, so every date-range read becomes a fan-out and merge.

The fix: bucket width as a per-tenant setting

Add a time component to the partition key:

PRIMARY KEY ((tenant_id, time_bucket), event_time, consent_ref, entry_id)

time_bucket is a short string derived from the row's own timestamp, at a width that is configured per website and can change over time. The width lives in a small registry table:

CREATE TABLE bucket_registry (
    tenant_id      text,
    effective_from timestamp,
    width          text,          -- '1mo' | '1w' | '1d' | '4h' | '2h'
    PRIMARY KEY ((tenant_id), effective_from)
) WITH CLUSTERING ORDER BY (effective_from DESC);

Rows are newest first. The width for a timestamp is the first row whose effective_from is at or before it. No rows means monthly, which is why most customers never appear in this table.

def width_for(registry_rows, ts):
    # registry_rows: this tenant's rows, newest effective_from first
    for row in registry_rows:
        if row.effective_from <= ts:
            return row.width
    return "1mo"


def time_bucket(ts, width):
    if width == "1mo":
        return ts.strftime("%Y-%m")                       # 2026-07
    if width == "1w":
        year, week, _ = ts.isocalendar()
        return f"{year}-W{week:02d}"                      # 2026-W28
    if width == "1d":
        return ts.strftime("%Y-%m-%d")                    # 2026-07-08
    hours = int(width.removesuffix("h"))                  # 4h, 2h
    return f"{ts:%Y-%m-%d}T{ts.hour // hours * hours:02d}"  # 2026-07-08T04

(The exact bucket string format and registry schema above are my reconstruction; the shape is what matters.)

Three properties make this work:

  • It's a pure function. Given a row and the cached registry, every service computes the same bucket. No counters, sequence numbers or coordination between writers.
  • It uses event time, never arrival time. A row stamped 8 July that arrives on 10 July still lands in 8 July's bucket, which is where readers will look for it.
  • History never changes. A width change only affects buckets written from effective_from onwards. Nothing is rewritten, so a bad config change can't corrupt existing data. (Rewriting old oversized partitions into the narrower layout is a separate job.)

The 131 GB website now uses 2-hour buckets and writes partitions of roughly 35 MB. A few thousand others got a registry row sized to their own volume; everyone else stays on monthly.

Reading a date range

Because the width can change mid-range, the reader walks the registry, not a fixed stride:

def buckets_for_range(registry_rows, start, end):
    # Split [start, end) at every effective_from inside it, then enumerate buckets
    # at that segment's width.
    cuts = sorted({start, end, *(r.effective_from for r in registry_rows
                                 if start < r.effective_from < end)})
    buckets = []
    for seg_start, seg_end in zip(cuts, cuts[1:]):
        width = width_for(registry_rows, seg_start)
        buckets += enumerate_buckets(seg_start, seg_end, width)  # one per width step
    return buckets

The dashboard opens just the current bucket. An export opens only the buckets in its range, and a small website's one-year export is still about twelve queries.

The materialized view needed a different key

Dashboard search by consent reference id is served by a materialized view keyed on the old layout. It can't reuse the new key: someone searching by reference id doesn't know when the record was created, so they can't compute time_bucket. Instead the view is partitioned by website plus a hash of the reference id:

PRIMARY KEY ((tenant_id, ref_bucket), consent_ref, event_time, time_bucket, entry_id)
-- ref_bucket = hash(consent_ref) % N

Hash the id you're searching for and you know exactly which partition to open.

Hashing was rejected for the main table, so why is it fine here? Because the view is only ever read by reference id. There's no date range for the hash to scatter, and each view row is a handful of columns, so partitions stay small.

Migrating without downtime

Partition keys can't be changed in place, so 18 months of data had to move to a new table while both tables served live traffic.

Dual-write migration to the time-bucketed consent log table A consent event reaches one of the writer services, which reads the website's bucket width from a cached registry. The writer always writes to the old table, and also writes to the new bucketed table when the event time is at or after a fixed cutoff. Rows before the cutoff are exported from the old table, given a time bucket by a Spark job and bulk-loaded into the new table. Readers resolve which buckets to open from the registry and then read the new table. consent event writer services 5 services, same bucket fn bucket registry per-site width, cached dashboard / export readers old table PK ((tenant_id), …) Spark backfill export, bucket, load new table PK ((tenant_id, time_bucket), …) always if event_time ≥ cutoff rows < cutoff width which buckets? 1 · A consent event reaches a writer service 2 · The writer looks up the site's bucket width (default: monthly) 3 · Dual write: old table always, new table if past the cutoff 4 · Everything before the cutoff is exported from the old table 5 · Spark computes each row's bucket and bulk-loads it 6 · Readers map a date range to buckets, then read the new table

Click the diagram to pause or play.

  • Live writes go to the old table always, and to the new table as well if event_time is at or after a fixed, hard-coded cutoff.
  • Deletes are dual-written too, so the new table never keeps rows the old one removed.
  • History before the cutoff was exported, bucketed and bulk-loaded by a Spark job.

Because the cutoff splits the data by event time, the live path and the backfill never write the same row. Five microservices had to learn to compute buckets identically for writing, reading and deleting — which is where "pure function of the row plus a cached config" pays off.

Where it landed

  • Deletes complete. A website's delete spreads over many small partitions, each compacting its tombstones on its own schedule.
  • Writes are faster and steadier. Load spreads across the cluster instead of pinning a few replicas that were also compacting a multi-GB partition.
  • Compaction and repair got cheap. Work is proportional to what changed, not to a 131 GB partition.
  • Reads are bounded by the query, not the history.

One table, one code path, and a dial set individually for each customer that needs it.

Takeaways

  • Partition by tenant alone is a time bomb for append-only data. It's correct on day one and wrong on day 500. Ask how big the largest partition gets in two years, not the average.
  • Skew means one global setting is wrong for someone. Making bucket width a per-tenant config with an effective_from history beats picking a compromise width.
  • Prefer derivable keys over coordinated ones. A bucket computed from the row's own timestamp needs no counters and no lookups beyond a cached config.
  • Hash buckets are fine when the read is a point lookup by the hashed field, and bad when the read is a range over something else.

Questions I'd still ask

  • Late events vs. the cutoff. The live path only writes rows at or after the cutoff to the new table. A late event stamped before the cutoff that arrives after the Spark export lands only in the old table. Was the export taken after a grace period, or reconciled with a final diff?
  • Registry cache staleness. If a writer's cached registry lags a width change, it writes to the old width's bucket while readers look in the new one. Presumably effective_from is set in the future, past the cache TTL, so every service sees the change before it takes effect.
  • Monitoring the dial. Partition-size alerts (e.g. from system.large_partitions or nodetool table histograms) are what would tell you a website needs a narrower width before it becomes the next 131 GB partition.