Reading and writing data
Load CSV, JSON and Parquet with the right options and schema, and write with the right save mode and layout.
On this page
Show code in
Every code block on the page follows this.
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.
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.
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
| Format | Options you will actually use | Notes |
|---|---|---|
| CSV | header, sep, quote, escape, multiLine, nullValue, dateFormat, mode | Everything 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 →. |
| JSON | multiLine (one document per file instead of one per line), mode, dateFormat | Default 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 → | mergeSchema | Columnar, compressed, typed. The schema is stored in the file, so no inference cost. |
| ORC | similar to Parquet | Columnar; common in Hive-era stacks. |
| Delta / Iceberg | versionAsOf, timestampAsOf | Parquet 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. |
| JDBC | url, dbtable or query, partitionColumn, lowerBound, upperBound, numPartitions | Without 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.
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:
| mode | What happens to a bad row | Use 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 |
DROPMALFORMED | The row is silently dropped | Rarely: you lose data without knowing how much |
FAILFAST | The job fails on the first bad row | The file must be perfect, for example a contract with a partner |
Save modes
| mode | If data already exists |
|---|---|
errorifexists (default) | Fail |
append | Add the new files next to the old ones. Running the job twice duplicates the data. |
overwrite | Replace the data. With partitionBy, it replaces every partition unless dynamic overwrite is on. |
ignore | Do 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
Step 2 · static overwrite
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
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
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
mode("overwrite") with partitionBy and the default settings
append in a job that can be retried
Reading a large JDBC table without partition options
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 questions1. 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.
Keep going
Up next · lesson 2 of 35 · 4 min readselect
Choose, rename and compute columns, and why select beats a chain of withColumn calls.
Related lessons
cast and schemasSpark internals · 6 min read
The small file problemData lake & lakehouse · 4 min read
Partitioning done right
Primary sources: Data sources · CSV options · DataFrameWriter