Skip to content
AI-grafen
EUniversityData handling· about 60 min· evolving, reviewed regularly· verified 2026-09-20· EN

Data storage and formats: Parquet, Arrow

Be able to choose a storage format for large datasets and understand columnar storage.

Prerequisites

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:

  1. 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.
  2. 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.
  3. 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:

FormatTypeGood atBad at
CSVrow, textreadable, works everywhereno typing, large, slow
JSON/JSONLrow, textnested structures, streamingeven larger, slow
Parquetcolumn, binaryanalysis, compression, schemanot human-readable
Arrowcolumn, in memoryzero-copy between toolsnot a storage format
Delta / IcebergParquet plus a transaction logversioning, ACID, time travelmore 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:

MechanismWhat it does
Column pruningreads only the requested columns
Row-group statisticsmin/max per block → skip blocks that cannot match
Dictionary encodingrepeated 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:

CodecSizeSpeed
snappylargerthe fastest — the default choice
zstdsmallernearly as fast; often the best today
gzipthe smallestslow

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

All the sources and licences