Data storage and formats: Parquet, Arrow
Be able to choose a storage format for large datasets and understand columnar storage.
Prerequisites
- EData pipelines and ETLrequired
Intuition
Row storage (CSV, JSON) saves one row at a time:
id,name,age,score
1,Ada,15,82
2,Bo,16,95
Columnar storage (Parquet, ORC) saves one column at a time:
id: 1, 2, 3, 4, ...
name: Ada, Bo, Cim, ...
age: 15, 16, 15, ...
Three consequences of that difference:
- Read only what you need.
SELECT avg(age)reads one column instead of the whole file. On wide tables that is a tenfold to hundredfold difference. - Better compression. Values in the same column resemble each other — the same type, often similar sizes, often repetitions. Columnar compression typically gives 5–10× against CSV.
- Worse for individual rows. If you have to read or change row 4 711, all the columns have to be fetched. Columnar formats are built for analysis, not for transactions.
A rule of thumb: CSV for exchange with people, Parquet for everything else that is larger than a few megabytes.
Formal
The formats and what they are for:
| Format | Type | Good at | Bad at |
|---|---|---|---|
| CSV | row, text | readable, works everywhere | no typing, large, slow |
| JSON/JSONL | row, text | nested structures, streaming | even larger, slow |
| Parquet | column, binary | analysis, compression, schema | not human-readable |
| Arrow | column, in memory | zero-copy between tools | not a storage format |
| Delta / Iceberg | Parquet plus a transaction log | versioning, ACID, time travel | more infrastructure |
Arrow is not a file format but an in-memory representation. The point is that pandas, Polars, DuckDB and Spark can share the same block of memory without serialising — which removes one of the biggest costs in data work.
Parquet's three layers of cleverness:
| Mechanism | What it does |
|---|---|
| Column pruning | reads only the requested columns |
| Row-group statistics | min/max per block → skip blocks that cannot match |
| Dictionary encoding | repeated strings are stored as integers plus a dictionary |
The second is called predicate pushdown and is the reason a filtered query against a large Parquet file often runs in a fraction of the time.
Partitioning splits the data into directories by a column:
data/year=2026/month=09/day=20/part-0.parquet
A query about September 2026 then touches only those files. But do not partition too finely — thousands of small files make everything slower (many opens, poor compression). The benchmark is files of 128 MB–1 GB.
Compression:
| Codec | Size | Speed |
|---|---|---|
| snappy | larger | the fastest — the default choice |
| zstd | smaller | nearly as fast; often the best today |
| gzip | the smallest | slow |
Code
import numpy as np, pandas as pd, time
from pathlib import Path
rng = np.random.default_rng(0)
n = 500_000
df = pd.DataFrame({
"id": np.arange(n),
"city": rng.choice(["Malmö", "Lund", "Göteborg", "Umeå"], n), # few unique → dictionary
"age": rng.integers(15, 80, n),
"score": rng.normal(70, 12, n).round(2),
"date": pd.to_datetime("2026-01-01") + pd.to_timedelta(rng.integers(0, 365, n), "D"),
})
out = Path("/tmp/format"); out.mkdir(exist_ok=True)
df.to_csv(out / "d.csv", index=False)
df.to_parquet(out / "d.snappy.parquet", compression="snappy", index=False)
df.to_parquet(out / "d.zstd.parquet", compression="zstd", index=False)
for f in sorted(out.glob("d.*")):
print(f"{f.name:<20} {f.stat().st_size / 1024**2:>7.2f} MB")
def ms(fn):
t = time.perf_counter(); fn(); return round((time.perf_counter() - t) * 1000)
print("read everything, csv ", ms(lambda: pd.read_csv(out / "d.csv")), "ms")
print("read everything, parquet", ms(lambda: pd.read_parquet(out / "d.zstd.parquet")), "ms")
print("read ONE column ",
ms(lambda: pd.read_parquet(out / "d.zstd.parquet", columns=["age"])), "ms")
# ↑ the last is dramatically faster: only one column is read from disk
# Partitioning
df["year"] = df["date"].dt.year
df["month"] = df["date"].dt.month
df.to_parquet(out / "part", partition_cols=["year", "month"], index=False)
# A query about one month touches only that directory
september = pd.read_parquet(out / "part", filters=[("month", "==", 9)])
print(len(september), "rows in September")
# Predicate pushdown with row-group statistics
high = pd.read_parquet(out / "d.zstd.parquet", filters=[("score", ">", 95)])
print(len(high), "rows with a score above 95 — blocks that cannot match are skipped")
Run the code yourself: the difference between «read everything» and «read one column» is the single clearest demonstration of why columnar storage exists.
Mastery means
- Chooses the format according to the use
- Explains the advantages of columnar storage
- Partitions and compresses deliberately
Sign in to do the exercises and build your mastery up.
Sources
- Apache Parquet — dokumentation (Apache-2.0) — Apache-2.0
- Apache Arrow (Apache-2.0) — Apache-2.0
- pandas — User Guide (BSD-3) — BSD-3-Clause