Review this configuration. A 60 TB events table is partitioned by ingest_date and written by a streaming job that commits every 60 seconds. A compaction job runs hourly over the last 24 hours with a 512 MB target and sorts each file by event_timestamp. Snapshot expiry runs monthly with 90-day retention. 85% of queries filter on customer_id over a 7-day range. What would you remove, what would you change and what would you leave alone?
Show the full answer Hide the answer
What is actually required
One dominant access path: customer over a week. Everything in the configuration should be judged against that, plus a recovery story for the streaming writer and a retention obligation nobody has written down yet.
The one change that matters
Sort by customer_id, not by event_timestamp. Compaction's sort key decides which column's
per-file minimum and maximum statistics are narrow enough to skip files with. Inside a partition
that is already bounded by ingest_date, sorting by event_timestamp is close to a no-op: the
timestamps in that partition already span roughly one day whatever order they are written in.
Meanwhile every file's customer_id range spans the whole key space, so the engine can skip
nothing and the 85% case reads every file in seven partitions.
Sorting by customer_id typically collapses that to a small fraction of files, because the
statistics the engine consults before reading are per file and only a narrow range lets it skip
one. On a table of this shape expect bytes scanned to fall by something on the order of 10 to 50
times for a single-customer week, with no change to storage cost.
The decision rule: choose the sort key from the most common filter that is not already the partition key. A table has one physical order, and every open table format has worked this way since the formats appeared around 2017, so the sort key is the single largest performance decision available after partitioning.
What I would remove
The 90-day snapshot retention, and the monthly cadence that goes with it. A writer committing every 60 seconds produces roughly 1,440 snapshots a day, so 90 days is on the order of 130,000 snapshots. Every compaction rewrite keeps its pre-compaction files alive for the whole retention window, so the table is paying to store roughly two copies of recent data plus a metadata tree the planner walks on every query. Seven days of retention, expired daily, is what a recovery story actually needs. Confirm what that story is before cutting it: if someone is relying on time travel for a monthly regulatory reproduction, the answer is a snapshot tag on the reporting date rather than 90 days of everything.
What I would change second
The compaction window. Hourly compaction across the last 24 hours rewrites yesterday's data 24 times over. Compact the current day frequently, compact each completed day once, and never touch it again. This is pure write amplification with no benefit, and on 60 TB it is a visible share of the platform's compute.
What I would leave alone
- The 512 MB target. It suits scan-heavy work, and raising it to 1 GB would hurt the 15% of queries doing narrower lookups while increasing rewrite cost.
- Partitioning by
ingest_date. It matches the retention boundary and the 7-day predicate, and replacing it with customer partitioning would create severe skew. - The 60-second commit. That is a freshness requirement, not a defect. The defect is failing to compensate for it.
How I would argue this in the review
Not with the principle. Take one representative customer-week query, record files opened and bytes scanned, change the sort key on a copy of three partitions, and show the same two numbers again. Layout arguments are won with a before-and-after on a real query and lost with a diagram.
When this is the wrong answer
If queries were genuinely time-range scans across all customers, the existing sort is correct and the change would make things worse. The configuration is not wrong in the abstract; it is wrong against a workload that presumably looked different when it was written, which is the usual way a maintenance setup becomes stale.