Databricks Performance and Optimization: Getting Spark and Delta Lake Working For You
Databricks performance problems tend to cluster around three sources: small files, poor partitioning/Z-ordering, and cluster configurations that don't match the actual workload shape. Here's the practical breakdown.
1. Small files are the silent killer
Streaming ingestion and frequent small batch writes create the "small file problem" — thousands of tiny Parquet files that force Spark to spend more time on file listing and task scheduling overhead than actual data processing.
SQL-- Compact small files into larger ones OPTIMIZE sales.orders; -- Check file count and average size before/after DESCRIBE DETAIL sales.orders;
Schedule
OPTIMIZE2. Z-ordering: pick columns based on actual filter patterns, not intuition
Z-order clusters related data together on disk so file-skipping is effective across multiple filter columns at once — but it only helps for the columns you actually specify, and its benefit degrades with too many columns in one Z-order.
SQLOPTIMIZE sales.orders ZORDER BY (customer_id, order_date);
- ▸Limit Z-order to 2-3 columns that are genuinely high-cardinality and frequently used together in filters or joins.
- ▸Combine with partitioning on a coarse column (like date) and Z-order on finer-grained columns within each partition — partitioning alone on high-cardinality columns creates too many small partitions and reintroduces the small-file problem.
3. Partition strategy: coarse-grained, not fine-grained
A common mistake carried over from Hive-era thinking: partitioning by too many columns or by high-cardinality columns (customer_id, for instance). This creates partition explosion — thousands of tiny partition directories, each with its own small-file problem.
- ▸Partition by date at day or month granularity for most fact tables — coarse enough to avoid explosion, fine enough to prune effectively.
- ▸For anything more granular than that, lean on Z-ordering within partitions instead of adding more partition columns.
4. Cluster sizing: match to workload shape, not just data volume
- ▸Autoscaling clusters are the right default for variable workloads (ad hoc notebooks, BI) — set a sensible min/max rather than a fixed size.
- ▸Job clusters (ephemeral, spun up per job and torn down after) are almost always cheaper and cleaner than running production ETL against an always-on interactive cluster — no idle cost, no risk of one job's resource usage affecting another's.
- ▸Check the Spark UI's stage timeline for skewed tasks (a few tasks taking dramatically longer than the rest) before assuming you need more nodes — skew usually means a join key or partition column needs rethinking, not more compute.
5. Caching: use it deliberately, not by default
.cache().persist()Codedf = spark.read.table("sales.orders").filter("order_date >= '2026-01-01'") df.cache() df.count() # materializes the cache
Unpersist explicitly once you're done with a cached DataFrame rather than relying on cluster-level eviction to clean up after you.
6. Photon and runtime version matter more than people expect
Enabling Photon (Databricks' native vectorized query engine) on SQL-heavy and Delta workloads is frequently a meaningful, no-code-change performance win — check whether it's enabled on your SQL warehouses and job clusters before reaching for more nodes. Similarly, staying current on Databricks Runtime versions matters: Delta Lake and Spark optimizer improvements ship regularly, and running an old runtime means missing free performance improvements.
7. Avoid wide transformations you don't need
- ▸Prefer joins explicitly for small dimension tables joined against large fact tables — don't rely on the optimizer's auto-broadcast threshold without checking it matches your actual table sizes.Code
broadcast() - ▸Watch for unnecessary calls that trigger a full shuffle when a narrowerCode
.repartition()would do (e.g., reducing output file count after a filter, where data volume already dropped).Code.coalesce() - ▸on any query you suspect is inefficient — look specifically forCode
explain()(shuffle) stages that seem disproportionate to the actual data movement needed.CodeExchange
Codedf.explain(mode="formatted")
8. Delta table maintenance is not optional
VACUUMMERGEUPDATEDELETESQLVACUUM sales.orders RETAIN 168 HOURS;
Keep the default 7-day retention unless you have a specific Time Travel requirement that needs longer — shorter retention windows reclaim storage faster but reduce how far back you can query historical versions.
The pattern across all of this
Most Databricks performance problems are storage-layout problems (small files, bad partitioning, missing Z-order) wearing a compute-problem costume. Before scaling a cluster up, check
OPTIMIZETags
Get the latest insights
Subscribe for new articles on cloud architecture, data platforms, and engineering leadership.