4 weeks ago
The following Query is running slowly even though it returns only ~ a million rows. It takes around 25 mins. I am trying Produce a pre-aggregated summary table (grouped by 9 dimension attributes) from a large fact table, to be consumed by a Power BI via native query. The goal is one row per (fiscal period, currency type, company, sales org, customer level, product category/package/size/container) combination, with several summed measures โ one of which is a SUM() of ~65 individual numeric columns added together, plus a couple of simple subtractions for a profit figure.
```SQL
WITH brand_categories AS (
SELECT category_brand_key
FROM catalog_prod.schema_a.dim_brand_category_mappings
WHERE brand_group = 'Brand X'
GROUP BY category_brand_key
),
brand_materials AS (
SELECT
m.material,
m.category_desc,
m.package_desc,
m.pack_size,
m.container_desc
FROM catalog_prod.schema_b.dim_material m
INNER JOIN brand_categories cat
ON cat.category_brand_key = m.category_key
),
-- customers that actually need a level lookup (fact value is null) within our date range
customers_needing_lookup AS (
SELECT DISTINCT f.customer
FROM catalog_prod.schema_c.fact_profitability_actuals f
WHERE f.customer_level IS NULL
AND f.customer IS NOT NULL
AND f.fiscper >= CONCAT(CAST(YEAR(CURRENT_DATE()) - 1 AS STRING), '001')
AND f.fiscper < CONCAT(CAST(YEAR(CURRENT_DATE()) + 1 AS STRING), '001')
),
-- shrunk lookup: only the customers we need, one row per customer (guards against SCD/history fan-out)
customer_lookup AS (
SELECT customer, customer_level
FROM (
SELECT
c.customer,
c.customer_level,
ROW_NUMBER() OVER (PARTITION BY c.customer ORDER BY c.customer) AS rn
FROM catalog_prod.schema_b.dim_customer c
INNER JOIN customers_needing_lookup n
ON n.customer = c.customer
WHERE c.customer_level IS NOT NULL
) t
WHERE rn = 1
),
final_agg AS (
SELECT
f.fiscper,
f.curr_type,
f.comp_code,
f.salesorg,
COALESCE(f.customer_level, c.customer_level) AS customer_level,
m.category_desc,
m.package_desc,
m.pack_size,
m.container_desc,
SUM(f.volume_in_uc) AS volume_in_uc,
SUM(
COALESCE(f.metric_001, 0)
+ COALESCE(f.metric_002, 0)
+ COALESCE(f.metric_003, 0)
+ COALESCE(f.metric_006, 0)
+ COALESCE(f.metric_007, 0)
+ COALESCE(f.metric_008, 0)
+ COALESCE(f.metric_009, 0)
+ COALESCE(f.metric_010, 0)
+ COALESCE(f.metric_011, 0)
+ COALESCE(f.metric_016, 0)
+ COALESCE(f.metric_029, 0)
+ COALESCE(f.metric_051, 0)
+ COALESCE(f.metric_053, 0)
+ COALESCE(f.metric_004, 0)
+ COALESCE(f.metric_005, 0)
+ COALESCE(f.metric_012, 0)
+ COALESCE(f.metric_018, 0)
+ COALESCE(f.metric_041, 0)
+ COALESCE(f.metric_023, 0)
+ COALESCE(f.metric_024, 0)
+ COALESCE(f.metric_025, 0)
+ COALESCE(f.metric_026, 0)
+ COALESCE(f.metric_028, 0)
+ COALESCE(f.metric_014, 0)
+ COALESCE(f.metric_021, 0)
+ COALESCE(f.metric_022, 0)
+ COALESCE(f.metric_013, 0)
+ COALESCE(f.metric_017, 0)
+ COALESCE(f.metric_019, 0)
+ COALESCE(f.metric_020, 0)
+ COALESCE(f.metric_027, 0)
+ COALESCE(f.metric_015, 0)
+ COALESCE(f.metric_030, 0)
+ COALESCE(f.metric_034, 0)
+ COALESCE(f.metric_035, 0)
+ COALESCE(f.metric_036, 0)
+ COALESCE(f.metric_037, 0)
+ COALESCE(f.metric_038, 0)
+ COALESCE(f.metric_039, 0)
+ COALESCE(f.metric_040, 0)
+ COALESCE(f.metric_060, 0)
+ COALESCE(f.metric_062, 0)
+ COALESCE(f.metric_063, 0)
+ COALESCE(f.metric_064, 0)
+ COALESCE(f.metric_112, 0)
+ COALESCE(f.metric_056, 0)
+ COALESCE(f.metric_044, 0)
+ COALESCE(f.metric_031, 0)
+ COALESCE(f.metric_032, 0)
+ COALESCE(f.metric_033, 0)
+ COALESCE(f.metric_042, 0)
+ COALESCE(f.metric_048, 0)
+ COALESCE(f.metric_049, 0)
+ COALESCE(f.metric_057, 0)
+ COALESCE(f.metric_058, 0)
+ COALESCE(f.metric_061, 0)
+ COALESCE(f.metric_111, 0)
+ COALESCE(f.metric_113, 0)
+ COALESCE(f.metric_043, 0)
+ COALESCE(f.metric_045, 0)
+ COALESCE(f.metric_046, 0)
+ COALESCE(f.metric_050, 0)
+ COALESCE(f.metric_055, 0)
+ COALESCE(f.metric_054, 0)
+ COALESCE(f.metric_065, 0)
) AS nnr,
SUM(COALESCE(f.gross_profit_total, 0))
- SUM(COALESCE(f.gross_profit_other, 0))
- SUM(COALESCE(f.net_effect_of_ic_sales_purchases, 0)) AS gp,
SUM(COALESCE(f.net_sales_revenue, 0)) AS NSR
FROM catalog_prod.schema_c.fact_profitability_actuals f
INNER JOIN brand_materials m
ON f.material = m.material
LEFT JOIN customer_lookup c
ON f.customer = c.customer
AND f.customer_level IS NULL
WHERE f.fiscper >= CONCAT(CAST(YEAR(CURRENT_DATE()) - 1 AS STRING), '001')
AND f.fiscper < CONCAT(CAST(YEAR(CURRENT_DATE()) + 1 AS STRING), '001')
GROUP BY
f.fiscper, f.curr_type, f.comp_code, f.salesorg,
m.category_desc, m.package_desc, m.pack_size, m.container_desc,
COALESCE(f.customer_level, c.customer_level)
)
SELECT
fiscper, curr_type, comp_code, salesorg,
category_desc, package_desc, pack_size, container_desc, customer_level,
volume_in_uc,
nnr,
nnr / NULLIF(volume_in_uc, 0) AS `NNR/UC`,
gp / NULLIF(volume_in_uc, 0) AS `GP/UC`,
gp AS `GP`,
NSR
FROM final_agg
```
4 weeks ago
Double scan of the fact table โ customers_needing_lookup and final_agg both scan fact_profitability_actuals with different predicates, so Spark can't reuse the plan. Fix: precompute a small deduped (customer, customer_level) lookup table once (refreshed on dim_customer changes) instead of deriving it live each run. Removes one full pass over the 1.4B-row table.
4 weeks ago
With 1.4B rows, Iโd first check join cardinality and whether filtering happens before the joins. Early reduction of the fact table could make a huge difference, especially before the final GROUP BY.
11 hours ago
Short version: the SQL is logically fine. The cost comes from scanning the 1.4B-row fact table twice and pushing every fact row through the customer left join before any aggregation happens. Aggregate first, join the small result afterwards, and the expensive part mostly goes away.
customers_needing_lookup reads the fact table just to find the NULL-level customers, then the main query reads it again.AND f.customer_level IS NULL in the ON clause doesn't stop the other ~99% of rows from going through the join. If customer_lookup is too big to broadcast, that's a sort-merge join that shuffles the 65-column-wide fact rows. I'd bet this is your dominant cost. Check Query Profile for a large shuffle or spill on that join.material, you still read most files in the date range before the join discards rows.ROW_NUMBER() ... ORDER BY c.customer partitions by customer and orders by the same column, so which row wins for a customer with multiple versions is arbitrary. That's a correctness risk, not just a perf one.SUMs are additive, so you can pre-aggregate the fact table down to a much smaller grain, then join the dimensions and re-aggregate. Only carry customer for rows that actually need a lookup:
WITH brand_categories AS (
SELECT DISTINCT category_brand_key
FROM catalog_prod.schema_a.dim_brand_category_mappings
WHERE brand_group = 'Brand X'
),
brand_materials AS (
SELECT m.material, m.category_desc, m.package_desc, m.pack_size, m.container_desc
FROM catalog_prod.schema_b.dim_material m
JOIN brand_categories cat
ON cat.category_brand_key = m.category_key
),
-- Single scan of the fact table, aggregated BEFORE any dimension join.
-- lookup_customer is populated only for rows that need the level backfill,
-- so the other ~99% of rows collapse to a handful of rows per material.
fact_pre AS (
SELECT
f.fiscper,
f.curr_type,
f.comp_code,
f.salesorg,
f.material,
f.customer_level,
CASE WHEN f.customer_level IS NULL THEN f.customer END AS lookup_customer,
SUM(f.volume_in_uc) AS volume_in_uc,
SUM(
COALESCE(f.metric_001, 0) + COALESCE(f.metric_002, 0) + COALESCE(f.metric_003, 0)
+ COALESCE(f.metric_006, 0) + COALESCE(f.metric_007, 0) + COALESCE(f.metric_008, 0)
+ COALESCE(f.metric_009, 0) + COALESCE(f.metric_010, 0) + COALESCE(f.metric_011, 0)
+ COALESCE(f.metric_016, 0) + COALESCE(f.metric_029, 0) + COALESCE(f.metric_051, 0)
+ COALESCE(f.metric_053, 0) + COALESCE(f.metric_004, 0) + COALESCE(f.metric_005, 0)
+ COALESCE(f.metric_012, 0) + COALESCE(f.metric_018, 0) + COALESCE(f.metric_041, 0)
+ COALESCE(f.metric_023, 0) + COALESCE(f.metric_024, 0) + COALESCE(f.metric_025, 0)
+ COALESCE(f.metric_026, 0) + COALESCE(f.metric_028, 0) + COALESCE(f.metric_014, 0)
+ COALESCE(f.metric_021, 0) + COALESCE(f.metric_022, 0) + COALESCE(f.metric_013, 0)
+ COALESCE(f.metric_017, 0) + COALESCE(f.metric_019, 0) + COALESCE(f.metric_020, 0)
+ COALESCE(f.metric_027, 0) + COALESCE(f.metric_015, 0) + COALESCE(f.metric_030, 0)
+ COALESCE(f.metric_034, 0) + COALESCE(f.metric_035, 0) + COALESCE(f.metric_036, 0)
+ COALESCE(f.metric_037, 0) + COALESCE(f.metric_038, 0) + COALESCE(f.metric_039, 0)
+ COALESCE(f.metric_040, 0) + COALESCE(f.metric_060, 0) + COALESCE(f.metric_062, 0)
+ COALESCE(f.metric_063, 0) + COALESCE(f.metric_064, 0) + COALESCE(f.metric_112, 0)
+ COALESCE(f.metric_056, 0) + COALESCE(f.metric_044, 0) + COALESCE(f.metric_031, 0)
+ COALESCE(f.metric_032, 0) + COALESCE(f.metric_033, 0) + COALESCE(f.metric_042, 0)
+ COALESCE(f.metric_048, 0) + COALESCE(f.metric_049, 0) + COALESCE(f.metric_057, 0)
+ COALESCE(f.metric_058, 0) + COALESCE(f.metric_061, 0) + COALESCE(f.metric_111, 0)
+ COALESCE(f.metric_113, 0) + COALESCE(f.metric_043, 0) + COALESCE(f.metric_045, 0)
+ COALESCE(f.metric_046, 0) + COALESCE(f.metric_050, 0) + COALESCE(f.metric_055, 0)
+ COALESCE(f.metric_054, 0) + COALESCE(f.metric_065, 0)
) AS nnr,
SUM(COALESCE(f.gross_profit_total, 0))
- SUM(COALESCE(f.gross_profit_other, 0))
- SUM(COALESCE(f.net_effect_of_ic_sales_purchases, 0)) AS gp,
SUM(COALESCE(f.net_sales_revenue, 0)) AS nsr
FROM catalog_prod.schema_c.fact_profitability_actuals f
WHERE f.fiscper >= CONCAT(CAST(YEAR(CURRENT_DATE()) - 1 AS STRING), '001')
AND f.fiscper < CONCAT(CAST(YEAR(CURRENT_DATE()) + 1 AS STRING), '001')
AND f.material IN (SELECT material FROM brand_materials)
GROUP BY
f.fiscper, f.curr_type, f.comp_code, f.salesorg, f.material,
f.customer_level,
CASE WHEN f.customer_level IS NULL THEN f.customer END
),
-- One deterministic row per customer. Replace MAX() with your real
-- tie-breaker (e.g. latest valid_from) if levels can differ across versions.
customer_lookup AS (
SELECT customer, MAX(customer_level) AS customer_level
FROM catalog_prod.schema_b.dim_customer
WHERE customer_level IS NOT NULL
GROUP BY customer
),
final_agg AS (
SELECT
p.fiscper,
p.curr_type,
p.comp_code,
p.salesorg,
COALESCE(p.customer_level, c.customer_level) AS customer_level,
m.category_desc,
m.package_desc,
m.pack_size,
m.container_desc,
SUM(p.volume_in_uc) AS volume_in_uc,
SUM(p.nnr) AS nnr,
SUM(p.gp) AS gp,
SUM(p.nsr) AS nsr
FROM fact_pre p
JOIN brand_materials m
ON p.material = m.material
LEFT JOIN customer_lookup c
ON p.lookup_customer = c.customer
GROUP BY
p.fiscper, p.curr_type, p.comp_code, p.salesorg,
COALESCE(p.customer_level, c.customer_level),
m.category_desc, m.package_desc, m.pack_size, m.container_desc
)
SELECT
fiscper, curr_type, comp_code, salesorg,
category_desc, package_desc, pack_size, container_desc, customer_level,
volume_in_uc,
nnr,
nnr / NULLIF(volume_in_uc, 0) AS `NNR/UC`,
gp / NULLIF(volume_in_uc, 0) AS `GP/UC`,
gp AS `GP`,
nsr AS `NSR`
FROM final_agg;
customers_needing_lookup CTE is gone. The NULL-level customers fall out of the same pass via lookup_customer.fact_pre instead of 1.4B wide rows, so the big shuffle on the customer join disappears.NNR/UC, GP/UC) are computed only in the last SELECT, from the summed totals, so
10 hours ago
The core reason this query takes 25 minutes is that it scans the 1.4B-row fact table twice and pushes 1.4B wide rows (with 65 additive metrics) through a LEFT JOIN with dim_customer before performing any GROUP BY aggregation.
Here is a breakdown of the main bottlenecks and how to fix them:
1. Double Fact Scan: customers_needing_lookup and final_agg both read fact_profitability_actuals. You can combine these into a single CTE pass by populating a conditional lookup column (CASE WHEN customer_level IS NULL THEN customer END).
2. Join Before Aggregate: Joining 1.4B rows to customer lookup forces Spark to shuffle high-width rows across nodes. Because SUMs are strictly additive, pre-aggregating the fact table first (fact_pre) reduces billions of rows to a fraction of that volume BEFORE joining dimension tables.
3. Non-Deterministic Order: ROW_NUMBER() OVER (PARTITION BY customer ORDER BY customer) does not guarantee consistent row selection across runs. Replacing this with a deterministic aggregation (e.g. MAX(customer_level) or filtering by valid_from) eliminates windowing overhead and fixes the correctness risk.
i think the query anuj has suggested would help in reducing the query time.