Parquet — Columnar Format Fundamentals
The dominant columnar file format for analytics. If you do any data work, you write to Parquet by default.
Why columnar matters
Row-oriented (CSV, JSON, traditional RDBMS):
row1: [user_id, name, email, amount, ts]
row2: [user_id, name, email, amount, ts]
row3: ...
Column-oriented (Parquet, ORC):
user_id column: [1, 2, 3, ...]
name column: [Alice, Bob, Carol, ...]
amount column: [10, 20, 30, ...]
For analytics queries (SELECT user_id, SUM(amount) FROM ... GROUP BY user_id), columnar wins because:
- Read only the columns you need. Scanning a 10-column row-store file reads everything; columnar reads just
user_idandamount. - Compression is much better. Similar values clustered together compress 5-10× better than mixed types in row order.
- Vectorized processing. Modern CPUs process columns in SIMD batches.
Typical analytics workload on Parquet vs CSV: 10-100× faster + 5-10× less storage.
Structure of a Parquet file
[ Row Group 1 ]
Column Chunk: user_id
Page 1 (dictionary)
Page 2 (data)
Column Chunk: amount
Page 1
Column Chunk: name
...
Footer (statistics: min/max per column chunk)
[ Row Group 2 ]
...
[ File Metadata ]
Schema
Row group offsets
- File — split into row groups (default ~128 MB).
- Row group — a horizontal slice of rows, with all columns for those rows.
- Column chunk — within a row group, the data for one column.
- Page — within a chunk, the actual encoded bytes (with stats).
- Footer — schema, row group offsets, column chunk statistics.
Each column chunk has min / max statistics per row group. Readers use this for predicate pushdown — skip entire row groups whose stats don’t match the filter.
Predicate pushdown
df = pl.scan_parquet("orders.parquet").filter(pl.col("year") == 2026).collect()
Polars reads the Parquet metadata; for each row group, checks the min/max of year column; if max < 2026 or min > 2026, skips the entire row group. No data read.
This is why Parquet + filters can be 100× faster than CSV + filters. CSV has to read every byte to find matching rows; Parquet skips whole chunks via metadata.
Projection pushdown
df = pl.scan_parquet("orders.parquet").select("user_id", "amount").collect()
Only user_id and amount columns read from disk. If the file has 100 columns, 98 are skipped at the I/O layer.
Combine projection + predicate pushdown for huge wins.
Compression
Parquet compresses each column chunk independently. Codecs:
| Codec | Speed | Ratio | Use |
|---|---|---|---|
snappy (default) |
fast | medium | balanced; default |
gzip |
slow | best | storage-cost-sensitive |
zstd |
medium | best | modern default for many shops |
brotli |
slow | best | rare |
lz4 |
fastest | weakest | speed-sensitive |
uncompressed |
n/a | none | small files; testing |
zstd has become the de facto best balance (compression close to gzip, decompression speed close to snappy). Most production pipelines now write zstd Parquet.
df.write_parquet("file.parquet", compression="zstd")
Encoding
Within a column chunk, Parquet picks an encoding per page:
- Dictionary — for low-cardinality columns (status, country). Stores values once, references via integer.
- Run-length — for sorted / clustered columns.
- Delta — for integers with small deltas (timestamps in order).
- Plain — fallback.
Combined with column-store layout, encodings give Parquet excellent compression for typical structured data.
Row group size
Default: ~128 MB. Tunable per writer.
| Row group size | Trade-off |
|---|---|
| Smaller (~16 MB) | finer-grained predicate pushdown, more metadata overhead |
| Larger (~256 MB) | better compression, less metadata, but coarser pushdown |
For interactive query engines (DuckDB, Trino, Spark): ~128 MB is the sweet spot. For very large datasets where compression matters more than fine-grained skip: ~256 MB.
Schema evolution
Parquet supports adding / removing columns; readers handle missing columns as null.
What works:
- Adding new optional columns.
- Removing columns (readers just skip them).
- Renaming via aliases at the catalog layer (Iceberg, Delta, Hive).
What breaks:
- Changing a column’s type (int → string).
- Required to optional or vice versa without care.
For production: layer a table format (Delta Lake, Iceberg, Hudi) on top of Parquet files. The table format handles schema evolution, transactions, time travel. Plain Parquet alone is just files.
Partitioning on disk
s3://bucket/orders/
year=2026/
month=01/
file1.parquet
file2.parquet
month=02/...
year=2025/...
Hive-style partitioning. Filters on year / month skip whole partitions without reading metadata. Combines with predicate pushdown within the read files.
Best practice: partition by columns commonly filtered. Don’t over-partition (millions of tiny files is worse than few big ones).
df.write_parquet("orders/", use_pyarrow=True, partition_cols=["year", "month"])
Reading and writing
Pandas
df.to_parquet("file.parquet", engine="pyarrow", compression="zstd")
df = pd.read_parquet("file.parquet", columns=["user_id", "amount"])
df = pd.read_parquet("file.parquet", filters=[("status", "=", "completed")])
filters enables predicate pushdown at the read level.
Polars
df.write_parquet("file.parquet", compression="zstd")
df = pl.read_parquet("file.parquet")
ldf = pl.scan_parquet("file.parquet").filter(...).select(...)
Lazy scan + filter is the canonical efficient pattern.
PyArrow directly
import pyarrow.parquet as pq
import pyarrow.dataset as ds
table = pq.read_table("file.parquet")
table = pq.read_table("file.parquet", columns=["user_id"], filters=[("status", "=", "x")])
# Partitioned dataset
dataset = ds.dataset("orders/", format="parquet", partitioning="hive")
scanner = dataset.scanner(columns=["amount"], filter=ds.field("year") == 2026)
for batch in scanner.to_batches():
process(batch.to_pandas())
PyArrow is the underlying library most others use. For streaming / iterative processing, PyArrow’s batch API is cleaner than reading the whole table.
DuckDB
-- SQL on Parquet files
SELECT user_id, SUM(amount) FROM 'orders/*.parquet' GROUP BY user_id;
SELECT * FROM 'orders/year=2026/*.parquet' WHERE amount > 100;
DuckDB reads Parquet directly via read_parquet() (also a glob). Predicate + projection pushdown via the optimizer.
For “ad-hoc analytics on data files” — DuckDB on Parquet is often faster than Pandas + Parquet.
Parquet vs alternatives
| Format | Row vs col | Speed | Schema | Streaming-friendly |
|---|---|---|---|---|
| Parquet | columnar | fast for analytics | rich, evolvable | OK |
| CSV | row | slow, text | none | yes |
| JSON / NDJSON | row | slow, text | weak | yes |
| ORC | columnar | similar to Parquet | rich | OK |
| Avro | row | fast, binary | strong, evolvable | yes (streaming-first) |
| Arrow IPC | columnar (in-memory) | fastest | strong | for IPC, not long-term storage |
- Parquet for analytics warehouse storage.
- Avro for streaming (Kafka, schema registry).
- Arrow for in-process data exchange (not for long-term file storage).
- ORC is comparable to Parquet, more common in Hive / Hadoop ecosystem.
Common gotchas
- Small files problem. Writing tons of tiny Parquet files (per-event) destroys query performance — per-file overhead dominates. Batch writes; use OPTIMIZE / compact passes.
objectdtype roundtrip. Pandas writes mixed-type object columns as bytes or fails. Cast to a real dtype before writing.- Schema drift across files. Two files with different schemas in the same folder → reader confused. Either enforce schema in writer or use a table format (Delta / Iceberg) that tracks schemas centrally.
- Predicate pushdown not happening. Some filters don’t push down (string regex, complex expressions). Inspect:
.explain()in Polars; check query times. - Compression mismatch. Some readers can’t read zstd-compressed Parquet on older versions. Default to snappy if compatibility matters; zstd if speed/ratio matters more.
Interview angle
- “Why is Parquet faster than CSV for analytics?” — columnar layout: only read needed columns (projection pushdown). Min/max statistics per row group enable skipping entire chunks (predicate pushdown). Compression is 5-10× better. Vectorized processing on CPU. Typical: 10-100× faster end-to-end.
- “What’s predicate pushdown?” — readers check row-group min/max stats against the query’s filter; skip row groups whose stats can’t match. No data read for skipped groups. Works in Spark, Polars, DuckDB, PyArrow.
- “How does Parquet handle schema evolution?” — supports adding / removing columns (readers tolerate missing as null). Doesn’t support type changes safely. Production: layer Delta Lake / Iceberg on top for proper schema management.
- “Parquet vs Avro — when each?” — Parquet for analytics (columnar, optimized for column scans + aggregations, large files). Avro for streaming (row-oriented, schema-evolvable, message-by-message). Different use cases; common to use both in one architecture (Avro for Kafka, Parquet for the warehouse).
- “How do you handle the ‘small files problem’ in a Parquet pipeline?” — batch writes (collect events before writing), or schedule a compaction job. In Delta Lake:
OPTIMIZE my_table. In raw Parquet: rewrite small files into bigger ones. - “What compression codec do you pick?” —
zstdis the modern default (close to gzip ratio, close to snappy speed).snappyif compatibility matters;gzipif storage cost dominates.lz4for speed-critical. - “Why partition by date?” — query workloads typically filter by date (recent data, this month, etc.). Partition pruning (skipping whole partition directories) is even cheaper than predicate pushdown. Don’t over-partition; millions of tiny files hurts.