Spark on Databricks is still Apache Spark: the same DataFrame API, the same Catalyst optimizer, the same stages, tasks and shuffles. But code that ran on a self-managed cluster can behave differently there, for good reasons and bad ones. The runtime is a vendor build with its own engine underneath, the cluster may be shared with other users under fine-grained access control, the default table format is Delta, and the unit you pay for is not a machine but a mix of instance hours and Databricks Units.
This article explains those differences from the Spark engineer's point of view: what the runtime adds, how compute types and access modes change which APIs you can call, what happens to a query between your notebook and object storage, how to build and schedule an incremental pipeline, and where performance, cost and reliability problems come from. For the Azure-specific network and workspace setup, see Azure Databricks, in depth; this page stays at the Spark layer.
A Databricks Runtime (DBR) is a versioned bundle: an Apache Spark release plus Databricks patches, Delta Lake, the optional Photon engine, a Python and Scala environment with pinned libraries, and connectors. Long-term support (LTS) versions receive fixes for an extended period, and every runtime's release notes state which Spark version it is based on. Pin the runtime explicitly in job definitions and upgrade deliberately; an unpinned or floating version is the most common cause of a pipeline that changed behaviour with no code change.
What the runtime adds that open-source Spark does not ship: Photon, a native vectorized engine written in C++ that executes supported SQL and DataFrame operators instead of the JVM code path (covered in depth in Spark Photon architecture); the disk cache, which keeps copies of remote Parquet data on executor-local SSD; Auto Loader for incremental file ingestion; Unity Catalog integration for governance; and a number of Delta features that appear on Databricks before or instead of the open-source project. Everything else, including adaptive query execution and the cost-based parts of Catalyst, is the Spark you already know.
Databricks offers several kinds of compute, and choosing one is the first performance and cost decision you make.
| Compute | Use it for | Notes |
|---|---|---|
| All-purpose (interactive) cluster | Notebook development, ad hoc analysis | Shared by users; highest DBU rate; set auto-termination |
| Jobs compute | Scheduled production jobs | Created for a run and torn down after; lower rate than all-purpose |
| Serverless compute | Notebooks, jobs and pipelines without cluster management | Starts in seconds, autoscaled by Databricks; fewer knobs, no custom instance types |
| SQL warehouse | SQL analytics and BI tools | Photon-based SQL endpoint, classic or serverless |
| Instance pools | Reducing classic cluster start time | Keep idle instances warm for reuse |
Classic clusters support autoscaling between a minimum and maximum number of workers. Autoscaling suits bursty, interactive or multi-tenant workloads; for a predictable batch job, a fixed size chosen from measured runs is often cheaper and more stable, because scale-up mid-job means waiting for nodes while tasks queue. Spot or preemptible workers cut cost for fault-tolerant batch work, but keep the driver on on-demand capacity: losing the driver loses the whole application, while losing an executor only re-runs its tasks.
Cluster policies let administrators constrain what users can create: allowed instance types, maximum workers, mandatory tags and auto-termination. They are the main lever against runaway interactive spend.
With Unity Catalog, every classic cluster has an access mode. Dedicated access mode (formerly called single user) assigns the cluster to one user or group; it supports the full Spark API surface, including RDDs and sparkContext. Standard access mode (formerly shared) lets many users share one cluster while Unity Catalog enforces table, row and column permissions for each of them. To make that isolation safe, standard mode runs user code through Spark Connect, the client-server protocol described in Spark Connect, and restricts APIs that bypass the planner. Databricks documents that RDD APIs are not supported in standard mode, and that sc, spark.sparkContext and sqlContext are unsupported for Scala. Serverless compute is built on the same model.
The practical consequence: older code that calls df.rdd.map, sc.parallelize or sets job groups will fail when moved to shared or serverless compute. Rewrite it with DataFrame operations or pandas UDFs, which also tend to be faster, or run it on dedicated compute. Spark Connect also defers analysis, so some errors that used to surface when a DataFrame was defined now surface when an action runs.
Unity Catalog names tables with three levels, catalog.schema.table. Prefer that over paths: access control, lineage and auditing attach to names, and code that reads s3://bucket/path directly bypasses governance and usually needs separately managed credentials.
Follow one query, SELECT region, sum(amount) FROM main.sales.orders WHERE order_date >= '2026-09-01' GROUP BY region, from a notebook to storage.
main.sales.orders against Unity Catalog, checks that the user has SELECT, and obtains short-lived storage credentials scoped to the table's location.Each step suggests a tuning lever: governance by name, a small transaction log through checkpointing, file skipping through clustering, Photon coverage, cache locality, and shuffle sizing.
A common production shape: JSON order events land in cloud storage, you ingest them incrementally into a bronze table, then upsert the latest state of each order into a silver table. Auto Loader tracks which files it has processed in the stream checkpoint, so every run reads only new files.
# Bronze: incremental ingestion with Auto Loader
from pyspark.sql import functions as F
(spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", "/Volumes/main/sales/chk/orders_schema")
.option("cloudFiles.inferColumnTypes", "true") # JSON columns are strings otherwise
.load("/Volumes/main/sales/landing/orders/")
.withColumn("_ingested_at", F.current_timestamp())
.writeStream
.option("checkpointLocation", "/Volumes/main/sales/chk/orders_bronze")
.trigger(availableNow=True) # process everything new, then stop
.toTable("main.sales.orders_bronze"))Then merge into silver. The source must contain at most one row per key, or MERGE fails with a multiple-match error, so deduplicate to the latest event first:
CREATE TABLE IF NOT EXISTS main.sales.orders_silver (
order_id STRING, customer_id STRING, status STRING,
amount DECIMAL(12,2), updated_at TIMESTAMP)
CLUSTER BY (customer_id);
MERGE INTO main.sales.orders_silver AS t
USING (
SELECT order_id, customer_id, status,
CAST(amount AS DECIMAL(12,2)) AS amount,
CAST(updated_at AS TIMESTAMP) AS updated_at FROM (
SELECT *, row_number() OVER (PARTITION BY order_id ORDER BY updated_at DESC) AS rn
FROM main.sales.orders_bronze
WHERE _ingested_at >= current_timestamp() - INTERVAL 1 DAY) WHERE rn = 1
) AS s
ON t.order_id = s.order_id
WHEN MATCHED AND s.updated_at > t.updated_at THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *;The bronze stream stamps each row with _ingested_at so the merge reads only the last day of arrivals, and the inner query casts to the target types and projects exactly the target's columns because UPDATE SET * and INSERT * match columns by name. CLUSTER BY declares liquid clustering, which lets OPTIMIZE co-locate rows by customer without fixed directory partitions. The s.updated_at > t.updated_at guard makes reprocessing harmless: replaying an old batch cannot overwrite newer state. Auto Loader puts fields that do not match the inferred schema in a _rescued_data column instead of dropping them; monitor it, because a non-empty rescued column is an upstream contract change.
For declarative pipelines with built-in data-quality expectations, Databricks offers Lakeflow Spark Declarative Pipelines (formerly Delta Live Tables); the hand-written version above makes each step explicit.
Notebooks are fine for exploration, but production jobs should be defined as code and deployed the same way to every environment. Databricks bundles describe jobs, pipelines and their compute in a databricks.yml file that the Databricks CLI validates and deploys:
bundle:
name: orders_pipeline
variables:
lts_runtime:
description: Pinned LTS runtime string, from databricks clusters spark-versions
node_type:
description: Worker and driver instance type for this cloud
resources:
jobs:
orders_daily:
name: orders_daily
job_clusters:
- job_cluster_key: etl
new_cluster:
spark_version: ${var.lts_runtime} # pin an LTS runtime explicitly
node_type_id: ${var.node_type}
autoscale: {min_workers: 2, max_workers: 8}
tasks:
- task_key: bronze
job_cluster_key: etl
notebook_task: {notebook_path: ./src/bronze.py}
- task_key: silver
depends_on: [{task_key: bronze}]
job_cluster_key: etl
notebook_task: {notebook_path: ./src/silver.py}
targets:
dev: {mode: development, default: true}
prod: {mode: production}Run databricks bundle validate in CI and databricks bundle deploy -t prod from the release pipeline. Both tasks share one job cluster, so the second task does not pay another start-up. databricks clusters spark-versions lists valid runtime strings for your workspace; keep the chosen value in a variable so upgrades are one reviewed change. Set retries and a timeout per task, and alerts on failure and on duration.
Most Spark tuning advice still applies; these levers are specific to, or matter more on, Databricks.
spark.sql.shuffle.partitions and let coalescing shrink it; Databricks also accepts auto for this setting.OPTIMIZE, liquid clustering on common filter and join keys, and predictive optimization where available so maintenance is scheduled for you; VACUUM old files within your time-travel retention needs.collect() and toPandas() on large results are the classic driver out-of-memory; write results to a table instead.Cost on Databricks is roughly cloud instance cost plus DBUs, where the DBU rate depends on compute type, tier and features such as Photon; serverless folds both into one price. Rates vary by cloud, region and contract, so measure rather than assume. The structural savings are consistent: run production on jobs compute or serverless rather than all-purpose clusters, auto-terminate interactive clusters, share job clusters across tasks, use spot workers for retry-tolerant batch work, right-size from measured runs, and tag everything so system billing tables can attribute spend to teams.
Failure modes that recur in production:
sparkContext. Port to DataFrames or use dedicated mode.pip install of unpinned versions changes behaviour. Pin dependencies in the job definition..rdd, sparkContext and sc. before moving to standard access mode or serverless.catalog.schema.table names under Unity Catalog.