小王算钱

Interview Prep · Data Infra

Big-data file formats:
from CSV to Lance

This course takes each format in order along a timeline: which problem of its predecessor it set out to solve, how the bytes are laid out, the steps of a read and a write, its strengths and weaknesses, and where it is still used. After that come a comparison table, staff-level interview questions, an exercise you can run, and a one-page cheat sheet.

Timeline

From text files to Lance: which problem of the previous step each one solves

  1. Origin

    Twitter and Cloudera

    The problem it set out to solve

    Twitter wanted a columnar format not tied to any query engine and able to store nested data like Thrift structs. Cloudera's Impala needed a columnar format too.
    Parquet, from the start of the file to the end
    • Magic PAR1
    • row group
      • Column chunk for column 1
        • Dictionary pageoptional
        • Data pagecommonly a 1 MB default limit, encoded and compressed on its own
        • Data page …
      • Column chunk for column 2
      • Column 3 …
    • row group …
    • Footerschema, position of every row group and column chunk, statistics for each column
    • Footer length + magic PAR1

    How it stores data

    • The file is cut into row groups. In each row group every column is a column chunk, and a column chunk is cut into pages. In parquet-java and pyarrow the default upper limit for a page is 1 MB uncompressed, and pages are often smaller in practice; Parquet's own documentation recommends 8 KB.
    • The page is the smallest unit of encoding and compression. Common encodings are dictionary, run-length and delta encoding, with a general-purpose compressor applied on top.
    • The footer is encoded with Thrift and records the schema, the position of every row group, and each column's minimum, maximum and null count.
    • Nested data uses Dremel's repetition level and definition level.

    Writing, step by step

    1. Gather a row group's worth of data, and encode and compress each column into pages written one after another.
    2. After all row groups are written, the footer is written last.
    3. Because the footer is at the end, the file cannot be read until it is finished.

    Reading, step by step

    1. Read the last 8 bytes of the file to get the footer's length.
    2. Read the footer to get the schema and statistics.
    3. Pick the row groups and columns the query needs, read only those column chunks, and decompress and decode them page by page.

    Strengths

    • Reads only the columns needed.
    • Compresses well.
    • Skips whole row groups using statistics.
    • Supports nested data.
    • Readable by almost every engine and language.

    Weaknesses

    • Reading one row at random is expensive: a whole page has to be read and decoded. The footer is usually cached, but the offset index used to locate a page is optional, and without it a reader has to walk the column chunk one page header at a time.
    • Very large values such as images and vectors fit badly: row groups are cut by row count, so with few rows the integer columns are too fragmented, and with many rows a single group of the image column is several GB.
    • With a very large number of columns (thousands), the footer has to be parsed in full, and just reading metadata is slow.
    • Filling in a new column for existing rows means rewriting the files.

    Where it is used, and where it still is

    Analytical queries in the data lake, and the data files of Iceberg, Delta and Hudi. Today's default choice.

    What replaced it, and why

    Not replaced. Newer formats challenge it where it is weak: random reads, very wide values, tables with thousands of columns.

    In an interview, one sentence

    The de facto columnar standard: three levels of row group, column chunk and page, with metadata in the footer.

    In an interview, two minutes

    Parquet has three levels: the file is cut into row groups, each column in a row group is a column chunk, and that is cut into pages. The page is the unit of encoding and compression. The footer stores the schema and each column's minimum and maximum, so a reader can read only the columns needed and skip whole row groups. Nested data uses Dremel's two levels. Its weak points are that fetching one row at random is expensive, very wide values are awkward, and filling in a column for existing rows means a rewrite. It became the standard mainly through its ecosystem: it is independent of any engine, Spark chose it as the default, and later Arrow, the cloud warehouses and the table formats all built around it.

    Worth adding

    It won mainly through its ecosystem, not through technical superiority. The section "Why Parquet all but took over" below covers this.

Five concepts first

Every format makes its trade-offs on these five things

Row and columnar storage
Whether the fields of one record sit next to each other, or the values of one column do.
Writes and whole-row reads favor row storage; analytical queries that read a few columns favor columnar. This is the first fork for every format.
Splittable
Whether a large file can be cut in the middle and read by several tasks at once.
A file that cannot be split can only be read through by one task. Sync markers, row groups and stripes all exist so that a file can be cut.
Encoding and compression
Encoding exploits patterns in the data itself (repetition, increasing values); compression squeezes it once more at the byte level.
Columnar storage compresses well because the values in a column share a type and often repeat. Dictionary or run-length encoding first, then compression, does far better than compressing a whole row directly.
Data skipping
The engine hands the query's conditions to the reader (this step is called predicate pushdown), and the reader uses each block's minimum and maximum to skip blocks that cannot match.
How fast a query is depends largely on how much data it can avoid reading. The finer the statistics, the more can be skipped, but the larger the metadata.
Schema evolution
Whether old data can still be read after fields are added, removed or change type.
Data is kept for many years and table structures always change. Avro aligns two schemas, Protobuf relies on field numbers, and Iceberg gives every column an ID that never changes.

In depth · 1

A SequenceFile example

Suppose there are a million small images on HDFS. Every file takes NameNode memory, and every file gets its own read task. The usual approach is to pack them into one SequenceFile: the key is the file name and the value is the image's bytes.

# PySpark: pack many small files into one SequenceFile
# the key is the file name, the value is the file's contents
pairs = sc.parallelize([
    ("cat-001.jpg", bytearray(open("cat-001.jpg", "rb").read())),
    ("cat-002.jpg", bytearray(open("cat-002.jpg", "rb").read())),
])
pairs.saveAsSequenceFile("hdfs:///datasets/cats.seq")

# read it back: (file name, contents) key-value pairs
sc.sequenceFile("hdfs:///datasets/cats.seq").take(1)

The file that is written looks roughly like this inside:

SEQ\x06                                   magic and version
org.apache.hadoop.io.Text                 key class
org.apache.hadoop.io.BytesWritable        value class
compressed: no   block-compressed: no     compression flags
<16-byte sync marker>                     end of header

[record length][key length] "cat-001.jpg" <image bytes>
[record length][key length] "cat-002.jpg" <image bytes>
<16-byte sync marker>                     appears at intervals
[record length][key length] "cat-003.jpg" <image bytes>

The header contains Java class names, which is why other languages find it hard to read. The sync marker lets a task start in the middle of the file and find the next complete record.

In depth · 2

Why Parquet all but took over

Not because it crushed ORC technically. The two appeared in the same year, have similar structures, and each wins on some measures. The difference was the ecosystem:

  1. 1

    Not tied to any one engine

    ORC grew up inside Hive. Parquet was built by Twitter and Cloudera together, with the goal from the start that several engines could use it, and it came with conversions to data models such as Avro, Thrift and Protocol Buffers.

  2. 2

    Spark made it the default format

    Spark's default data source is Parquet. In the years when Spark became the mainstream compute engine, a great deal of new data was naturally written as Parquet.

  3. 3

    It handles nested data well

    It uses Dremel's approach, which lines up with nested data models such as Protocol Buffers and Thrift.

  4. 4

    Arrow and the Python ecosystem

    pandas, DuckDB and Polars can each read and write Parquet in one line of code. A data scientist can use it without setting up Hadoop.

  5. 5

    Cloud warehouses and table formats all built around it

    The query services of the major clouds can read Parquet directly. Delta Lake uses only Parquet, and Iceberg's default data file is Parquet too.

Once a format becomes the one everyone can read, new tools support it first, and so more data is written in it. That is a network effect, and it has little to do with the details of the format itself.

In depth · 3

What Lance actually changed

  1. Why fetching one row from Parquet is slow

    To fetch one value from row N, a reader first reads the footer, finds which row group and which page it is in, reads the whole page (commonly up to 1 MB by default), decompresses and decodes it, and finally takes out that one value. This is repeated for every column wanted. The offset index used to locate a page is optional, and without it the reader walks the page headers one by one. In a sequential scan these costs are spread thin; in random reads every row pays them again.

  2. How Lance makes random reads fast

    It drops row groups and lets each column cut its own pages. Version 2.1 chooses a layout by the width of a value: narrow values such as integers and short strings are cut into small blocks, each holding at most 4,096 values by default and required by the specification to be no more than 32 KiB compressed, so fetching one value reads one small block; wide values such as images and vectors are laid end to end and can be located directly. The dividing line in the current specification is 256 bytes per value.

  3. Why adding a column needs no rewrite

    A table is made of fragments, and a fragment can have several data files, each holding different columns. To add an embedding column, you only write one file per fragment containing just the new column, then commit a new version. The files that hold the images are not touched at all.

  4. It is a table format as well

    Every write produces a new manifest recording which fragments this version has. So it comes with historical versions. A delete writes a deletion file as a marker. Vector, scalar and full-text indexes are recorded in the table too.

  5. Where to be careful

    The ecosystem is far smaller than Parquet's; the format went through 2.0, 2.1 and 2.2 within a year or two; and as versions and small files pile up they need regular compaction and cleanup. LanceDB has published performance numbers comparing it with Parquet, but I could not confirm their test conditions or find independent verification, so they are not quoted here. The Lance paper itself also notes that Parquet, configured appropriately, can do random reads much faster than with its default settings.

Seven follow-ups

What an interviewer tends to ask next

1Why can a file compressed as a whole with gzip not be split?

The DEFLATE algorithm that gzip uses writes instructions like "go back 300 bytes and copy 20 bytes", reaching back as far as 32 KB. A task that starts reading in the middle of the file does not have the earlier data, so the instruction cannot be carried out. In addition, DEFLATE's internal blocks can start at any bit, with no marker to search for in between, and each block's code table is written at the start of the block. Put together, the file can only be read from the beginning to the end.

Top: the whole file compressed into one gzip stream. Bottom: a container format compressing block by block internally. The red line is the split point.
Earlier data
← go back 300 bytes, copy 20 bytes

A task that starts at the red line does not have the 20 bytes to copy; they are to the left of the line.

Block 1, compressed on its own
S
Block 2
Block 2
S
Block 3, compressed on its own

The task that starts at the red line: scan forward to the next S and start decompressing at block 3. Block 2 belongs to the previous task.

Two kinds of compression can be split. One is bzip2: each block is compressed independently and the block header has a fixed marker, so the next block can be found from the middle. The other is a container format that compresses in blocks internally, which is what SequenceFile, Avro and Parquet all do.

2When a SequenceFile is read from the middle, what happens to the records before the sync marker?

They are not dropped; they belong to the previous task. Two rules work together: each task first scans forward to the first sync marker in its own section and starts reading after it; and on reaching the end of its section it does not stop at once, but keeps reading until the first sync marker past the end. That way every record is read exactly once.

R is one record, S is a sync marker, and the red line is the boundary between the two tasks.
H
R
R
S
R
R
R
S
R
R
S
R
Task A reads these
—
Task B reads these

Task A reads past the red line and stops only at the next S. The record right after the red line is in task B's range but is read by task A; task B skips it and starts after the S.

3Why does SequenceFile store key-value pairs? It is not a hash table.

It really is not a hash table: you cannot look up by key, keys can repeat, and it is just "a sequence of key-value pairs". It is stored this way because every step of MapReduce works on key-value pairs, and one job's output has to serve directly as the next job's input.

How data flows through MapReduce in the word-count example. Key-value pairs go in and come out at every step.
  1. Input

    (0, "a b a")

  2. map output

    (a, 1) (b, 1) (a, 1)

  3. Group by key

    (a, [1, 1]) (b, [1])

  4. reduce output

    (a, 2) (b, 1)

Often the key carries no real meaning. When storing a table, the usual practice is to leave the key empty and put the whole row in the value. For real lookups by key, Hadoop has MapFile: a sorted SequenceFile plus an index.

4Why did RCFile not carry types?

Because it left types to the layer above. RCFile was written for Hive, where the work was divided into three layers at the time, and RCFile did only the bottom one. The goals the paper set for it were also about placement: fast loading, fast queries, efficient use of space and adaptability to changing query patterns.

The three layers at work when Hive reads an RCFile table. Types live in the upper two; the file holds only bytes.
  1. metastore: the table's schema

    name string, age int

  2. SerDe: interprets bytes as types

    bytes → "alice", 30

  3. RCFile: only how bytes are placed

    each column is a run of compressed bytes

The cost showed up later: a file that does not know types cannot use delta encoding for integers or dictionary encoding for strings, and cannot record a minimum and maximum. ORC put types into the file and fixed exactly this. The three-layer division described here follows Hive's architecture at the time; I did not find a direct statement in the paper of why types were not stored.

5How exactly do ORC and Parquet differ in storing nested data?

Take the same data. The schema is struct<name: string, tags: array<string>>, with three rows:

Rownametags
1a[x, y]
2b[] (empty array)
3cnull

ORC treats every level of the structure as a column, including the list in the middle. Each column has several streams: PRESENT marks whether each value is null, LENGTH records how many elements each row has (or how long a string is), and DATA is the values themselves.

ORC: three columns, each with its own streams. Where a parent is null, nothing is stored for its children.
ColumnPRESENTLENGTHDATA
name1, 1, 11, 1, 1abc
tags (the list level)1, 1, 02, 0—
elements of tags1, 11, 1xy

Parquet stores only leaf columns, and the list level has no column of its own. All the structural information goes into the two levels of the leaf column.

Parquet: the single leaf column "elements of tags". A place with no value still takes a row, and the levels say why.
Valuerepetition leveldefinition levelMeaning
x03New row, value present
y13Still in the same list
—01New row, the list exists but is empty
—00New row, the list itself is null

The trade-off: ORC is intuitive, but reading a deeply nested leaf means reading PRESENT and LENGTH for every level along the path. Parquet reconstructs the structure from one column, but every leaf carries two levels and the logic to assemble records is more complex.

6Dictionary, run-length and delta encoding: one example of each

Dictionary encoding: for columns with few distinct values. Each value is replaced by a very small code.
Original
USCNUSUSINCN
Dictionary
0 = US1 = CN2 = IN
Encoded
010021
Run-length encoding: for columns where identical values sit together. Only "which value, how many in a row" is recorded.
Original
okokokokokfailfailok
Encoded
ok × 5fail × 2ok × 1
Delta encoding: for numbers that grow slowly. Only the first value and each later step's difference are recorded.
Original
1700000000170000000317000000041700000009
Encoded
1700000000+3+1+5

The three can be stacked. Parquet commonly applies dictionary encoding first, then run-length encodes the codes or packs them tightly by bit. Run-length encoding works best when the data is sorted and is useless when values alternate; the country column in the exercise below is the latter case.

7Before Arrow, how slow was it for Spark to hand data to Python?

Suppose Spark holds ten million rows of numbers and you want Python to add 1 to each. The computation itself takes a few tens of milliseconds in NumPy, and almost all the time goes on moving the data.

Without Arrow: data crosses in batches, but every row is converted and passed to the function on its own, ten million times in all.
  1. JVM: converts this row from its own memory format to pickle bytes

  2. Sends the rows to the Python process over a socket, a batch at a time

  3. Python: turns the bytes back into a Python object

  4. Calls your function and computes x + 1

  5. The result is pickled again and sent back, and the JVM converts it to its own format

With Arrow: a batch at a time (say ten thousand rows), and both sides use the same memory layout.
  1. JVM: writes a batch of data in Arrow's columnar layout

  2. Python: receives arrays pandas can use directly, with no value-by-value conversion

  3. Computes x + 1 on the whole array and sends the result back in the same layout

In code, the only difference is replacing @udf with @pandas_udf, and the function's argument changes from one number to a pandas.Series. When Databricks released the feature in 2017, it tested three functions on ten million rows and found speedups of 3x to more than 100x over row-at-a-time UDFs. That is Databricks' own test.

Side by side

The same dimensions, in one table

FormatRow or columnarWhere the schema isSplittableCan skip dataBest for
CSV / JSONRow, textNoneMostly, when uncompressedNoExchanging data, reading by eye
SequenceFileRow, binaryJava class names onlyYesNoLegacy systems, packing small files
AvroRow, binaryIn the file headerYesNoKafka, data ingestion, metadata files
Thrift / ProtobufSingle message, not a file formatIn the IDL file; data carries field numbersNot applicableNot applicableRPC, messages, metadata of other formats
RCFileColumnar, by row groupNo type informationYes, by row groupNo; it can only leave columns unreadLegacy Hive tables
ORCColumnarIn the file footerYes, by stripeYes; one index entry per 10,000 rows by defaultWarehouses in the Hive ecosystem
ParquetColumnarIn the footerYes, by row groupYesAnalytical queries in the data lake
ArrowColumnar, in memoryTravels with the dataNot applicableNot applicablePassing data between systems, in-memory computation
Iceberg / Delta / HudiTable format; data files mostly ParquetIn the table's metadataBy data fileYes; whole files are skipped firstWarehouses on a data lake
LanceColumnar, no row groupsIn the manifest and the file footerYes, by fragmentYes, plus indexesTraining data, vector search, multimodal data

Where the industry is heading · checked October 2026

What people are working on now

  • Formats designed for AI data

    Training and retrieval need random reads, and the data contains images and vectors. Lance is the most watched in this direction, and in early 2026 the Apache Polaris catalog added support for Lance tables.

  • Separating encodings from file structure

    Nimble, Vortex and F3 all make encodings swappable and nestable. Vortex joined LF AI & Data, under the Linux Foundation, in August 2025.

  • Parquet itself keeps changing

    parquet-format 2.11.0 in March 2025 added two geospatial types, GEOMETRY and GEOGRAPHY, and included Variant for the first time; 2.12.0 in August 2025 finalized the Variant specification. Variant stores semi-structured data whose structure is not fixed.

  • Table formats are filling in row-level updates

    Iceberg's v3 specification introduced deletion vectors: a bitmap marks which rows of a data file have been deleted, and the bitmap is stored in a Puffin file. This makes frequent updates and deletes cheaper.

Common misconceptions

These sound right but are not

  • Misconception: Parquet won because its technology is better than ORC's.

    Technically each leads in some respects. A paper published at VLDB 2024 found that Parquet decodes faster and its files are slightly smaller, while ORC skips data more effectively. Parquet won mainly because it is not tied to an engine and because Spark made it the default format.

  • Misconception: Nobody uses Avro any more.

    It has left only analytical storage. Kafka messages, data ingestion and Iceberg's manifest files still use Avro.

  • Misconception: Iceberg, Delta and Hudi are file formats.

    They are table formats, a layer of metadata on top of data files. The data files are overwhelmingly Parquet: Delta uses only Parquet, and Iceberg also allows ORC and Avro.

  • Misconception: Columnar is always better than row storage.

    For record-by-record writes, whole-row reads and fetching one row by primary key, row storage fits better.

  • Misconception: Arrow is a replacement for Parquet.

    One governs memory and the other disk, and they are usually used together.

  • Misconception: A compressed file can never be split.

    Only a file compressed with gzip into a single stream cannot be split. SequenceFile, Avro and Parquet all compress block by block inside the file, so they can still be split.

  • Misconception: Lance will replace Parquet.

    They target different ways of reading. Full-table scans and aggregation remain Parquet's home ground.

Staff-level interview questions

Answer first, then open

▶Why does columnar storage compress better than row storage?
  • The values in one column share a type and often repeat or follow a pattern. They can be encoded first: dictionary encoding where values repeat a lot, run-length encoding for consecutive identical values, delta encoding for increasing numbers.
  • General-purpose compression is applied after encoding. In row storage a row mixes integers, strings and timestamps, so adjacent bytes follow no pattern and the only option is to compress them directly.
  • How much better depends on the data: columns with few distinct values and some ordering gain the most, and random strings that differ on every row gain almost nothing.
▶A 10 GB gzip-compressed CSV and a 10 GB Parquet file: what is different about reading them?
  • The gzipped CSV cannot be split and has to be decompressed start to finish by one task. Parquet can be divided among tasks by row group.
  • CSV has to read every column and parse text. Parquet reads only the columns used and can skip whole row groups using statistics.
  • How big the difference is depends on the query: it is largest when reading a few columns with a filter, and much smaller for a full export that reads every column.
▶On the path from Kafka to the data lake, which format is used at each layer, and why?
  • In Kafka, Avro (or Protobuf) with a Schema Registry: records are written one at a time, messages are small, and the schema changes often.
  • Landing in the data lake, Parquet: from then on it is analytical queries that read a few columns.
  • On top of the Parquet files, Iceberg or Delta: writes become atomic, the table's structure can change, and historical versions can be read.
▶Iceberg's data files are Parquet. Why are its manifests Avro?
  • A manifest is a series of "file records". Early on they had few fields and a reader used essentially the whole record, which suits row storage.
  • It needs dependable schema evolution, since Iceberg adds fields to the manifest when the format itself is upgraded.
  • A manifest is never modified once written, and writing a row format in one pass is simple. That was the trade-off at the time. As statistics grew wider, reading only the partition and a few statistics still meant decoding the whole record, and the drawback of row storage showed. In discussions of v4, the Iceberg community has proposed moving manifests to Parquet.
▶How large should a row group be?
  • Larger: better compression, less metadata, faster sequential scans.
  • Smaller: finer statistics, more data skipped, and less memory used when reading.
  • Parallelism matters too: a row group is usually the smallest unit of a read task. There is no standard answer; the point is to state the trade-off clearly.
▶When is Parquet the wrong choice?
  • Fetching a few rows by primary key: each fetch reads and decodes a whole page.
  • Columns with very wide values such as images and vectors: no row group size works for both.
  • Wide tables with thousands of columns: the footer has to be parsed in full.
  • Frequently adding columns and backfilling them: the data files are rewritten every time.
▶Adding an embedding column to a table of one billion rows: what does Parquet with Iceberg have to do, and what does Lance?
  • Parquet with Iceberg: changing the table's structure is itself quick, since only metadata changes. Iceberg v3 also lets a new column carry a fixed default value without a rewrite. But when every row has a different value, putting the values in means rewriting every data file, and the existing columns are rewritten along with it.
  • Lance: write one data file per fragment containing only the new column, then commit a new version. The original files are untouched.
  • The cost is that a Lance fragment then has several files, reading a whole row opens more files, and compaction may be needed later.
  • There is a third way: store the embeddings in a separate table and join it to the original by primary key. Which to choose depends on how often the column is recomputed and whether queries always read it together with the original table.
▶What is the small-files problem, and how is it solved?
  • Streaming writes, or partitions cut too finely, produce huge numbers of very small files. On HDFS every file takes NameNode memory. On object storage, the cost of listing files and opening them one by one exceeds the cost of reading the data.
  • The fix is regular compaction. A table format lets compaction run safely in the background, because readers always see one complete version.
  • One of SequenceFile's original uses was packing small files into large ones.
  • How often to compact is a trade-off: compacting often makes reads fast, but compaction itself costs compute and can conflict with jobs that are writing. It depends on whether reads or writes matter more.
▶Design question: choosing a storage format for a machine-learning training platform. How do you think about it?
  • Ask first how the data is read: sequential scans of whole tables, or random sampling? Are there images, video or vectors? How often are feature columns added?
  • Ask about the surroundings too: can the existing engines and tools read it? Can the team carry the maintenance risk of a newer format?
  • If the workload is mostly sequential scans of tabular features, Parquet with a table format is the safe choice. Only when random access, multimodal data and frequent column additions dominate is a format like Lance worth considering.
  • It can be layered: raw data and reporting stay in Parquet, and training datasets are stored separately.

Hands-on exercise · about 5 minutes

Look at a Parquet footer yourself

Requires Python and pyarrow (pip install pyarrow). Save the code below as demo.py and run it.

import os
import pyarrow as pa
import pyarrow.csv as pacsv
import pyarrow.parquet as pq

# 1. Build a table of one million rows: an increasing integer column,
#    a column with only 4 distinct values, and a decimal column
n = 1_000_000
table = pa.table({
    "user_id": pa.array(range(n), pa.int64()),
    "country": pa.array(["US", "CN", "IN", "BR"] * (n // 4)),
    "amount": pa.array([(i * 7919) % 100_000 / 100 for i in range(n)]),
})

# 2. Write it as CSV and as Parquet, and compare the sizes
pacsv.write_csv(table, "demo.csv")
pq.write_table(table, "demo.parquet", row_group_size=250_000)
print("CSV    ", os.path.getsize("demo.csv") // 1024, "KB")
print("Parquet", os.path.getsize("demo.parquet") // 1024, "KB")

# 3. Read the footer: the file's structure, without reading any data
meta = pq.ParquetFile("demo.parquet").metadata
print("row groups:", meta.num_row_groups, " total rows:", meta.num_rows)

# 4. For each column in the first row group: bytes used and encodings
for i in range(meta.num_columns):
    col = meta.row_group(0).column(i)
    print(col.path_in_schema, col.total_compressed_size // 1024, "KB", col.encodings)

# 5. The minimum and maximum of user_id in each row group
for g in range(meta.num_row_groups):
    stats = meta.row_group(g).column(0).statistics
    print("row group", g, "user_id:", stats.min, "to", stats.max)

# 6. Read one column, and only rows with user_id >= 900000
result = pq.read_table(
    "demo.parquet", columns=["amount"], filters=[("user_id", ">=", 900_000)]
)
print("read", result.num_rows, "rows")

Actual output on pyarrow 19.0:

CSV     18221 KB
Parquet 8991 KB
row groups: 4  total rows: 1000000
user_id 1249 KB ('PLAIN', 'RLE', 'RLE_DICTIONARY')
country 2 KB ('PLAIN', 'RLE', 'RLE_DICTIONARY')
amount 994 KB ('PLAIN', 'RLE', 'RLE_DICTIONARY')
row group 0 user_id: 0 to 249999
row group 1 user_id: 250000 to 499999
row group 2 user_id: 500000 to 749999
row group 3 user_id: 750000 to 999999
read 100000 rows
  • The country column has 250,000 values in the first row group and takes only 3,049 bytes (the output rounds down to KB and shows 2 KB). It has only 4 distinct values, so dictionary encoding turns each value into a 2-bit code, about 62 KB for 250,000 values; the codes cycle in a fixed order, and Snappy then squeezes that to about 3 KB. Run-length encoding did nothing here, because no two adjacent values are the same.
  • To see run-length encoding at work, sort the country column and write the file again: identical values sit together and the column shrinks to a few dozen bytes.
  • The minimum and maximum printed in step 5 are stored in the footer. The filter in step 6 is user_id >= 900000, and the maximum of each of the first three row groups is below 900000, so those groups are skipped entirely and only the fourth is read.
  • Try changing things: set row_group_size to 50000 and run again to see how the number of row groups and the file size change; replace country with a string that differs on every row and see how large that column becomes.

Cheat sheet

Ten minutes before the interview

CSV / JSON
Text, universal, untyped, not splittable once compressed as a whole.
SequenceFile
Binary key-value pairs; sync markers make it splittable; tied to Java.
Thrift / Protobuf
Not file formats. Field numbers give schema evolution.
Avro
Row format with the schema in the header. Mainstream on the write side: Kafka, Iceberg manifests.
Dremel
Repetition level and definition level: nested data stored by column.
RCFile
Record Columnar: row groups first, columns inside. No types, no index.
ORC
Stripes, indexes and a file footer; knows types. The Hive ecosystem.
Parquet
Row group, column chunk, page; statistics in the footer. The de facto standard.
Arrow
Columnar layout in memory, zero copy. Paired with Parquet.
Iceberg / Delta / Hudi
Table formats: atomic commits, historical versions, schema changes.
Lance
No row groups, fast random reads, columns added without a rewrite, built-in indexes.
Nimble / Vortex / F3
Wide tables, extensible encodings, GPUs. All still new.

Sources

The parts on newer formats (Lance, Nimble, Vortex and F3, and the new features of Parquet and Iceberg) were checked in October 2026. They may have changed since.

DISCUSSION

Join the discussion

Share a thought or question. Comments publish immediately; add your email only if you would like reply notifications.

Be the first to start the conversation.