How would you choose partition granularity and clustering columns for a multi-terabyte BigQuery event table holding five years of history?
answer
- start from the query log, not the schema
- two limits: how many, and how small
- coarsest that still prunes
- four slots, most-filtered first
- one choice is reversible, one is not
basics
~20 sPick the coarsest granularity that still prunes for the dominant queries: daily for five years of history keeps partition count and partition size sane, hourly usually does not. Then cluster on the most-filtered high-cardinality columns, leading with the one queries name most, and enforce a required partition filter plus expiration.
solid answer
~50 sStart from the query log, not the schema. Two limits bound the choice: BigQuery caps partitions per table, and very small partitions read inefficiently because per-partition overhead dominates. Five years of hourly partitions is far too many; five years of daily partitions is a manageable count and, at multi-terabyte scale, gives partitions large enough to be worth reading. If a single day were tiny, monthly with clustering on the timestamp would be better; hourly is justified only for high-volume tables whose queries genuinely ask for hour windows and whose retention is short. Then spend the four clustering slots on the columns queries actually filter, ordered most-filtered first — typically `tenant_id`, then a secondary key. Finish with `require_partition_filter = TRUE` so nobody scans five years by accident, and `partition_expiration_days` so cold history stops existing. Revisit the clustering list against `INFORMATION_SCHEMA.JOBS` as query patterns drift; the partition axis you must get right up front, because changing it means rewriting the table.
code
sql · 13 linesCREATE TABLE ds.events (
event_ts TIMESTAMP,
tenant_id STRING,
user_id STRING,
event_type STRING,
amount NUMERIC
)
PARTITION BY DATE(event_ts)
CLUSTER BY tenant_id, user_id
OPTIONS (
partition_expiration_days = 1825,
require_partition_filter = TRUE
);go deeper
Know the shape of the standard answer: partition by date, cluster on the columns queries filter, and set an expiration so old data does not accumulate.
Explain why granularity is bounded from both sides — a partition cap above and inefficient small partitions below — and why clustering column order follows how often each column is filtered.
Derive the design from real query history, justify a non-daily granularity when the data warrants it, and put the required-partition-filter and expiration guardrails in place as part of the table definition.
Own the irreversibility: the partition axis is a rewrite to change, so justify it explicitly, plan the review cadence for the clustering list, and decide when a secondary access pattern earns its own materialized view or table rather than distorting the primary layout.
## Frame the decision as two constraints and one input The input is the workload: which columns queries filter, over what time windows, and how often. Everything below is downstream of reading real query history rather than guessing. The constraints are structural. First, a BigQuery table has a hard maximum number of partitions, so granularity multiplied by retention must stay under it — that alone eliminates hourly partitioning for a five-year table. Second, partitions that are too small cost more than they save: each partition carries metadata, and reads of many tiny partitions do not amortise as well as reads of a few substantial ones. The commonly cited guidance is that a table under roughly a gigabyte gains little from partitioning at all, and the same logic applies per partition. ## Choosing granularity Work out the average partition size at each candidate granularity: total size divided by number of periods. For a multi-terabyte table over five years, daily gives partitions on the order of a gigabyte or more — comfortably worth reading and pruning — and a partition count that is large but within limits. That is why daily is the default answer and the one an interviewer expects you to reach, with reasoning. Deviate deliberately: - **Monthly** when the table is large in total but thin per day, or when queries genuinely ask month-at-a-time questions. You lose day-level pruning, so recover it by clustering on the timestamp: within a monthly partition, sorted-by-time data still lets block skipping narrow a single day. - **Hourly** only when volume per hour is already substantial *and* the dominant queries ask for hour windows *and* retention is short enough to stay under the partition cap. Hourly partitioning of long-retention data is the classic self-inflicted wound: an enormous partition count, tiny partitions, and no query that benefits. Granularity is irreversible in place. Changing it means `CREATE TABLE ... PARTITION BY ... AS SELECT * FROM old` and a rename, which at multi-terabyte scale is a real project. Spend the design time here. ## Choosing clustering columns You get four, and their order behaves like a composite key prefix. Rank candidates by how often queries filter on them, from the query log: ```sql SELECT query, total_bytes_billed FROM `region-us`.INFORMATION_SCHEMA.JOBS WHERE job_type = 'QUERY' AND creation_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 30 DAY) ORDER BY total_bytes_billed DESC LIMIT 200; ``` The column that appears in nearly every `WHERE` clause goes first — usually a tenant, account or customer id in a multi-tenant table. A second column that frequently co-occurs goes next. Adding a third and fourth is often ceremonial: if almost no query filters on them, they contribute little and they lengthen the sort. Cardinality is a secondary consideration; a boolean or a three-value enum as the leading clustering column wastes the slot because it barely narrows anything. Clustering is cheap to revise: `ALTER TABLE ... SET OPTIONS` changes the specification, new writes adopt it, and background reclustering converges existing data for free. So the honest strategy is to be decisive about the partition axis and iterative about the clustering list. ## The guardrails that come with the design ```sql CREATE TABLE ds.events ( event_ts TIMESTAMP, tenant_id STRING, user_id STRING, event_type STRING, amount NUMERIC ) PARTITION BY DATE(event_ts) CLUSTER BY tenant_id, user_id OPTIONS ( partition_expiration_days = 1825, require_partition_filter = TRUE ); ``` `require_partition_filter` turns "someone ran an unbounded query over five years" from a surprise invoice into an immediate error. It does not guarantee the filter prunes well, but it removes the worst class of accident. `partition_expiration_days` drops aged partitions as a metadata operation — free, and far cheaper than a scheduled `DELETE`. ## What to do about the queries this layout will not serve One partition axis and four clustering columns cannot satisfy every access pattern. When a second, unrelated pattern matters — say, lookups by `event_type` across all time — do not distort the primary layout for it. The alternatives are a materialized view or a scheduled summary table shaped for that pattern, or a second physical table maintained by the pipeline. That is a cost trade you make explicitly: extra storage and refresh work in exchange for not compromising the layout that serves ninety percent of the traffic. ## How you would review it later Commit to re-examining the design on a cadence. The signals: the top jobs by billed bytes and whether their filters match the clustering prefix; whether `__NULL__` or `__UNPARTITIONED__` partitions are accumulating rows; and partition sizes and counts from `INFORMATION_SCHEMA.PARTITIONS`. Query patterns drift as products change, and a clustering list that was right two years ago frequently is not. Because the clustering specification is alterable and the partitioning is not, the review is mostly about clustering — which is exactly why the partition granularity deserves the careful call up front.
- What goes wrong if you partition five years of data hourly?You collide with the per-table partition cap, and even where it fits you get an enormous number of small partitions whose metadata and read overhead outweigh the finer pruning. Almost no query asks for a single hour five years back. Hourly makes sense only for high-volume, short-retention tables whose dominant queries really are hour-scoped.
- The workload has a second access pattern the layout cannot serve. What do you do?Do not compromise the primary layout. Build a materialized view or a scheduled summary table shaped for the secondary pattern, or maintain a second physical table from the pipeline. You are trading storage and refresh work for keeping the layout optimal for the dominant traffic, and that trade should be made explicitly with the cost written down.
- How do you decide when the clustering columns need revisiting?Watch the top jobs by billed bytes over the last month and check whether their filters name the leading clustering column. When the expensive queries consistently filter something the clustering does not lead with, the list is stale. Clustering is alterable in place with background reclustering converging existing data, so acting on that signal is cheap.
saying these in an interview costs you the question
- Choosing hourly partitions for long retention
- Filling all four clustering slots regardless of query patterns
- Leading clustering with a low-cardinality flag column
- Assuming partition granularity can be changed later
- Designing the layout without reading the query log