Every Spark SQL query is planned before a single row is read, and the planner has to guess. How many rows survive this filter? Is the join output bigger or smaller than its inputs? Which table is small enough to broadcast? Which of five joins should run first? Spark's answers come from statistics, and the part of the optimizer that uses detailed statistics to compare alternatives is the cost-based optimizer (CBO). It is switched off by default, and most tables carry nothing but a file size, so most clusters plan from very little information.

This article covers which statistics Spark keeps, how to collect them, the formulas that turn them into row estimates, join reordering, and where Adaptive Query Execution takes over. It also covers the behaviour that surprises nearly everyone: writing to a table throws away its statistics. Configuration names, defaults and algorithms were read from the Spark source on 2026-10-03. For the optimizer as a whole, start with the Spark SQL optimizer architecture.

How statistics reach the plan

From ANALYZE to a physical planANALYZE TABLEscan, HLL++, min/maxINSERT / overwritedrops or resets statsCatalog table propsspark.sql.statistics.*File listingsizeInBytes fallbackwritesclearsLogical plan statisticsper operator: rows, size, column statssize onlyFilterEstimationselectivityJoinEstimationoutput rowsJoin reorderDP over inner joinsJoin strategybroadcast or shufflePhysical planfirst stage runs on estimatesAQEreal shuffle sizesstage doneAQE can switch join strategy and coalesce or split partitions,but it never reorders joins: join order is fixed by the plan.
ANALYZE writes statistics into catalog table properties; estimators propagate them through the plan; join order and strategy are chosen from estimates; AQE corrects strategy at stage boundaries but never reorders joins.

Two levels of statistics

Spark keeps two levels of statistics. Every relation has sizeInBytes. It comes from the catalog if someone stored it. Otherwise a non-partitioned data source table is measured from its files, while partitioned tables fall back to spark.sql.defaultSizeInBytes, which defaults to Long.MaxValue so that an unknown table is never broadcast by accident.

The second level is row count and column statistics: per column, the number of distinct values (NDV), min, max, null count, average and maximum length, and optionally a histogram. These exist only if you ran ANALYZE TABLE. They are stored as table properties in the catalog, under keys that start with spark.sql.statistics.

Which level the planner uses depends on one flag. With spark.sql.cbo.enabled=false (the default), statistics flow through the plan by size only, and a filter is assumed to keep its whole input. With CBO on, operators estimate rows and column ranges from the stored statistics, so a selective filter shrinks the estimate and a join's output is estimated from key NDVs.

Collecting and inspecting statistics

Statistics are collected explicitly. The variants differ a lot in cost:

-- Size only: lists files, no scan. Cheap.
ANALYZE TABLE sales.orders COMPUTE STATISTICS NOSCAN;

-- Size and row count: one full scan.
ANALYZE TABLE sales.orders COMPUTE STATISTICS;

-- Column statistics for the columns the planner needs: one scan, with aggregates.
ANALYZE TABLE sales.orders COMPUTE STATISTICS
  FOR COLUMNS customer_id, region_id, order_date, status;

-- Everything (expensive on wide tables).
ANALYZE TABLE sales.orders COMPUTE STATISTICS FOR ALL COLUMNS;

-- One partition's size and row count.
ANALYZE TABLE sales.orders PARTITION (order_date = '2026-10-02') COMPUTE STATISTICS;

NDV is computed with HyperLogLog++, and spark.sql.statistics.ndv.maxError (default 0.05) sets its relative error. Histograms are off by default. Set spark.sql.statistics.histogram.enabled=true before running ANALYZE to build equi-height histograms with spark.sql.statistics.histogram.numBins bins (default 254). Bin edges come from approximate percentiles, controlled by spark.sql.statistics.percentile.accuracy (default 10000). Histograms cost an extra pass over the data, so add them only for columns with skewed filter predicates.

Since Spark 4.0, spark.sql.statistics.updatePartitionStatsInAnalyzeTable.enabled (default false) makes a table-level ANALYZE also refresh per-partition statistics, at extra cost. To inspect what is stored:

DESCRIBE TABLE EXTENDED sales.orders;              -- "Statistics: 81920000000 bytes, 2000000000 rows"
DESCRIBE TABLE EXTENDED sales.orders customer_id;  -- min, max, num_nulls, distinct_count, avg_col_len, histogram
EXPLAIN COST SELECT ...;                           -- Statistics(sizeInBytes=..., rowCount=...) per operator

The settings that matter

These are the switches that matter, with defaults read from the Spark source:

SettingDefaultEffect
spark.sql.cbo.enabledfalseRow and column estimation through the plan
spark.sql.cbo.joinReorder.enabledfalseDynamic-programming join reordering (needs CBO on)
spark.sql.cbo.joinReorder.dp.threshold12Maximum number of joined items the DP will consider
spark.sql.cbo.joinReorder.card.weight0.7Weight of row count versus size when comparing plans
spark.sql.cbo.planStats.enabledfalseFetch row counts and column stats from the catalog for logical plans
spark.sql.cbo.starSchemaDetectionfalseFact and dimension heuristic for star joins
spark.sql.statistics.size.autoUpdate.enabledfalseRecompute table size after data-changing commands
spark.sql.statistics.fallBackToHdfsfalseFor non-partitioned Hive tables without stored stats, read size from the file system
spark.sql.autoBroadcastJoinThreshold10MBLargest estimated side that is broadcast

You do not have to enable everything at once: cbo.enabled alone improves filter and join estimates, and therefore broadcast decisions. Join reordering is a separate opt-in that only acts when every input has a row count.

How estimates are computed

The estimators are simple formulas, and knowing them makes EXPLAIN COST readable.

Equality filters. For col = literal with the literal inside [min, max], selectivity is 1 / NDV. A literal outside the range estimates zero rows. With a histogram, Spark estimates from the bins the literal falls in instead, which matters when a few values dominate.

Range filters. For col < literal on a numeric or date column, Spark interpolates linearly between min and max, so selectivity is roughly (literal - min) / (max - min). Histograms replace that uniform assumption with per-bin densities. Conjunctions multiply selectivities, which assumes the predicates are independent. A comparison between two columns whose min/max ranges partially overlap gets a fixed guess of one third.

Inner equi-joins. JoinEstimation uses the textbook formula T(A join B) = T(A) × T(B) / max(V(A.k), V(B.k)), where V is NDV. It assumes every key on the side with fewer distinct values finds a match. With several join keys, Spark does not multiply the denominators. It takes only the largest max(V), a deliberately conservative choice because join keys are usually correlated. With histograms on both keys, it estimates bin by bin over the overlapping range.

Worked numbers: orders has 2,000,000,000 rows and 50,000,000 distinct customer_id values. customers has 50,000,000 rows. A filter c.segment = 'enterprise' on a column with 5 distinct values leaves 10,000,000 customers, and that filtered side now has at most 10,000,000 distinct keys. The join estimate is 2,000,000,000 × 10,000,000 / max(50,000,000, 10,000,000) = 400,000,000 rows. Without CBO the planner sees only byte sizes and believes the filter kept everything.

Join reordering

Join reordering lives in CostBasedJoinReorder. It acts on a run of inner joins that have join conditions and no join hints. The run must contain more than two items and at most dp.threshold (12 by default), and every item must have a row count. Anything else is left in the order you wrote it.

The search is a dynamic program in the style of the classic System R optimizer. Level 1 holds each table. Level k builds every plan that joins k items by combining two smaller plans, and keeps the cheapest plan for each set of items. Combinations with no join condition between them are skipped, because cartesian products are almost always a mistake. A plan's cost is the total estimated cardinality and size of its intermediate results. Two plans are compared by a weighted geometric mean of the ratio of their row counts and the ratio of their sizes, with card.weight = 0.7 on rows.

# Simplified CostBasedJoinReorder search (Python pseudocode)
def reorder(items, conds, w=0.7, threshold=12):
    if not (2 < len(items) <= threshold) or any(i.row_count is None for i in items):
        return None                                   # keep the written order
    best = {frozenset([i]): Plan(i, cost=(0, 0)) for i in items}
    for level in range(2, len(items) + 1):
        for left, right in pairs_of_disjoint_sets(best, level):
            on = conditions_between(left, right, conds)
            if not on:
                continue                              # no cartesian products
            cand = join_plan(best[left], best[right], on)   # rows, size from JoinEstimation
            key = left | right
            if key not in best or better(cand, best[key], w):
                best[key] = cand
    return best.get(frozenset(items))

def better(a, b, w):
    rel_rows = a.cost_rows / b.cost_rows
    rel_size = a.cost_size / b.cost_size
    return rel_rows ** w * rel_size ** (1 - w) < 1

Worked example: a four-table star query

Take a query over orders (2 billion rows), customers (50 million), products (2 million) and regions (200 rows), filtered to enterprise customers and one product category, and written in the order an analyst typed it:

SELECT r.name, p.category, sum(o.amount)
FROM orders o
JOIN customers c ON o.customer_id = c.customer_id
JOIN regions   r ON c.region_id   = r.region_id
JOIN products  p ON o.product_id  = p.product_id
WHERE c.segment = 'enterprise' AND p.category = 'storage'
GROUP BY r.name, p.category;

Before, with CBO off, Catalyst still pushes both filters down to the scans, and regions is broadcast because its file size alone is tiny. But customers and products are far over 10 MB, and size-only estimates assume a filter keeps everything. So the plan joins in written order: a sort-merge join of all of orders with enterprise customers, then a second sort-merge join of that large intermediate with products.

After running ANALYZE ... FOR COLUMNS on the join and filter columns of all four tables and enabling cbo.enabled and joinReorder.enabled, EXPLAIN COST shows the filtered products estimated at a few thousand rows. Reordering moves it next to orders and broadcasts it, so the stream shrinks to the storage category before it meets customers. One large shuffle remains instead of two. AQE alone might convert the products join to a broadcast at runtime, but only after the big intermediate was already shuffled. Check the plan before trusting it. If the estimate is wrong, you have replaced a slow plan with a broadcast that can run executors out of memory, which is why the next two sections matter. When you need to force a strategy regardless of estimates, see Spark join strategy hints.

Where AQE takes over

Adaptive Query Execution is on by default since Spark 3.2, and it is often described as making statistics unnecessary. It does not. AQE re-plans at stage boundaries using real shuffle output sizes. It can turn a planned sort-merge join into a broadcast join, coalesce small shuffle partitions and split skewed ones (see data skew in Spark).

AQE cannot change anything decided before the first shuffle, and above all AQE does not reorder joins. Join order is fixed in the logical plan, so a bad order chosen from missing statistics stays bad. CBO picks the order and the first strategy; AQE corrects strategy and partitioning as real sizes arrive.

Failure modes

  • Writes erase statistics. In CommandUtils.updateTableStats, a data-changing command on a table that has statistics either replaces them with a fresh size only, when size.autoUpdate.enabled is on and the size changed, or removes them entirely, which is the default. Row counts and column statistics are gone after the next INSERT, and join reordering quietly stops applying. Re-run ANALYZE as the last step of each load.
  • Writes from outside Spark leave stale statistics. If another engine or a raw file copy changes the table's files, Spark does not know, and the old row counts and min/max survive. Range filters on new dates then fall outside the stored max and are estimated at zero rows. A zero estimate invites a broadcast of a table that is actually huge.
  • NDV hides skew. 1/NDV assumes uniform values. If 30 percent of rows have status = 'shipped', the estimate is wildly low. A histogram on that column fixes it.
  • Correlated predicates. Multiplying selectivities assumes independence. city = 'Paris' AND country = 'FR' is estimated far too small.
  • Connectors differ. DataSource V2 tables report their own statistics, and ANALYZE support varies by format. Check EXPLAIN COST for each format you use.

Running it in production

Treat statistics as part of the data product. Add an ANALYZE step to each pipeline that writes a table, limited to the columns that appear in joins, filters and GROUP BY. That keeps the scan affordable. Large partitioned tables usually need table-level column statistics refreshed after each daily load, with partition pruning still doing most of the work at scan time. A small job can report which tables have lost their row counts:

# PySpark: find tables in a schema that have no row count (stats erased or never collected)
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()

def missing_row_counts(schema):
    out = []
    for t in spark.catalog.listTables(schema):
        if t.tableType == "VIEW":
            continue
        rows = spark.sql(f"DESCRIBE TABLE EXTENDED {schema}.{t.name}").collect()
        stats = [r.data_type for r in rows if r.col_name == "Statistics"]
        if not stats or "rows" not in stats[0]:
            out.append(t.name)
    return out

for name in missing_row_counts("sales"):
    spark.sql(f"ANALYZE TABLE sales.{name} COMPUTE STATISTICS FOR COLUMNS customer_id, order_date")

The column list is illustrative; drive it from per-table configuration. Roll CBO out one query family at a time, comparing EXPLAIN COST and runtimes with it off and on.

Trade-offs

Collection cost versus plan quality. Column statistics cost a scan, and histograms an extra pass. For a rarely queried table, AQE alone may be enough; for a fact table behind hundreds of dashboards, the scan pays for itself.

Better estimates versus riskier plans. CBO makes aggressive plans possible: broadcasts and reorderings based on estimates. With fresh statistics that is a large win. With stale statistics it is a new failure mode. If you cannot keep statistics fresh, leave CBO off and rely on AQE and hints.

Planning time. Join reordering grows exponentially with the number of items. Above the threshold of 12 the written order is kept, so order very wide queries deliberately.

What to do next

  1. Run DESCRIBE TABLE EXTENDED on your five most-joined tables and record whether each has a row count.
  2. Pick one slow multi-join query and capture its EXPLAIN COST output and runtime as a baseline.
  3. Run ANALYZE TABLE ... COMPUTE STATISTICS FOR COLUMNS on the join, filter and grouping columns of each table it touches.
  4. Enable spark.sql.cbo.enabled and spark.sql.cbo.joinReorder.enabled for that session, compare the plan and runtime, and check every new broadcast against real table sizes.
  5. Add histograms only on columns where a filter estimate was clearly wrong.
  6. Make ANALYZE the last step of every pipeline that writes those tables, because each write erases the statistics.
  7. Schedule the missing-row-count check, and alert when a hot table loses its statistics.
Key takeaway: Spark plans from statistics, and by default it has only file sizes. ANALYZE TABLE adds row counts and column statistics, CBO turns them into filter and join estimates, and join reordering uses those estimates to search for a cheaper order of inner joins. AQE fixes join strategy and partitioning at runtime but never reorders joins. Writes erase statistics, so make ANALYZE part of every load, check plans with EXPLAIN COST, and enable CBO only where you can keep statistics fresh.