Spark on Databricks, in depth: the runtime, compute types and access modes, the life of a query, jobs as code, tuning and cost

By Sandeep Belgavi · 2026-10-03 · Category: Apache Spark
Advertisement
Control planeworkspace, jobs, Unity CatalogNotebook / job / bundlesubmits workCompute (classic cluster or serverless)DriverCatalyst plan, AQE, UC checksExecutorsPhoton or JVM operatorsLocal SSDdisk cache, shuffle filestasksDatabricks RuntimeApache Spark + Delta Lake + Photon + libraries, versioned and patched togethercredentials, policyCloud object storageDelta tables: Parquet files + _delta_logread / write via vended credentials
Where Spark runs on Databricks. The control plane schedules work and governs access; the driver and executors run a vendor-built runtime on your compute or on serverless compute; tables are Delta files in object storage.

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.

What the Databricks Runtime adds to Spark

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.

Advertisement

Choosing compute

Databricks offers several kinds of compute, and choosing one is the first performance and cost decision you make.

ComputeUse it forNotes
All-purpose (interactive) clusterNotebook development, ad hoc analysisShared by users; highest DBU rate; set auto-termination
Jobs computeScheduled production jobsCreated for a run and torn down after; lower rate than all-purpose
Serverless computeNotebooks, jobs and pipelines without cluster managementStarts in seconds, autoscaled by Databricks; fewer knobs, no custom instance types
SQL warehouseSQL analytics and BI toolsPhoton-based SQL endpoint, classic or serverless
Instance poolsReducing classic cluster start timeKeep 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.

Access modes, Spark Connect and Unity Catalog

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.

The life of a query

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.

  1. The driver resolves main.sales.orders against Unity Catalog, checks that the user has SELECT, and obtains short-lived storage credentials scoped to the table's location.
  2. Delta reads the transaction log, rebuilding the current snapshot from the latest checkpoint plus subsequent JSON commits, and gets the list of live Parquet files with per-file column statistics (see Delta Lake architecture).
  3. Catalyst pushes the date filter down; file statistics let Delta skip files whose maximum order_date is before September. Clustering decides how effective this is.
  4. The physical plan is a scan, partial aggregation, shuffle by region, final aggregation. If Photon is enabled, supported operators run in Photon; unsupported expressions fall back to the JVM path for that part of the plan, which the query profile shows.
  5. Executors read files from object storage, or from the local disk cache if a previous query read them on the same node.
  6. After the shuffle, adaptive query execution (Spark AQE) inspects real partition sizes and coalesces small shuffle partitions or splits skewed ones before the final stage runs.

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.

Worked example: Auto Loader to an idempotent MERGE

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.

Jobs as code

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.

Tuning that matters on Databricks

Most Spark tuning advice still applies; these levers are specific to, or matter more on, Databricks.

Cost and failure modes

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:

What to do next

  1. Inventory jobs that still run on all-purpose clusters and move them to jobs compute or serverless.
  2. Pin an LTS runtime in every job definition and schedule upgrades with a test run.
  3. Grep your code for .rdd, sparkContext and sc. before moving to standard access mode or serverless.
  4. Replace path-based reads with catalog.schema.table names under Unity Catalog.
  5. Define one pipeline as a bundle with dev and prod targets, validated and deployed from CI.
  6. Add liquid clustering on your most filtered table and compare files scanned before and after.
  7. Compare cost per run with and without Photon on your heaviest job.
  8. Read Spark AQE and Delta Lake architecture to understand the two layers most of this tuning acts on.
Key takeaway: Spark on Databricks is Spark plus a managed runtime, Photon, Delta and governance. Pin the runtime, run production on jobs or serverless compute, write DataFrame code that works under standard access mode, address tables by Unity Catalog name, make merges idempotent, define jobs as bundles, and measure cost per run before and after every tuning change.