cancel
Showing results forย 
Search instead forย 
Did you mean:ย 
Data Engineering
Join discussions on data engineering best practices, architectures, and optimization strategies within the Databricks Community. Exchange insights and solutions with fellow data engineers.
cancel
Showing results forย 
Search instead forย 
Did you mean:ย 

Slow Running SQL Query

Sherbo
New Contributor II

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.

  • fact_profitability_actuals โ€” the main fact table, ~1.4B rows in scope after filtering to 2 fiscal years 
  • dim_category_mapping โ€” small mapping table (brand โ†’ category key)
  • dim_material โ€” ~950K rows, used to filter the fact table down to a specific product brand via an inner join on material key
  • dim_customer โ€” ~11M row customer dimension, used only as a fallback lookup: a small percentage of fact rows are missing a "customer level" attribute directly, so I left-join this dimension (filtered down to just the customers that actually need it via a semi-join first) to backfill it via COALESCE
  • Joins/conditions
  • INNER JOIN fact โ†’ material dimension (on material key) โ€” filters fact rows to the target brand.
  • LEFT JOIN fact โ†’ customer dimension (on customer key, only when the fact-level attribute is NULL) โ€” backfills a missing attribute, doesn't filter rows.
  • WHERE filter on fiscal period range (2 years), confirmed to prune partitions correctly.
  • Final GROUP BY on 9 columns, aggregating volume, a 65-column additive SUM, and a profit calc (SUM(a) - SUM(b) - SUM(c)).

 ```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 REPLIES 4

aayush_410
New Contributor III

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.

Aayush Sharma

ThiamLee
Contributor

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.

anuj_lathi
Databricks Employee
Databricks Employee

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.

What's making it slow

  1. Two scans of the fact table. customers_needing_lookup reads the fact table just to find the NULL-level customers, then the main query reads it again.
  2. The left join runs on 1.4B rows. The 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.
  3. Brand filter via join only. The material join shrinks the data, but unless the fact table's file layout lines up with material, you still read most files in the date range before the join discards rows.
  4. Non-deterministic dedupe. 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.

Rewrite: aggregate before you join

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:

Full query

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;

Why this shape is faster

  • One fact scan instead of two. The customers_needing_lookup CTE is gone. The NULL-level customers fall out of the same pass via lookup_customer.
  • The join now happens after aggregation. Dimension joins run against the pre-aggregated fact_pre instead of 1.4B wide rows, so the big shuffle on the customer join disappears.
  • Re-aggregation is safe. Every measure is a plain SUM. The ratios (NNR/UC, GP/UC) are computed only in the last SELECT, from the summed totals, so

Anuj Lathi
Solutions Engineer @ Databricks

Aravind_Reddy
Databricks Partner

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.