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.