Parquet is the default storage format for Spark tables, and almost every Spark performance problem that is not a shuffle problem turns out to be a Parquet layout problem. A job that reads 4 TB to answer a query that touches 30 GB, a write that produces 200,000 files of 3 MB each, a filter that should skip most of a table but skips nothing: all of these come from how files were written and how the reader is allowed to use them.

This page explains Parquet from the bytes up and then follows Spark through both directions: how spark.read.parquet turns files into tasks and decides what it can skip, and how df.write.parquet decides how many files, row groups and pages to produce. It ends with schema evolution, timestamp traps, failure modes and a checklist. How decoded columns flow through Spark as ColumnarBatch objects is covered in Spark columnar processing, and how output is committed safely to object stores in Spark S3 optimisation.

Anatomy of a Parquet file

A Parquet file is columnar inside horizontal slices. Rows are grouped into row groups; within a row group, each column is stored contiguously as a column chunk; each column chunk is divided into pages, typically around 1 MB, which are the unit of encoding and compression. At the end of the file sits the footer: the schema, the byte offset of every column chunk, and statistics for every chunk, namely minimum, maximum and null count. A reader opens a file by reading the last 8 bytes to find the footer length, reading the footer, and then seeking straight to the chunks it needs.

Anatomy of one Parquet fileMagic bytes PAR1Row group 0 (rows 0 to N)Column chunk: user_iddictionary page + data pagesColumn chunk: event_tsdata pages (RLE / delta)Column chunk: payloaddata pages, compressedRow group 1 ... row group Ksame columns, next block of rowsFooter (FileMetaData, Thrift)schema, row-group offsets, per-chunk min/max/null counts, encodings, optional column and offset indexes4-byte footer length + magic bytes PAR1Readers fetch the footer first, then only the column chunks and row groups they need.
Rows are sliced into row groups, each row group stores one chunk per column, and the footer at the end indexes all of it.

Two layers of size reduction apply to each page. Encoding exploits the data's structure: dictionary encoding replaces repeated values with small integer ids, run-length and bit-packing encoding compresses runs and small integers, and delta encodings suit sorted numbers and timestamps. Compression (Snappy, Zstandard, Gzip, LZ4) is then applied to the encoded bytes. Encoding is why sorting matters so much: a sorted column has long runs and small deltas, so it encodes to a fraction of its unsorted size before the compressor even starts. The broader theory of columnar layouts is in columnar database architecture.

The read path, step by step

When you call spark.read.parquet(path), Spark does four things before any task reads a data page.

  1. Lists files and resolves the schema. By default Spark takes the schema from one file's footer (or from the table catalog). With mergeSchema on, it reads every footer and merges them, which is why that option is off by default: on a table of 100,000 files it is a distributed job of its own.
  2. Prunes partitions. For directory-partitioned tables such as dt=2026-10-03/, filters on partition columns remove whole directories before files are even listed in detail.
  3. Plans splits. Files are cut into byte ranges of up to spark.sql.files.maxPartitionBytes (128 MB by default), and small files are packed together, each file costing spark.sql.files.openCostInBytes so that thousands of tiny files do not land in one task. A byte range does not align with row groups; each task reads the footer and processes the row groups whose midpoint falls inside its range. A file with one huge row group therefore gives one task real work and leaves the others idle.
  4. Pushes down projection and filters. Only the referenced columns are read. Filters that the Parquet reader can evaluate are pushed into it, where they are compared against row-group statistics.

Inside each task, the vectorized reader (spark.sql.parquet.enableVectorizedReader, on by default) decodes pages directly into column vectors, spark.sql.parquet.columnarReaderBatchSize rows at a time, 4,096 by default. In Spark 4.2 nested types are vectorized too (enableNestedColumnVectorizedReader is on); older releases fall back to the slower row-based reader for structs, arrays and maps, which is a common reason one table reads several times slower than another of similar size.

Filter pushdown and what defeats it

Filter pushdown (spark.sql.parquet.filterPushdown, on by default) is the mechanism that lets Spark skip data inside files. For each row group, the reader checks the pushed predicate against the chunk's min and max. If event_ts >= '2026-10-01' and the row group's maximum is 2026-09-12, the whole row group is skipped without reading a byte of it. Parquet can also use dictionary pages to skip a row group when an equality value is absent from the dictionary, and optional Bloom filters, written per column, to rule out values for high-cardinality equality lookups that min/max cannot help with.

What pushdown cannot do matters as much. Predicates wrapped in functions or casts, such as to_date(event_ts) = '2026-10-03' or a string column compared to an integer, often cannot be pushed. Pushed filters are still re-applied by Spark after the scan, because skipping works at row-group granularity, not row granularity. And statistics only skip anything if values are clustered: an unsorted column whose every row group spans the full value range has min and max that match every predicate.

df = spark.read.parquet("s3a://lake/events/")
q = df.where("event_ts >= TIMESTAMP '2026-10-01 00:00:00' AND country = 'IN'") \
      .select("user_id", "event_ts", "amount")
q.explain("formatted")
# Look in the FileScan node for:
#   PartitionFilters: [...]    -> directories removed at planning time
#   PushedFilters:   [IsNotNull(event_ts), GreaterThanOrEqual(event_ts,...), EqualTo(country,IN)]
#   ReadSchema:      struct<user_id:bigint,event_ts:timestamp,amount:double,country:string>
# If your predicate is missing from PushedFilters, rewrite it so it compares a bare column.

After running the query, compare the Spark UI's scan metrics for files read, bytes read and row groups or rows output with what the predicate should need. If the bytes read are close to the table size, the filter was pushed but the layout gave it nothing to skip.

Worked example: making one filter skip 99 percent

The numbers below are arithmetic on stated assumptions, not a benchmark, but they show the mechanism. An events table holds 1 TB of compressed Parquet: 30 days of data, 8,000 row groups of 128 MB, events arriving roughly in time order but written by jobs that shuffle on user_id. A dashboard filters country = 'IN' for the last day.

  • As written: the shuffle on user_id scatters countries evenly, so every row group contains Indian rows and every row group's min/max on country covers the whole alphabet. Nothing can be skipped by country. If time is not a partition column and rows from different days are mixed by the shuffle, time statistics do not help either: the query reads the needed columns from all 8,000 row groups.
  • Partition by day: writing with partitionBy("dt") removes 29 of 30 directories at planning time, so the query touches about 267 row groups.
  • Sort within partitions by country: adding sortWithinPartitions("country", "event_ts") before the write makes each row group cover a narrow range of countries. If India is 12% of rows, roughly 12% of each day's row groups, plus one or two at each boundary, contain it, so about 35 row groups are read instead of 267.

Sorting also shrinks the files, because a sorted country column is long runs of identical values that dictionary and run-length encoding collapse. The cost is a sort at write time, which is paid once per write against savings on every read.

The write path: files, row groups and partitions

On the write side, each task writes its own files: at least one per task, one per partition directory that the task has rows for, and a new one whenever maxRecordsPerFile is reached. Inside each file, the Parquet writer buffers rows until the row group reaches parquet.block.size (128 MB by default in parquet-mr), then encodes and flushes it, and each column is split into pages of about parquet.page.size (1 MB). Compression is set by spark.sql.parquet.compression.codec, which defaults to Snappy; Zstandard usually gives noticeably smaller files for a modest CPU cost and is worth testing on your own data.

(events
   .repartition("dt")                         # one task per day -> few, large files
   .sortWithinPartitions("country", "event_ts")  # cluster values so statistics prune
   .write
   .mode("overwrite")
   .option("compression", "zstd")
   .option("maxRecordsPerFile", 20_000_000)   # cap file size inside a big partition
   .option("parquet.bloom.filter.enabled#user_id", "true")  # for point lookups by user
   .option("partitionOverwriteMode", "dynamic")  # replace only the days in this DataFrame
   .partitionBy("dt")
   .parquet("s3a://lake/events/"))

The repartition("dt") before partitionBy("dt") is the important line. Without it, every one of, say, 400 tasks may hold rows for every one of 30 days and write 12,000 small files. With it, each day goes to one task and maxRecordsPerFile splits it into reasonably sized files. A heavily skewed day can make one task huge; repartitioning on dt plus a salt column trades a few more files for balanced tasks. The general trade-offs are in coalesce versus repartition.

Dynamic partition overwrite matters for correctness. In the default static mode, mode("overwrite") with partitionBy deletes every existing partition of the table before writing, so a job that recomputes one day silently deletes the other 29.

Schema evolution and merging

Parquet stores the schema in every file, so a table's files can disagree. Adding a nullable column is safe: old files simply lack it and the reader fills nulls, provided the reader's schema includes it, which means either a catalog schema, an explicit .schema(...), or mergeSchema. Renaming a column is not safe in a plain Parquet table, because columns are matched by name; the old data appears as nulls under the new name. Changing a type is the worst case: widening an int to a long is handled by recent readers in some cases, but other changes fail at read time with a conversion error, sometimes only when a task reaches the one old file.

Table formats such as Delta Lake and Apache Iceberg solve this by keeping the schema in table metadata and mapping columns by id rather than name. Spark can also read Parquet field ids (spark.sql.parquet.fieldId.read.enabled, off by default) for files written with them. If you need renames or type changes on a large table, a table format is the right tool rather than a convention layered on raw Parquet.

Timestamps, calendars and other type traps

Timestamps cause more Parquet incidents than any other type. Spark writes timestamps as the legacy INT96 type by default (spark.sql.parquet.outputTimestampType is INT96), for compatibility with old Hive and Impala readers. INT96 is deprecated in the Parquet format and has no min/max statistics in many writers, so time predicates cannot skip row groups. Set the type to TIMESTAMP_MICROS for new tables unless an old consumer needs INT96.

Very old dates cause a second class of failure. Spark 3 moved from the hybrid Julian/Gregorian calendar to the proleptic Gregorian calendar, so dates before 1582 and timestamps before 1900 can mean different instants in files written by Spark 2 or by legacy systems. The rebase modes spark.sql.parquet.datetimeRebaseModeInRead and int96RebaseModeInRead default to EXCEPTION, so such a value fails the read rather than being shifted silently. Choose LEGACY if the files came from the old calendar and CORRECTED if they did not, and only after checking one known value.

Decimals have a similar compatibility switch, spark.sql.parquet.writeLegacyFormat, which only matters when Hive or Impala versions that predate the standard decimal encodings read your output.

Failure modes

  • Small files. Thousands of tiny files mean listing cost, one footer read per file and poor compression. Repartition before writing and compact old partitions.
  • Giant row groups or giant files with one row group. Only one task can process a row group, so parallelism collapses. Keep row groups near the 128 MB default.
  • Executor out-of-memory on write. Each open file buffers a full row group per column. A task that writes to many partition directories at once holds many buffers. Repartition by the partition column so each task writes few directories.
  • Filters that do not prune. Check PushedFilters and the layout before blaming the engine.
  • Schema drift. A column written as int in one file and string in another fails reads at an unpredictable point. Enforce schemas at write time.
  • Silent data loss on overwrite. Static partition overwrite deletes partitions you did not intend to replace.

What to do next

  1. Run explain("formatted") on your three most expensive queries and confirm their predicates appear in PushedFilters.
  2. Inspect a few files with a Parquet metadata tool and record row-group count, row-group size, encodings and whether min/max statistics are present.
  3. Repartition by the partition column and sort within partitions by your most selective filter column before writing.
  4. Set partitionOverwriteMode to dynamic for every job that rewrites a subset of partitions.
  5. Switch new tables to TIMESTAMP_MICROS and test Zstandard compression on a representative day.
  6. If you need renames, type changes or atomic multi-file commits, move the table to a table format such as Delta Lake or Iceberg, and read about the reader interfaces in DataSource V2.
Key takeaway: Parquet performance in Spark is decided at write time. Readers skip work by column projection, partition pruning and row-group statistics, and statistics only help when values are clustered, so repartition by the partition column, sort within partitions by the main filter column, keep row groups near 128 MB and use dynamic partition overwrite. Check PushedFilters in the plan, prefer TIMESTAMP_MICROS over INT96, treat schema changes as a reason to adopt a table format, and handle calendar rebase errors deliberately rather than by flipping a setting.