Reading and writing DataFrames

kept in this browsersaved to your account

How spark.read and DataFrameWriter load and persist data on Databricks, and why saving to a Unity Catalog table beats saving to a path.

What it is

spark.read builds a DataFrame from files or an external source; df.write persists a DataFrame somewhere. Databricks defaults everything to Delta Lake (see Delta Lake, the lakehouse table format), so spark.read.parquet(...) and friends exist mostly for reading data that arrived from outside the platform, not for your own tables.

In SQL, the equivalent is read_files, a table-valued function that reads a directory of files with the same format options as the Python reader.

Why it exists

Bronze ingestion needs to read whatever format a source system produces — CSV exports, JSON from an API, Parquet from another warehouse — while everything you write for silver and gold should land as Delta so downstream tools get ACID guarantees, schema enforcement, and time travel. spark.read/spark.write cover both jobs with one API, switching behavior through a format(...) call and a handful of options instead of a different library per file type.

How it works

Reading

spark.read.format("csv"|"json"|"parquet"|"delta").options(...).load(path) reads files directly. spark.table("catalog.schema.table") (or the shorthand spark.read.table(...)) reads a table already registered in Unity Catalog by name — this is what you should reach for once data is past bronze, since it doesn’t require knowing the storage path.

Writing: save versus saveAsTable

df.write.save(path)df.write.saveAsTable("catalog.schema.table")
Registers a table in Unity CatalogNoYes
Addressed byStorage pathThree-level name
Typical useOne-off files, external interchangeAnything other pipelines or users query

On Databricks, saveAsTable is almost always the right call: it makes the output governed, discoverable, and queryable from SQL without anyone needing to know where the files live. Writing to a bare path produces data nobody but you can find (see Managed and external tables for the managed/external distinction that still applies once a table is registered).

Save modes

.mode(...) controls what happens if the target already has data:

ModeBehavior
appendAdd rows to what exists
overwriteReplace existing data entirely
error / errorifexists (default)Fail if the target already exists
ignoreDo nothing if the target already exists

Schema evolution and partitioning

overwrite fails by default if the DataFrame’s schema doesn’t match the existing table. Two options relax that: .option("mergeSchema", "true") adds new columns instead of failing, and .option("overwriteSchema", "true") replaces the table’s schema outright — use it deliberately, since it can silently drop columns that aren’t in the new DataFrame.

.partitionBy("column") writes separate directories per partition value. It sounds like free performance but usually isn’t the right call on Delta tables: partitioning by a low-cardinality column you always filter on (like country) can help, but partitioning by something high-cardinality (like event_date at hourly grain, or a customer ID) creates too many small files and hurts more than it helps. Delta’s liquid clustering and file-level statistics do most of what manual partitioning used to do — reach for explicit partitionBy only when you have a specific, measured reason.

Example

CREATE TABLE shop.silver.orders
USING DELTA
AS SELECT * FROM read_files(
  '/Volumes/shop/bronze/orders_csv',
  format => 'csv',
  header => true
);
raw = (
    spark.read
    .format("csv")
    .option("header", "true")
    .option("inferSchema", "true")
    .load("/Volumes/shop/bronze/orders_csv")
)

(
    raw.write
    .format("delta")
    .mode("overwrite")
    .option("mergeSchema", "true")
    .saveAsTable("shop.silver.orders")
)

# Downstream code reads by name, not by path.
orders = spark.table("shop.silver.orders")

Common mistakes

  • Using .save(path) for tables that other people or jobs need: they end up with no discoverable name, no lineage, no grants in Unity Catalog.
  • Forgetting that errorifexists is the default mode: a rerun of a notebook fails with a confusing error instead of appending or overwriting.
  • Adding mergeSchema everywhere out of habit: it hides real schema drift (a renamed or dropped source column) instead of surfacing it.
  • Partitioning a table “for performance” without checking whether the query patterns and cardinality actually justify it — too many small files makes reads slower, not faster.
  • Reading with spark.read.load(path) when the data is Delta and already a registered table: spark.table(...) is simpler and doesn’t depend on the physical path staying put.

Where this sits

Resources

1All resources
Report a problem with this page
What kind of problem?

Reports about "Reading and writing DataFrames" go to the maintainer, not to a public thread.

Suggest a resource
What kind?

Nothing appears on the site automatically. A person reads every suggestion, checks the link and writes the note that goes with it.