Menu

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.

  • Beginner
  • 25 min read
  • Updated Oct 2026
On this page
  1. Sample data and helpers
  2. Reading and writing Parquet
  3. What it is and why it is the default
  4. Save modes
  5. Pitfalls
  6. In interviews
  7. Reading and writing CSV and JSON
  8. CSV
  9. JSON
  10. Pitfalls
  11. In interviews
  12. ORC and Avro
  13. ORC
  14. Avro
  15. Choosing a format
  16. Pitfalls and interview angle
  17. Handling corrupt records
  18. What it is and why it matters
  19. Parse modes
  20. Capturing the raw text
  21. What about badRecordsPath?
  22. Pitfalls
  23. In interviews
  24. Reading from JDBC sources
  25. What it is
  26. A plain read runs as one task
  27. Parallel reads with partitionColumn
  28. Writing back
  29. In interviews
  30. Partition discovery
  31. What it is
  32. Reading a sub-folder
  33. Overwriting one partition: dynamic overwrite
  34. Pitfalls
  35. In interviews
  36. Practice questions
  37. 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 enforceSchema is set to false; 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-avro package 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

  • PERMISSIVE without 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. lowerBound and upperBound only 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; query MIN and MAX first.
  • Choose an evenly distributed column. A column with large gaps or a skewed distribution gives uneven partitions.
  • numPartitions is 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.
  • query vs dbtable: query pushes arbitrary SQL (here a GROUP BY) to the database, so only the result crosses the network. It cannot be combined with partitionColumn; for that, use a subquery alias in dbtable, for example "(SELECT ...) AS t".
  • Filter pushdown: simple filters become a WHERE clause in the generated SQL (PushedFilters with a * above). Aggregate, limit and top-N pushdown exist too, controlled by pushDownAggregate, pushDownLimit and related options; check the plan to confirm what was pushed.
  • fetchsize controls 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 of WHERE clauses, 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=01234 is 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) = 2026 when 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-avro package.
  • Capture malformed CSV and JSON rows with columnNameOfCorruptRecord and quarantine them; badRecordsPath is Databricks-only.
  • JDBC reads run as one task unless you set partitionColumn, bounds and numPartitions; bounds set the stride, they do not filter.
  • Partition discovery turns key=value folders into columns and enables pruning; use dynamic partition overwrite to replace only what you rewrite.

By Data Career Hub Editorial · Last reviewed Oct 2026 · All examples run on PySpark 4.2.0 in local mode with Java 21. Avro uses the external spark-avro_2.13:4.2.0 package, downloaded from Maven Central through spark.jars.packages. The JDBC examples use the embedded Apache Derby database whose driver ships with PySpark, so they run without a server; the PostgreSQL variant is shown for reference and was not executed. badRecordsPath was tested on open-source Spark 4.2 and is silently ignored there.

Progress is saved in this browser only. No account needed.

Search
Filter by type