Skip to content
Everyone knows 6 min read · I/O 2 practice problems ↓

Reading and writing data

Load CSV, JSON and Parquet with the right options and schema, and write with the right save mode and layout.

You will learn

  • How to read CSV, JSON and Parquet, and the options that matter for each
  • Why you should pass a schema instead of inferring one
  • What the four save modes do, and which one is dangerous
  • How partitionBy lays out files on disk, and how to overwrite one partition safely

Read first

Nothing. This lesson starts from scratch.

TL;DR Read with spark.read.format(...).option(...).schema(...).load(path) and write with df.write.mode(...).partitionBy(...).format(...).save(path). Pass an explicit schema for CSV and JSON, prefer Parquet or a table format for anything you read twice, and never use mode("overwrite") without knowing whether it replaces the whole table or one partition.

The shape of every read and write

Every source and sink in Spark goes through the same two builders. Reading is lazy: load only reads file footers or a sample to learn the schema; the data is read when an action runs. Writing is an action: save starts a job immediately.

PySparkSpark SQL · Read and write
orders = (spark.read
    .format("csv")
    .option("header", True)
    .schema("order_id INT, customer STRING, amount DOUBLE, order_date DATE")
    .load("s3://raw/orders/"))

(orders.write
    .mode("overwrite")
    .partitionBy("order_date")
    .format("parquet")
    .save("s3://clean/orders/"))
CREATE TABLE clean_orders
USING parquet
PARTITIONED BY (order_date)
AS SELECT * FROM csv.`s3://raw/orders/`;

Shortcuts like spark.read.csv(path, header=True), spark.read.parquet(path) and df.write.parquet(path) are the same builders with the format filled in. saveAsTable("db.orders") writes and registers the table in the catalog, so others can query it by name.

The formats and their options

FormatOptions you will actually useNotes
CSVheader, sep, quote, escape, multiLine, nullValue, dateFormat, modeEverything is a string unless you give a schema. Slowest to read; no column pruningcolumn pruning: Reading only the columns a query uses, which columnar files like Parquet make possible. Learn more →.
JSONmultiLine (one document per file instead of one per line), mode, dateFormatDefault is JSON Lines: one object per line. Nested objects become structs.
ParquetParquet: A columnar file format: values of each column are stored together, so queries read only the columns they need. Learn more →mergeSchemaColumnar, compressed, typed. The schema is stored in the file, so no inference cost.
ORCsimilar to ParquetColumnar; common in Hive-era stacks.
Delta / IcebergversionAsOf, timestampAsOfParquet files plus a transaction logtransaction log: The ordered list of commits that defines which files make up a Delta table at each version. Learn more →. See the Lakehouse track.
JDBCurl, dbtable or query, partitionColumn, lowerBound, upperBound, numPartitionsWithout the partitionpartition: A chunk of a DataFrame's rows. Spark processes each partition as one task, so partitions decide how much work runs in parallel. Learn more → options, the whole table is read by one tasktask: The work for one partition in one stage, run on one CPU core. Learn more →.

Inferring vs passing a schema

inferSchema=True on CSV makes Spark read the data an extra time just to guess types, and the guess depends on the data it sees: a column of zip codes becomes an integer and loses its leading zeros, and a column that is empty in today's file becomes a string. JSON inference scans everything too. Passing the schema is faster and makes the job behave the same every day.

PySpark · Two ways to write a schema
from pyspark.sql.types import StructType, StructField, IntegerType, StringType, DoubleType

# A DDL string: short and readable
schema = "order_id INT, customer STRING, amount DOUBLE"

# The same schema as objects, useful when you build it in code
schema = StructType([
    StructField("order_id", IntegerType(), nullable=False),
    StructField("customer", StringType()),
    StructField("amount", DoubleType()),
])
orders = spark.read.schema(schema).option("header", True).csv("s3://raw/orders/")

Bad records

CSV and JSON readers have three mode settings for rows that do not fit the schema:

modeWhat happens to a bad rowUse it when
PERMISSIVE (default)Bad fields become null. Add columnNameOfCorruptRecord to keep the raw line in a column.You want to load everything and quarantine bad rows afterwards
DROPMALFORMEDThe row is silently droppedRarely: you lose data without knowing how much
FAILFASTThe job fails on the first bad rowThe file must be perfect, for example a contract with a partner
Tip: in PERMISSIVE mode, count the rows where the corrupt-record column is not null and alert when it is above zero. Silent nulls are the most common cause of "the numbers look wrong" incidents.

Save modes

modeIf data already exists
errorifexists (default)Fail
appendAdd the new files next to the old ones. Running the job twice duplicates the data.
overwriteReplace the data. With partitionBy, it replaces every partition unless dynamic overwrite is on.
ignoreDo nothing, silently

partitionBy and overwriting one partition

partitionBy("order_date") writes one folder per value, such as order_date=2025-03-01/. Readers that filter on order_date skip the other folders entirely (partition pruning). Pick a low-cardinalitycardinality: How many distinct values a column has. User ids are high cardinality; country codes are low. Learn more → column that queries filter on; partitioning by a user id creates millions of folders of tiny files.

Step 1 · the table

The table holds 365 date folders. Today's job recomputes only 2025-03-01 because late data arrived.

Step 2 · static overwrite

With the default spark.sql.sources.partitionOverwriteMode=static, mode("overwrite") deletes all 365 folders and writes one. A year of data is gone.

The command

fixed_day.write.mode("overwrite").partitionBy("order_date").parquet(path)   # deletes every other day

Step 3 · dynamic overwrite

With dynamic, Spark replaces only the partitions present in the DataFrame being written. The other 364 folders are untouched.

The command

spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
fixed_day.write.mode("overwrite").partitionBy("order_date").parquet(path)
INSERT OVERWRITE TABLE orders PARTITION (order_date)
SELECT * FROM fixed_day;

Step 4 · with Delta

Delta Lake makes the intent explicit with replaceWhere, and checks that every row written matches the condition.

The command

fixed_day.write.format("delta").mode("overwrite") \
    .option("replaceWhere", "order_date = '2025-03-01'").save(path)

How many files a write produces

A write produces at least one file per task, and with partitionBy one file per task per partition value it holds. 200 tasks writing 30 dates can produce 6,000 files. repartition("order_date") before the write gives one file per date; .option("maxRecordsPerFile", n) caps the size of each file. The Spark internals lesson on small files covers this in depth.

Common mistakes

inferSchema in production

An extra pass over the data, and types that change when the data changes. Pass a schema.

mode("overwrite") with partitionBy and the default settings

Wipes every partition, not just the ones you wrote. Turn on dynamic overwrite or use Delta replaceWhere.

append in a job that can be retried

A retry writes the same rows again. Make the write idempotentidempotent: Safe to run again: running a step twice gives the same result as running it once. Learn more →: overwrite the partition, or MERGE.

Reading a large JDBC table without partition options

One task, one connection, and the job crawls. Set partitionColumn, bounds and numPartitions.

Key takeaways

  • Reading is lazy; writing starts a job.
  • Pass a schema for CSV and JSON instead of inferring it.
  • PERMISSIVE mode plus a corrupt-record column lets you load and then quarantine bad rows.
  • overwrite with partitionBy replaces everything unless dynamic overwrite is on.
  • Partition by a low-cardinality column that queries filter on.

Check yourself

3 questions

1. Which save mode duplicates data when a job is re-run?

Show the answer

append. append adds new files every time the job runs.

2. What does PERMISSIVE mode do with a row that does not fit the schema?

Show the answer

Sets the bad fields to null and keeps the row. It keeps the row with nulls, and can store the raw line in a corrupt-record column.

3. Why is inferSchema=True slow on CSV?

Show the answer

It reads the data an extra time to guess the types. Inference is a separate pass over the data before the real read.

Practice it

Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.

Solve: Combine Batches with Schema Drift →

Keep going

Up next · lesson 2 of 35 · 4 min read
select
Choose, rename and compute columns, and why select beats a chain of withColumn calls.

Related lessons

Primary sources: Data sources · CSV options · DataFrameWriter