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:
- Dashboard — recent activity for a website, plus search by consent reference id.
- 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¶
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:
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_fromonwards. 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.
Click the diagram to pause or play.
- Live writes go to the old table always, and to the new table as well if
event_timeis 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_fromhistory 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_fromis 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_partitionsor nodetool table histograms) are what would tell you a website needs a narrower width before it becomes the next 131 GB partition.