PySpark courseLesson 5 of 8
PySpark course · Lesson 5 of 8
Reading and Writing Data in PySpark: Parquet, CSV, JSON, ORC, Avro and JDBC
Read and write Parquet, CSV, JSON, ORC and Avro in PySpark, capture corrupt records, load JDBC tables in parallel and use partition discovery to prune folders.
On this page
- Sample data and helpers
- Reading and writing Parquet
- What it is and why it is the default
- Save modes
- Pitfalls
- In interviews
- Reading and writing CSV and JSON
- CSV
- JSON
- Pitfalls
- In interviews
- ORC and Avro
- ORC
- Avro
- Choosing a format
- Pitfalls and interview angle
- Handling corrupt records
- What it is and why it matters
- Parse modes
- Capturing the raw text
- What about badRecordsPath?
- Pitfalls
- In interviews
- Reading from JDBC sources
- What it is
- A plain read runs as one task
- Parallel reads with partitionColumn
- Writing back
- In interviews
- Partition discovery
- What it is
- Reading a sub-folder
- Overwriting one partition: dynamic overwrite
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
Every pipeline starts by reading something and ends by writing something, and most production incidents start at those edges: a schema guessed wrongly, a malformed row turned silently into NULLs, a JDBC read that runs on one core, or an overwrite that wiped every partition instead of one. This lesson walks through the reader and writer API for each common format, then the three topics interviewers probe most: corrupt records, parallel JDBC reads and partition discovery.
Sample data and helpers
The setup creates a session with the external Avro package, a small orders DataFrame and two helpers: files() lists what a write produced (hiding the random UUID Spark puts in file names), and plan() prints a physical plan with the temporary path shortened.
import os, tempfile, glob
from pyspark.sql import SparkSession, functions as F
spark = (
SparkSession.builder.master("local[2]").appName("io")
.config("spark.jars.packages", "org.apache.spark:spark-avro_2.13:4.2.0")
.getOrCreate()
)
spark.sparkContext.setLogLevel("ERROR")
base = tempfile.mkdtemp()
orders = spark.createDataFrame(
[(1, "asha", "IN", 120.0, "2026-01-03"), (2, "ben", "UK", 15.0, "2026-01-04"),
(3, "chen", "UK", 35.5, "2026-01-04"), (4, "dara", "IN", 210.0, "2026-01-05")],
"order_id INT, customer STRING, country STRING, amount DOUBLE, order_date STRING",
).withColumn("order_date", F.to_date("order_date"))
import re, io, contextlib
def files(path):
"""List data files under a folder, hiding checksum files and the random UUID in names."""
names = sorted(os.path.relpath(p, path) for p in glob.glob(path + "/**", recursive=True)
if os.path.isfile(p) and not p.endswith(".crc"))
return [re.sub(r"[0-9a-f]{8}-[0-9a-f-]{27}", "<uuid>", n) for n in names]
def plan(df):
"""Print the physical plan with the temporary folder path shortened."""
buf = io.StringIO()
with contextlib.redirect_stdout(buf):
df.explain()
print(re.sub(r"\[file:[^\]]*\]", "[...]", buf.getvalue()).strip())
The general shape of every read and write is the same:
# reading
spark.read.format("<fmt>").schema(...).option("key", "value").load(path)
# writing
df.write.format("<fmt>").mode("<save mode>").partitionBy(...).option("key", "value").save(path)
Shortcuts such as spark.read.parquet(path) and df.write.csv(path) are the same calls with the format filled in.
Reading and writing Parquet
What it is and why it is the default
Parquet is a columnar file format: values of one column are stored together, compressed, and described by statistics (minimum, maximum, null count) in the file footer. It stores the schema inside the file. That gives Spark three big wins:
- Column pruning: only the columns the query needs are read from disk.
- Predicate pushdown: row groups whose statistics cannot match a filter are skipped.
- No schema guessing: types come from the footer, not from parsing text.
Parquet is Spark’s default format (spark.sql.sources.default = parquet) and its default compression codec is Snappy. For a fuller comparison see Parquet vs Avro vs ORC.
pq = os.path.join(base, "orders_parquet")
orders.write.mode("overwrite").parquet(pq)
print(files(pq))
back = spark.read.parquet(pq)
back.printSchema()
plan(back.filter("amount > 100").select("order_id", "amount"))
['_SUCCESS', 'part-00000-<uuid>-c000.snappy.parquet', 'part-00001-<uuid>-c000.snappy.parquet']
root
|-- order_id: integer (nullable = true)
|-- customer: string (nullable = true)
|-- country: string (nullable = true)
|-- amount: double (nullable = true)
|-- order_date: date (nullable = true)
== Physical Plan ==
*(1) Filter (isnotnull(amount#9) AND (amount#9 > 100.0))
+- *(1) ColumnarToRow
+- FileScan parquet [order_id#6,amount#9] Batched: true, DataFilters: [isnotnull(amount#9), (amount#9 > 100.0)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[...], PartitionFilters: [], PushedFilters: [IsNotNull(amount), GreaterThan(amount,100.0)], ReadSchema: struct<order_id:int,amount:double>
Read the plan from the bottom. ReadSchema lists only order_id and amount (column pruning), and PushedFilters shows the amount > 100 filter handed to the Parquet reader (predicate pushdown). The output is a folder, not a file: one part- file per task (two here, because the DataFrame had two partitions), plus an empty _SUCCESS marker written when the job commits.
Notice that the schema came back with every column nullable = true, even though nothing here was declared otherwise: Spark treats file-sourced columns as nullable for safety.
Save modes
| Mode | If the target exists |
|---|---|
errorifexists (default, also error) |
Fail |
append |
Add new files beside the existing ones |
overwrite |
Replace the existing data (see dynamic overwrite under partition discovery) |
ignore |
Do nothing, silently |
for mode in ["append", "ignore"]:
orders.write.mode(mode).parquet(pq)
print(mode, spark.read.parquet(pq).count())
try:
orders.write.parquet(pq)
except Exception as e:
print(type(e).__name__, str(e).split(" Path")[0])
append 8
ignore 8
AnalysisException [PATH_ALREADY_EXISTS]
Save modes take no locks and are not atomic: with overwrite, the existing data is deleted before the new data is written, so a failure halfway leaves the target empty or partial.
Pitfalls
- Appends are not idempotent. Rerunning an append job after a partial failure duplicates data. Plain file writes are also not atomic for readers: a reader listing the folder mid-write can see some new files. Table formats such as Delta Lake fix both; see Delta Lake transactions.
- Schema drift across files. Spark reads the schema from a sample of footers unless you set
mergeSchema(covered in DataFrames and schemas). - Too many tiny files destroy Parquet’s advantages. See the writing efficient output lesson.
In interviews
“Why Parquet over CSV?” Strong answers mention columnar layout, compression, column pruning, predicate pushdown via footer statistics and an embedded schema, and also its weakness: files are immutable, so row-level updates need a table format on top.
Reading and writing CSV and JSON
CSV
CSV has no types and no schema, so every read is a parse. Always give a schema in pipelines, set header, and be explicit about separators, quoting and null markers.
csv_path = os.path.join(base, "orders_csv")
orders.coalesce(1).write.mode("overwrite").option("header", True).csv(csv_path)
print(files(csv_path))
print(open(glob.glob(csv_path + "/part-*.csv")[0]).read())
schema = "order_id INT, customer STRING, country STRING, amount DOUBLE, order_date DATE"
spark.read.schema(schema).option("header", True).csv(csv_path).show()
['_SUCCESS', 'part-00000-<uuid>-c000.csv']
order_id,customer,country,amount,order_date
1,asha,IN,120.0,2026-01-03
2,ben,UK,15.0,2026-01-04
3,chen,UK,35.5,2026-01-04
4,dara,IN,210.0,2026-01-05
+--------+--------+-------+------+----------+
|order_id|customer|country|amount|order_date|
+--------+--------+-------+------+----------+
| 1| asha| IN| 120.0|2026-01-03|
| 2| ben| UK| 15.0|2026-01-04|
| 3| chen| UK| 35.5|2026-01-04|
| 4| dara| IN| 210.0|2026-01-05|
+--------+--------+-------+------+----------+
coalesce(1) produced a single file here only so it could be printed; on real data a single output file means a single writing task.
Real-world CSV is messier. This file uses ; as a separator, a quoted field containing a comma, doubled quotes inside a quoted field, NA for missing values and a European decimal format:
raw = os.path.join(base, "messy.csv")
with open(raw, "w") as f:
f.write('id;name;note;amount\n1;"Smith, J";"said ""hi""";1.234,50\n2;NA;;99\n')
df = (spark.read
.option("header", True).option("sep", ";").option("quote", '"').option("escape", '"')
.option("nullValue", "NA")
.schema("id INT, name STRING, note STRING, amount STRING")
.csv(raw))
df.show(truncate=False)
+---+--------+---------+--------+
|id |name |note |amount |
+---+--------+---------+--------+
|1 |Smith, J|said "hi"|1.234,50|
|2 |NULL |NULL |99 |
+---+--------+---------+--------+
The amount is read as a string on purpose: 1.234,50 is not a valid Spark number, so you clean it with string functions before casting. The default quote character is " but the default escape character is a backslash, so files that escape quotes by doubling them ("", the usual spreadsheet convention) need escape set to " as above. Other options worth knowing: multiLine (quoted values containing newlines), dateFormat / timestampFormat, encoding, comment, and emptyValue.
JSON
Spark’s JSON reader expects JSON Lines by default: one complete JSON object per line. It infers a schema (including nested structs and arrays) by scanning the data, and fields missing from some records become NULL.
json_path = os.path.join(base, "events.jsonl")
with open(json_path, "w") as f:
f.write('{"id": 1, "user": {"name": "asha", "tier": "gold"}, "items": ["a", "b"]}\n')
f.write('{"id": 2, "user": {"name": "ben"}, "items": [], "coupon": "X1"}\n')
events = spark.read.json(json_path)
events.printSchema()
events.select("id", "user.name", "coupon").show()
root
|-- coupon: string (nullable = true)
|-- id: long (nullable = true)
|-- items: array (nullable = true)
| |-- element: string (containsNull = true)
|-- user: struct (nullable = true)
| |-- name: string (nullable = true)
| |-- tier: string (nullable = true)
+---+----+------+
| id|name|coupon|
+---+----+------+
| 1|asha| NULL|
| 2| ben| X1|
+---+----+------+
Inferred columns come back in alphabetical order, and integers are inferred as long. A file holding one pretty-printed JSON array is not JSON Lines; without multiLine every line is a corrupt record:
multi = os.path.join(base, "array.json")
with open(multi, "w") as f:
f.write('[\n {"id": 1, "name": "asha"},\n {"id": 2, "name": "ben"}\n]\n')
print(spark.read.json(multi).columns)
spark.read.option("multiLine", True).json(multi).show()
out = os.path.join(base, "orders_json")
orders.coalesce(1).write.mode("overwrite").json(out)
print(open(glob.glob(out + "/part-*.json")[0]).read())
['_corrupt_record', 'id', 'name']
+---+----+
| id|name|
+---+----+
| 1|asha|
| 2| ben|
+---+----+
{"order_id":1,"customer":"asha","country":"IN","amount":120.0,"order_date":"2026-01-03"}
{"order_id":2,"customer":"ben","country":"UK","amount":15.0,"order_date":"2026-01-04"}
{"order_id":3,"customer":"chen","country":"UK","amount":35.5,"order_date":"2026-01-04"}
{"order_id":4,"customer":"dara","country":"IN","amount":210.0,"order_date":"2026-01-05"}
multiLine files cannot be split, so each file is read by one task. Prefer JSON Lines for large inputs.
Pitfalls
- Inference on JSON reads the whole input to build the schema (tune with
samplingRatio, or better, supply a schema). - NULLs are not written to JSON. The writer omits NULL fields by default (
ignoreNullFields), so a reader of the output may infer a different schema. - CSV headers are not checked against your schema by name unless
enforceSchemais set tofalse; by default columns are matched by position. - Compressed CSV/JSON (
.gz) cannot be split, so one large gzip file is one task. Bzip2 is splittable but slow.
In interviews
Expect “how would you ingest a messy CSV feed?” Cover an explicit schema, separator/quote/escape/null options, a parse mode, capturing corrupt rows and converting to Parquet or a table format as the first step.
ORC and Avro
ORC
ORC is another columnar format, from the Hive ecosystem. It offers the same column pruning and predicate pushdown as Parquet, plus built-in lightweight indexes and stripe-level statistics. In Spark 4 the default ORC codec is zstd (Parquet stays on Snappy), which you can see in the file names below. In practice, choose Parquet unless an existing Hive estate or tool standardises on ORC.
Avro
Avro is a row-based format with a JSON-defined schema embedded in each file. Rows are written one after another, which makes it cheap to append and well suited to messages and record-at-a-time systems; it is the most common format for Kafka payloads with a schema registry. Its strength is schema evolution: a reader schema can differ from the writer schema, and Avro resolves the difference using field names and defaults.
Avro support is an external module in Spark. You add the spark-avro package matching your Spark and Scala versions (here org.apache.spark:spark-avro_2.13:4.2.0, set in the session above) and use format("avro").
orc = os.path.join(base, "orders_orc")
orders.write.mode("overwrite").orc(orc)
print(files(orc))
spark.read.orc(orc).filter("country = 'UK'").orderBy("order_id").show()
avro = os.path.join(base, "orders_avro")
orders.write.mode("overwrite").format("avro").save(avro)
print(files(avro))
spark.read.format("avro").load(avro).orderBy("order_id").show()
['_SUCCESS', 'part-00000-<uuid>-c000.zstd.orc', 'part-00001-<uuid>-c000.zstd.orc']
+--------+--------+-------+------+----------+
|order_id|customer|country|amount|order_date|
+--------+--------+-------+------+----------+
| 2| ben| UK| 15.0|2026-01-04|
| 3| chen| UK| 35.5|2026-01-04|
+--------+--------+-------+------+----------+
['_SUCCESS', 'part-00000-<uuid>-c000.snappy.avro', 'part-00001-<uuid>-c000.snappy.avro']
+--------+--------+-------+------+----------+
|order_id|customer|country|amount|order_date|
+--------+--------+-------+------+----------+
| 1| asha| IN| 120.0|2026-01-03|
| 2| ben| UK| 15.0|2026-01-04|
| 3| chen| UK| 35.5|2026-01-04|
| 4| dara| IN| 210.0|2026-01-05|
+--------+--------+-------+------+----------+
Schema evolution in action: read the same Avro files with a reader schema that drops columns and adds a new channel field with a default. Avro resolves fields by name and fills the default for data written before the field existed.
import json
avro_schema = json.dumps({
"type": "record", "name": "Order",
"fields": [
{"name": "order_id", "type": "int"},
{"name": "customer", "type": ["null", "string"], "default": None},
{"name": "channel", "type": "string", "default": "web"},
],
})
spark.read.format("avro").option("avroSchema", avro_schema).load(avro).orderBy("order_id").show()
+--------+--------+-------+
|order_id|customer|channel|
+--------+--------+-------+
| 1| asha| web|
| 2| ben| web|
| 3| chen| web|
| 4| dara| web|
+--------+--------+-------+
For Avro-encoded Kafka values, the same module provides from_avro and to_avro functions (in pyspark.sql.avro.functions).
Choosing a format
| Need | Format |
|---|---|
| Analytics, scans of a few columns | Parquet (or ORC) |
| Row-level updates, ACID, time travel | A table format over Parquet (Delta Lake, Iceberg, Hudi) |
| Message payloads, append-heavy, evolving schemas | Avro |
| Exchange with spreadsheets or legacy systems | CSV |
| Semi-structured API payloads, landing zone | JSON Lines, converted early |
Pitfalls and interview angle
- Forgetting the
spark-avropackage gives a “failed to find data source: avro” error; the package version must match your Spark and Scala build. - Avro is not columnar, so it is a poor choice for analytical scans.
- Interviewers ask “row vs columnar format, when would you use each?” Answer with access patterns: write-heavy record streams favour row formats, read-heavy analytics favour columnar ones.
Handling corrupt records
What it is and why it matters
When a CSV or JSON value does not fit the schema (a string in a number column, a truncated line, broken JSON), Spark has to decide what to do with the row. The default is to keep going, which is convenient and dangerous: bad data becomes NULLs that look like legitimately missing data. A production pipeline should detect, count and quarantine bad records.
Parse modes
bad = os.path.join(base, "bad.csv")
with open(bad, "w") as f:
f.write("order_id,customer,amount\n1,asha,10.5\n2,ben,abc\n3,chen\n4,dara,40\n")
schema = "order_id INT, customer STRING, amount DOUBLE"
for mode in ["PERMISSIVE", "DROPMALFORMED"]:
print(mode)
spark.read.schema(schema).option("header", True).option("mode", mode).csv(bad).show()
try:
spark.read.schema(schema).option("header", True).option("mode", "FAILFAST").csv(bad).collect()
except Exception as e:
print("FAILFAST:", type(e).__name__)
PERMISSIVE
+--------+--------+------+
|order_id|customer|amount|
+--------+--------+------+
| 1| asha| 10.5|
| 2| ben| NULL|
| 3| chen| NULL|
| 4| dara| 40.0|
+--------+--------+------+
DROPMALFORMED
+--------+--------+------+
|order_id|customer|amount|
+--------+--------+------+
| 1| asha| 10.5|
| 4| dara| 40.0|
+--------+--------+------+
FAILFAST: Py4JJavaError
| Mode | Behaviour | Use when |
|---|---|---|
PERMISSIVE (default) |
Keep the row, set unparseable fields to NULL, optionally keep the raw text in a corrupt-record column | You want to quarantine bad rows |
DROPMALFORMED |
Silently drop bad rows | Rarely: data loss without a trace |
FAILFAST |
Fail the job on the first bad row | Bad input should stop the pipeline |
Capturing the raw text
Add a string column named by columnNameOfCorruptRecord (default _corrupt_record) to your schema. In PERMISSIVE mode Spark stores the whole raw line there for any row that failed to parse. Split the result into a good set and a quarantine set:
captured = (spark.read
.schema("order_id INT, customer STRING, amount DOUBLE, _corrupt_record STRING")
.option("header", True)
.option("mode", "PERMISSIVE")
.option("columnNameOfCorruptRecord", "_corrupt_record")
.csv(bad)
.cache())
captured.show(truncate=False)
good = captured.filter("_corrupt_record IS NULL").drop("_corrupt_record")
quarantine = captured.filter("_corrupt_record IS NOT NULL")
print("good:", good.count(), "bad:", quarantine.count())
+--------+--------+------+---------------+
|order_id|customer|amount|_corrupt_record|
+--------+--------+------+---------------+
|1 |asha |10.5 |NULL |
|2 |ben |NULL |2,ben,abc |
|3 |chen |NULL |3,chen |
|4 |dara |40.0 |NULL |
+--------+--------+------+---------------+
good: 2 bad: 2
Write quarantine to its own location with a load timestamp, and fail or alert when the bad count crosses a threshold. The cache() matters: Spark disallows queries over raw CSV/JSON files that reference only the corrupt-record column, and caching the parsed result first is the documented workaround; it also avoids parsing the file twice for the two counts.
JSON works the same way. Note how a type mismatch ("three" for an int) keeps the other fields, while a broken line gives NULLs everywhere:
jbad = os.path.join(base, "bad.jsonl")
with open(jbad, "w") as f:
f.write('{"id": 1, "amount": 10}\n{"id": 2, "amount": \n{"id": "three", "amount": 5}\n')
spark.read.schema("id INT, amount DOUBLE, _corrupt_record STRING").json(jbad).show(truncate=False)
+----+------+----------------------------+
|id |amount|_corrupt_record |
+----+------+----------------------------+
|1 |10.0 |NULL |
|NULL|NULL |{"id": 2, "amount": |
|NULL|5.0 |{"id": "three", "amount": 5}|
+----+------+----------------------------+
What about badRecordsPath?
badRecordsPath is a Databricks Runtime option: on Databricks it writes bad records and unreadable files as JSON under a given path and carries on. It is not part of open-source Apache Spark. On Spark 4.2 the option is accepted and silently ignored: the read behaves as plain PERMISSIVE and nothing is written. Even on Databricks its own documentation notes that it is non-transactional, and it recommends the corrupt-record column approach instead. If you meet badRecordsPath in an interview, explain both points.
For whole files that cannot be read (for example a truncated Parquet file), open-source Spark has ignoreCorruptFiles and ignoreMissingFiles (as reader options or spark.sql.files.* settings). Both default to false; enabling them skips the file, so log or count what was skipped.
Pitfalls
PERMISSIVEwithout a corrupt-record column loses the evidence.- Declaring the corrupt-record column with a type other than string fails the read.
- Parquet and ORC carry their own schema, so parse modes do not apply to them; type mismatches there surface as schema errors.
In interviews
“How do you handle bad records in Spark?” Strong answer: explicit schema, PERMISSIVE with columnNameOfCorruptRecord, split good and bad, write the quarantine set somewhere durable, alert on thresholds, and use FAILFAST where any bad row should stop the job. Mention that badRecordsPath is Databricks-specific.
Reading from JDBC sources
What it is
Spark’s JDBC data source reads a table or query from any database with a JDBC driver (PostgreSQL, MySQL, SQL Server, Oracle, and so on) into a DataFrame, and can write a DataFrame back. The driver JAR must be on the classpath of the driver and executors (for example via spark.jars.packages). The examples here use Apache Derby in embedded mode, purely because its JDBC driver ships with PySpark and needs no server; the Spark documentation lists Derby support as deprecated, so do not take it as a production choice.
A plain read runs as one task
db = os.path.join(base, "shopdb")
url = f"jdbc:derby:{db};create=true"
props = {"driver": "org.apache.derby.iapi.jdbc.AutoloadedDriver"}
big = spark.range(1, 1001).select(
F.col("id").cast("int").alias("ORDER_ID"),
(F.col("id") % 7).cast("int").alias("STORE_ID"),
(F.col("id") * 1.5).alias("AMOUNT"))
big.write.jdbc(url, "ORDERS", mode="overwrite", properties=props)
whole = spark.read.jdbc(url, "ORDERS", properties=props)
print(whole.count(), whole.rdd.getNumPartitions())
1000 1
Without partitioning options, Spark opens one connection and pulls the whole table through one task. On a large table that is slow and can run one executor out of memory.
Parallel reads with partitionColumn
Give Spark a numeric, date or timestamp column plus lowerBound, upperBound and numPartitions. Spark splits the range into equal strides and issues one query per partition, each with a WHERE clause on that column.
parallel = (spark.read.format("jdbc")
.option("url", url).option("driver", props["driver"])
.option("dbtable", "ORDERS")
.option("partitionColumn", "ORDER_ID").option("lowerBound", 1).option("upperBound", 1000)
.option("numPartitions", 4).option("fetchsize", 500)
.load())
print(parallel.rdd.getNumPartitions())
print(parallel.groupBy(F.spark_partition_id().alias("p")).count().orderBy("p").collect())
pushed = (spark.read.format("jdbc")
.option("url", url).option("driver", props["driver"])
.option("query", "SELECT STORE_ID, SUM(AMOUNT) AS TOTAL FROM ORDERS GROUP BY STORE_ID")
.load())
pushed.orderBy("STORE_ID").show(3)
plan(whole.filter("AMOUNT > 1400"))
4
[Row(p=0, count=251), Row(p=1, count=249), Row(p=2, count=249), Row(p=3, count=251)]
+--------+--------+
|STORE_ID| TOTAL|
+--------+--------+
| 0|106606.5|
| 1|106821.0|
| 2|107035.5|
+--------+--------+
only showing top 3 rows
== Physical Plan ==
*(1) Scan JDBCRelation(ORDERS) [numPartitions=1] [ORDER_ID#665,STORE_ID#666,AMOUNT#667] PushedFilters: [*IsNotNull(AMOUNT), *GreaterThan(AMOUNT,1400.0)], ReadSchema: struct<ORDER_ID:int,STORE_ID:int,AMOUNT:double>
Key points:
- Bounds do not filter.
lowerBoundandupperBoundonly set the stride. Rows outside the range still come back, in the first and last partitions. If the bounds are far from the real min and max, partitions are badly skewed; queryMINandMAXfirst. - Choose an evenly distributed column. A column with large gaps or a skewed distribution gives uneven partitions.
numPartitionsis also the maximum number of concurrent connections. Agree it with the database owner; twenty parallel full scans can hurt a production OLTP database. Read from a replica where possible.queryvsdbtable:querypushes arbitrary SQL (here aGROUP BY) to the database, so only the result crosses the network. It cannot be combined withpartitionColumn; for that, use a subquery alias indbtable, for example"(SELECT ...) AS t".- Filter pushdown: simple filters become a
WHEREclause in the generated SQL (PushedFilterswith a*above). Aggregate, limit and top-N pushdown exist too, controlled bypushDownAggregate,pushDownLimitand related options; check the plan to confirm what was pushed. fetchsizecontrols rows fetched per round trip. Some drivers default to very small values (Oracle) or to fetching everything at once (PostgreSQL without a fetch size), so set it explicitly.- Alternative:
predicates.spark.read.jdbc(url, table, predicates=[...])takes a list ofWHEREclauses, one per partition, useful when no numeric column fits (for example one predicate per month).
The same read against PostgreSQL only changes the URL, driver package and credentials. This was not executed here:
spark = (SparkSession.builder
.config("spark.jars.packages", "org.postgresql:postgresql:42.7.13")
.getOrCreate())
orders_pg = (spark.read.format("jdbc")
.option("url", "jdbc:postgresql://db-replica:5432/shop")
.option("dbtable", "public.orders")
.option("user", "reader").option("password", os.environ["PG_PASSWORD"])
.option("partitionColumn", "order_id").option("lowerBound", 1).option("upperBound", 50_000_000)
.option("numPartitions", 16).option("fetchsize", 10_000)
.load())
Writing back
df.write.jdbc(...) inserts with batched statements (batchsize). Each task writes in its own transaction, so a failed job can leave partial data. Common safe patterns are writing to a staging table and swapping or merging in the database, or making the target idempotent with an upsert keyed on a business key. mode("overwrite") drops and recreates the table unless you set truncate to true.
In interviews
“How do you read a 500 million row table from Postgres efficiently?” Expected answer: partitionColumn with real min/max bounds, a sensible numPartitions the database can handle, fetchsize, pushing filters and projections to the database, reading from a replica, and for repeated loads, incremental extraction by an updated-at column or change data capture instead of full reads.
Partition discovery
What it is
When you write with partitionBy("country", "order_date"), Spark creates one folder per distinct value, in Hive style: country=IN/order_date=2026-01-03/. The partition values are not stored in the data files, only in the folder names. On read, partition discovery walks the folders, turns each key=value into a column, infers its type (spark.sql.sources.partitionColumnTypeInference.enabled, default true) and appends the partition columns at the end of the schema.
part = os.path.join(base, "orders_by_country")
orders.write.mode("overwrite").partitionBy("country", "order_date").parquet(part)
print(files(part)[:4])
p = spark.read.parquet(part)
p.printSchema()
plan(p.filter("country = 'UK'"))
['_SUCCESS', 'country=IN/order_date=2026-01-03/part-00000-<uuid>.c000.snappy.parquet', 'country=IN/order_date=2026-01-05/part-00001-<uuid>.c000.snappy.parquet', 'country=UK/order_date=2026-01-04/part-00000-<uuid>.c000.snappy.parquet']
root
|-- order_id: integer (nullable = true)
|-- customer: string (nullable = true)
|-- amount: double (nullable = true)
|-- country: string (nullable = true)
|-- order_date: date (nullable = true)
== Physical Plan ==
*(1) ColumnarToRow
+- FileScan parquet [order_id#596,customer#597,amount#598,country#599,order_date#600] Batched: true, DataFilters: [], Format: Parquet, Location: InMemoryFileIndex(1 paths)[...], PartitionFilters: [isnotnull(country#599), (country#599 = UK)], PushedFilters: [], ReadSchema: struct<order_id:int,customer:string,amount:double>
The filter on country appears under PartitionFilters: Spark lists only the country=UK folders and never opens the others. This is partition pruning, the main reason to partition. order_date was inferred as a date from the folder names, and partition columns moved to the end of the schema.
Reading a sub-folder
If you point the reader at one partition folder directly, discovery starts there and the country column disappears. Set basePath to keep it:
spark.read.parquet(os.path.join(part, "country=UK")).printSchema()
(spark.read.option("basePath", part).parquet(os.path.join(part, "country=UK"))
.select("order_id", "country").orderBy("order_id").show())
root
|-- order_id: integer (nullable = true)
|-- customer: string (nullable = true)
|-- amount: double (nullable = true)
|-- order_date: date (nullable = true)
+--------+-------+
|order_id|country|
+--------+-------+
| 2| UK|
| 3| UK|
+--------+-------+
Overwriting one partition: dynamic overwrite
By default (spark.sql.sources.partitionOverwriteMode = static), mode("overwrite") on a partitioned path deletes every partition before writing. To reprocess a single day and leave the others untouched, switch to dynamic: only partitions present in the DataFrame being written are replaced.
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
fix = spark.createDataFrame([(3, "chen", "UK", 36.0, "2026-01-04")],
"order_id INT, customer STRING, country STRING, amount DOUBLE, order_date STRING"
).withColumn("order_date", F.to_date("order_date"))
fix.write.mode("overwrite").partitionBy("country", "order_date").parquet(part)
spark.read.parquet(part).orderBy("order_id").show()
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "static")
+--------+--------+------+-------+----------+
|order_id|customer|amount|country|order_date|
+--------+--------+------+-------+----------+
| 1| asha| 120.0| IN|2026-01-03|
| 3| chen| 36.0| UK|2026-01-04|
| 4| dara| 210.0| IN|2026-01-05|
+--------+--------+------+-------+----------+
The UK / 2026-01-04 partition was replaced (order 2 is gone and order 3 has the corrected amount), and both IN partitions survived. With static mode, the IN rows would have been deleted too. You can also set it per write with .option("partitionOverwriteMode", "dynamic").
Pitfalls
- Over-partitioning. Partitioning by a high-cardinality column (user ID, timestamp) creates thousands of folders with tiny files and makes listing slow. Partition by low-cardinality columns that queries filter on, typically a date. Aim for partitions holding at least hundreds of megabytes.
- Type inference surprises. A folder value
zip=01234is inferred as an integer and loses its leading zero. Disable inference or provide a schema that includes the partition column. - Filters on expressions of the partition column (for example
year(order_date) = 2026when the folder is by date) may not prune; filter on the raw column. - Listing cost on object stores. Discovery lists every folder at read time. Table formats and metastores avoid that by keeping the file list in metadata.
In interviews
Expect “how does partition pruning work?” and “how would you rerun one day of a partitioned table?” The strong answers mention Hive-style folders, PartitionFilters in the plan, dynamic partition overwrite, and the cost of over-partitioning.
Practice questions
Your job reads a CSV feed in PERMISSIVE mode and the daily totals are suddenly low. What happened and how would you prevent it?
Some values probably failed to parse (a format change upstream) and became NULL, which aggregates ignore. Add a corrupt-record column to the schema, split good and bad rows, write the bad rows to a quarantine location and alert or fail when the bad count exceeds a threshold. Use FAILFAST where any malformed row should stop the load.
Is badRecordsPath a good way to capture bad records in open-source Spark?
No. It is a Databricks Runtime option; open-source Spark ignores it without an error. Use PERMISSIVE mode with columnNameOfCorruptRecord, and ignoreCorruptFiles only when skipping unreadable files is acceptable and logged.
A JDBC read of a large table takes hours and one task does all the work. Why, and how do you fix it?
Without partitionColumn, lowerBound, upperBound and numPartitions, Spark reads through a single connection in one task. Add those options on an evenly distributed numeric or date column with bounds from the real min and max, set fetchsize, push filters to the database, and keep numPartitions within what the database can serve.
Do lowerBound and upperBound filter the rows returned by a JDBC read?
No. They only decide how the range is split into partitions. Rows below the lower bound go to the first partition and rows above the upper bound go to the last. Use a WHERE in a subquery or a DataFrame filter to restrict rows.
You need to reprocess yesterday’s partition of a table partitioned by date. What setting matters, and what goes wrong without it?
partitionOverwriteMode = dynamic. In the default static mode, mode("overwrite") on the table path deletes all partitions and writes only yesterday’s data, wiping the history. Dynamic mode replaces only the partitions present in the written DataFrame.
When would you pick Avro over Parquet?
For row-oriented, append-heavy workloads such as event messages (often on Kafka with a schema registry), where whole records are written and read and schema evolution with defaults matters. Parquet is better for analytical scans that read a few columns from many rows.
Why does reading a partition sub-folder lose the partition column, and how do you keep it?
Partition discovery treats the path you pass as the root, so key=value folders above it are not seen. Set the basePath option to the table root so Spark still derives the column from the folder names.
Key takeaways
- Every read and write follows
format,schema,option,load/save; give file readers explicit schemas and pick save modes deliberately. - Parquet is the default for analytics because of column pruning, predicate pushdown and an embedded schema; Avro suits row-oriented records with evolving schemas and needs the
spark-avropackage. - Capture malformed CSV and JSON rows with
columnNameOfCorruptRecordand quarantine them;badRecordsPathis Databricks-only. - JDBC reads run as one task unless you set
partitionColumn, bounds andnumPartitions; bounds set the stride, they do not filter. - Partition discovery turns
key=valuefolders into columns and enables pruning; use dynamic partition overwrite to replace only what you rewrite.
Progress is saved in this browser only. No account needed.