Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Optimize a slow GROUP BY on a big table

SQL · Performance & Internals

Optimize a slow GROUP BY on a big table

Mediumsql-68
group-byperformancepre-aggregationapproximate

Question

How do you speed up a GROUP BY over hundreds of millions of rows?

Solution

First, make it read less data. Then make it do less work per row. Most GROUP BY slowness comes from scanning too many rows and columns, not from the grouping itself.

Reduce what goes in

  • Filter early on the partition column so pruning removes most of the data. A report for the last 30 days should not read 3 years.
  • Select only the columns you need. In a columnar warehouse, each extra column is more bytes read.
  • Join after aggregating when you can. Aggregating the 500-million-row fact table to one row per customer first, then joining to a small dimension, is far cheaper than joining first and grouping later.

Do the work once, not on every query

If many people run the same aggregation, build a summary table or a materialised view at the grain they need:

CREATE TABLE daily_revenue AS
SELECT order_date, region, SUM(amount) AS revenue, COUNT(*) AS orders
FROM orders
GROUP BY order_date, region;

A dashboard that reads 2 million summary rows beats one scanning 800 million. Keep it up to date incrementally.

Cheaper functions

COUNT(DISTINCT user_id) on very high cardinality is expensive, because the engine must track every distinct value across all workers. If an estimate within about 1 to 2 percent is acceptable, use APPROX_COUNT_DISTINCT (available in BigQuery, Snowflake, Databricks and others, built on the HyperLogLog idea). Make sure stakeholders know it is approximate. Do not use it for billing numbers.

Layout

Cluster or sort the table on the grouping and filtering keys, so rows with the same key are stored together.

Look at the profile

If the profile shows spill to disk, the groups do not fit in memory. A skewed group key makes one worker do most of the work, and the fix is to filter or split that hot key, or to pre-aggregate in two stages (partial aggregate first, then final). Check the profile before you try random changes.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext