In BigQuery the cost and speed of most queries are set by how many bytes they read. The storage is columnar, so selecting fewer columns helps, but filters on rows only help if the table is laid out so BigQuery can skip whole chunks of rows without reading them. Partitioning and clustering are the two tools for that layout. Chosen well, they turn a query over two years of events into a read of a few gigabytes. Chosen badly, they add limits and maintenance without saving anything.

This article explains how each one works, the exact conditions under which BigQuery can prune, how to choose columns from your query patterns, and how to estimate the saving before you change anything. It ends with migration steps for a live table and a checklist. The wider BigQuery architecture, including Dremel, Capacitor and slots, is covered in BigQuery architecture.

Advertisement

The cost model in one paragraph

Under on-demand pricing BigQuery bills the bytes a query processes, at a list price of $6.25 per TiB at the time of writing, with the first TiB each month free. Bytes processed are the sum over referenced columns of their stored size in the rows that survive pruning. Charges round up to the nearest MB, with a minimum of 10 MB per table referenced and per query, so tiny queries over many tables are not free. Under capacity pricing you pay for slots instead, but bytes read still drive slot time, so the same layout work makes queries faster and frees capacity. Prices change; check Google's pricing page before you budget.

Partitioning: splitting a table into independent pieces

A partitioned table is divided into partitions that BigQuery stores and tracks separately. There are three kinds.

  • Time-unit column partitioning on a DATE, TIMESTAMP or DATETIME column. TIMESTAMP and DATETIME columns can be partitioned by hour, day, month or year; DATE columns by day, month or year.
  • Ingestion-time partitioning, where BigQuery assigns rows by the time they arrived and exposes the _PARTITIONTIME and _PARTITIONDATE pseudocolumns.
  • Integer-range partitioning on an integer column with a start, end and interval, written with RANGE_BUCKET.

Rows whose partitioning column is NULL go to a partition named __NULL__, and integer values outside the declared range go to __UNPARTITIONED__. The number of partitions per table is capped; the limit was 10,000 at the time of writing, which is about 27 years of daily partitions but only about 416 days of hourly ones. Check the quotas page for the current figure. A table's partitioning is fixed when it is created; you cannot convert an unpartitioned table in place.

Two table options make partitioning operationally useful. partition_expiration_days deletes partitions older than the given age automatically, which is the cheapest retention policy available. require_partition_filter rejects any query that does not filter on the partitioning column, which stops an accidental full scan from reaching production.

CREATE TABLE analytics.events (
  event_ts    TIMESTAMP NOT NULL,
  tenant_id   STRING    NOT NULL,
  event_type  STRING,
  user_id     STRING,
  payload     JSON
)
PARTITION BY DATE(event_ts)
CLUSTER BY tenant_id, event_type
OPTIONS (
  partition_expiration_days = 730,
  require_partition_filter  = TRUE
);
Advertisement

Clustering: sorting rows so blocks can be skipped

Clustering sorts the data in each partition, or in the whole table if it is not partitioned, by up to four columns. BigQuery stores the data in blocks and keeps metadata about the range of values each block holds. At query time it compares your filters against that metadata and reads only blocks that could match. Google calls this block pruning.

Clustering columns must be top-level and non-repeated, of type STRING, INT64, NUMERIC, BIGNUMERIC, BOOL, DATE, DATETIME, TIMESTAMP, GEOGRAPHY or RANGE. Their order matters: data is sorted by the first column, then the second within it, and so on, so a filter benefits most when it includes the first clustering column. A filter only on the third column skips far fewer blocks.

New data arrives unsorted relative to what is already stored, so BigQuery reclusters in the background. Automatic reclustering has no effect on your query capacity and is not billed as query work. Because the number of blocks that will be skipped is only known when the query runs, a dry run against a clustered table reports an upper bound rather than an exact estimate.

Partition pruning, then block pruning: what a filtered query actually readsWHERE DATE(event_ts) BETWEEN '2026-09-24' AND '2026-09-30' AND tenant_id = 't-417'2024-10-01skipped...skipped2026-09-23skipped2026-09-24read...read2026-09-30readPartition pruning: only 7 daily partitions remainInside one partition: blocks sorted by tenant_idt-001..t-180skippedt-401..t-433readt-434..t-610skippedt-181..t-400skippedt-611..t-799skippedt-800..t-999skippedBlock metadata holds min and max per clustering columnBytes billedcolumns referencedx rows in surviving partitionsx fraction of blocks keptPartitioning is known before the query runs.Block pruning is only known at run time,so a dry run reports an upper bound.Partition pruning removes whole partitions using constant filters on the partitioning column.Clustering sorts data inside each partition so blocks can be skipped by their column ranges.
A query on a table partitioned by day and clustered by tenant. Constant date filters remove all but seven partitions before execution; at run time, block metadata removes the blocks whose tenant range cannot match.

When pruning works, and when it does not

Pruning is the whole point, and it is easy to defeat without noticing. Google's documentation on querying partitioned tables gives the rules, summarised here.

  • Filter the partitioning column with a constant expression. WHERE DATE(event_ts) >= '2026-09-24' prunes. Comparing it with a value computed by a subquery over another table is a dynamic expression, and the documentation shows such a query scanning all partitions.
  • Only some functions keep pruning. Functions such as DATE_ADD, DATE_SUB, DATE_TRUNC, TIMESTAMP_TRUNC and TIMESTAMP_SUB on the partitioning column preserve pruning when their other arguments are constant. EXTRACT(MONTH FROM ...) on the column does not.
  • OR defeats it. WHERE DATE(event_ts) = '2026-09-30' OR tenant_id = 't-417' cannot skip any partition, because rows from any date may match the second condition.
  • Keep pseudocolumns alone on one side. For ingestion-time tables, write _PARTITIONTIME >= TIMESTAMP('2026-09-24') rather than wrapping _PARTITIONTIME in an expression.
  • Clustering follows the same spirit. Equality and range filters on the leading clustering columns prune blocks; wrapping the column in a function or filtering only later columns prunes little.

Do not trust a rule of thumb for your exact query shape. Run it with a dry run or look at the bytes processed afterwards, and compare with the same query without the filter. If they are the same, pruning did not happen.

Choosing the layout

Start from the queries, not the data. Collect the filters used by the queries that cost the most, from INFORMATION_SCHEMA.JOBS, and follow this procedure.

  1. If nearly every expensive query filters on a time range, partition by that time column. Choose the coarsest granularity that still lets typical queries skip most data: day for most event tables, month for slowly growing tables, hour only when queries truly look at hours and volume is high.
  2. Check partition size. Google's guidance is that if partitioning leaves roughly less than 10 GB per partition, clustering alone is often the better choice, because many tiny partitions add metadata overhead without much pruning benefit.
  3. Cluster by the columns that appear most often in equality filters after the time filter, highest-value first. Tenant, customer, country and event type are common. Clustering also helps aggregation and joins on those columns.
  4. Avoid partitioning on a column that queries rarely filter. If you must partition by ingestion time for loading reasons, but queries filter on event time, those are different columns and the partitions will not prune.
  5. Turn on require_partition_filter for large tables and set partition_expiration_days to the retention you actually need.

Worked example: an events table

A product analytics table receives about 40 GiB a day and holds two years, about 730 days or 28.5 TiB. It has 20 columns of roughly equal size. The most common expensive query asks for one tenant's events of a few types over the last seven days and reads 4 of the 20 columns.

Unpartitioned, unclustered. The query reads 4/20 of every row: 0.2 x 28.5 TiB = 5.7 TiB, about $35.60 at $6.25 per TiB. Run hourly by a dashboard, that is about $855 a day.

Partitioned by day. Seven partitions of 40 GiB, of which 20% is referenced: 56 GiB, about 0.055 TiB or $0.34 per query. That is roughly a hundredfold saving from one DDL clause.

Also clustered by tenant_id, event_type. If the tenant holds 1% of rows, perfect clustering would read about 0.56 GiB. In practice blocks straddle tenant boundaries and recent data may not yet be reclustered, so expect a few times that, perhaps 1 to 3 GiB. The cost per query falls to around a cent or two, and the query finishes faster because fewer slots are busy.

Daily partitions here are 40 GiB, comfortably above the 10 GB guidance, so partitioning is justified. A table growing at 500 MB a day would be better served by clustering on the time column and tenant, possibly with monthly partitions.

Changing the layout of a live table

Partitioning cannot be added or changed in place, so the standard path is to create a new table with the desired layout from the old one, then swap readers and writers.

-- 1. Build the new layout (bills a full read of the source once)
CREATE TABLE analytics.events_v2
PARTITION BY DATE(event_ts)
CLUSTER BY tenant_id, event_type
OPTIONS (require_partition_filter = TRUE, partition_expiration_days = 730)
AS SELECT * FROM analytics.events;

-- 2. Point writers at events_v2, backfill the rows written during the copy,
--    then repoint readers through a view such as events_current.
CREATE OR REPLACE VIEW analytics.events_current AS
SELECT * FROM analytics.events_v2;

Clustering is more flexible. You can change or remove a table's clustering specification with bq update --clustering_fields=tenant_id,event_type analytics.events or the tables update API, which is useful for tables receiving continuous streaming writes that are hard to swap. Existing data is not reclustered automatically under the new specification; only new data is. To rewrite everything, Google documents running an update that touches every row, such as UPDATE analytics.events SET tenant_id = tenant_id WHERE true, which is billed as a full rewrite, so do it once and deliberately.

Measuring the result

Two views tell you whether the layout is working. INFORMATION_SCHEMA.JOBS records bytes processed and billed for every query, and INFORMATION_SCHEMA.PARTITIONS shows how data is spread across partitions.

-- Most expensive queries against the table in the last 7 days
SELECT user_email, query, total_bytes_processed / POW(1024, 4) AS tib,
       total_slot_ms
FROM `region-us`.INFORMATION_SCHEMA.JOBS
WHERE creation_time >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 7 DAY)
  AND job_type = 'QUERY'
  AND query LIKE '%analytics.events%'
ORDER BY total_bytes_processed DESC
LIMIT 20;

-- Partition sizes: look for tiny partitions, skew and the __NULL__ partition
SELECT partition_id, total_rows, total_logical_bytes / POW(1024, 3) AS gib
FROM analytics.INFORMATION_SCHEMA.PARTITIONS
WHERE table_name = 'events'
ORDER BY partition_id DESC;

Track bytes processed per query for the top query shapes before and after the change. A dashboard filtering on a derived column that is not the partitioning column will show no improvement, and these views make that obvious within a day. BI Engine, which caches hot data in memory, is a separate lever covered in BigQuery BI Engine.

Failure modes

  • Pruning silently off. A view or a BI tool wraps the partitioning column in a function or compares it with a subquery result. Every query scans the table. Check bytes processed, and use require_partition_filter so the worst cases fail loudly.
  • Too many small partitions. Hourly partitioning on a low-volume table hits the partition limit within a little over a year and adds overhead. Use daily or monthly partitions plus clustering.
  • Wrong time column. Partitioning by ingestion time while users filter by event time. Late-arriving data spreads across partitions and queries cannot prune.
  • NULL pile-up. Rows with a NULL partitioning column collect in __NULL__ and are scanned by queries that cannot exclude it. Make the column NOT NULL or default it at load.
  • Leading clustering column rarely filtered. Ordering clustering columns by cardinality instead of by query use. Reorder to match the filters.
  • Expensive DML. UPDATE and MERGE statements that touch many partitions rewrite them. Batch DML by partition and constrain MERGE conditions on the partitioning column.
  • Expiration surprises. partition_expiration_days deletes data that someone needed for a yearly report. Agree retention with data owners before setting it.

Trade-offs

ChoiceGainCost
Partition by dayExact, pre-execution pruning; cheap retention and per-partition DMLFixed at creation; partition count limit
Partition by hourFiner pruning for very high volumeHits the partition limit in about 14 months; many small partitions
Cluster onlyFlexible, can change spec later, no partition limitDry-run estimates are only an upper bound; pruning depends on sort quality
Partition and clusterCoarse plus fine pruningTwo designs to get right; layout changes need a rebuild
require_partition_filterNo accidental full scansAd-hoc users must learn to add the filter

The same principle, laying data out so the engine can skip it, applies in Spark and other engines; compare with partition pruning in Spark.

What to do next

  1. Pull the 20 most expensive queries on your largest tables from INFORMATION_SCHEMA.JOBS and list the columns they filter on.
  2. For each table, check whether a time column is filtered in nearly all of them and whether daily partitions would be at least around 10 GB.
  3. Choose up to four clustering columns, ordered by how often they appear in equality filters after the time filter.
  4. Estimate the saving with the arithmetic above, then prove it by building a copy of a recent slice with the new layout and comparing bytes processed.
  5. Rebuild the table with CTAS, enable require_partition_filter and a retention-based partition_expiration_days, and swap readers through a view.
  6. Audit views and BI queries for functions or subqueries on the partitioning column.
  7. Re-check bytes processed for the top query shapes a week later, and add an alert for any query that scans more than a set number of TiB.
Key takeaway: BigQuery charges for bytes read, so table layout decides cost. Partition by the time column most queries filter on, at a granularity that keeps partitions large, and filter it with constant expressions so BigQuery can prune before the query starts. Cluster by up to four frequently filtered columns, most important first, so blocks can be skipped at run time. Prove each change with bytes processed from INFORMATION_SCHEMA.JOBS, guard large tables with require_partition_filter, and remember that partitioning is fixed at creation while clustering can be changed later.