skip to content

Why do columnar formats split each column into row groups rather than one contiguous run per column?

level: middleimportance: should knowfreq 55%

answer

  1. a chunk cannot span the whole table
  2. the writer has to buffer something
  3. one slice, many workers, no coordination
  4. statistics are only useful at fine granularity
  5. too small costs metadata, too large costs memory

basics

~20 s

Row groups bound the memory a writer and reader need, give the engine an independently readable unit it can hand to a parallel worker, and give each column chunk its own metadata so whole chunks can be skipped.

solid answer

~50 s

A single unbroken run per column would mean the writer must buffer the entire table's values for one column before it can start on the next, and a reader could not process any part of the data without addressing the whole file. Cutting the table into **row groups** — horizontal slices of tens of thousands to a few million rows — fixes all three problems at once. Memory is bounded: the writer buffers one slice, encodes each column chunk, flushes, and moves on. Parallelism has a unit: a row group is self-describing, so different workers or threads can read different slices with no coordination. And per-chunk metadata (value counts, offsets, min/max) exists at a granularity fine enough to be useful for skipping. The sizing is a genuine trade-off: too small and you pay metadata overhead, more seeks and worse compression per chunk; too large and you burn memory and skip at a coarse granularity.

code

text · 10 lines
text
file
  row group 0  (500,000 rows)
    col user_id   offset 4         len 812331     min/max, count
    col event_ts  offset 816335    len 401004     min/max, count
    col payload   offset 1217339   len 9884112    min/max, count
  row group 1  (500,000 rows)
    col user_id   ...
    col event_ts  ...
    col payload   ...
  metadata: schema + per-chunk offsets, lengths, counts, statistics

go deeper

for a junior

Know the vocabulary: a columnar table is cut into horizontal slices, and inside each slice one column's values form a chunk. Chunks are addressed through metadata rather than scanned for.

for a middle

Explain all three motivations — bounded write memory, an independent unit of parallel work, and statistics at a granularity fine enough to skip on — and state the sizing trade-off in both directions.

for a senior

Connect sizing to symptoms you have seen: tiny slices from micro-batch ingestion producing poor compression and request-bound reads, or huge slices starving a large cluster of parallel work units.

for a principal

Own the target-size policy across the platform: how ingestion batch size, table width and storage-layer request economics jointly determine the slice budget, and what you standardise so teams do not each rediscover it.

## The layout question Once you decide to store a column's values contiguously, a second decision follows immediately: contiguous over *what span*? The two extremes are one run per column for the entire table, and a fresh chunk every few rows. Every practical columnar format lands in between, by cutting the table into horizontal slices — variously called row groups, stripes, blocks or granule sets — and writing one **column chunk** per column within each slice. ## Reason 1: bounded memory on the write side Encoding a column well requires seeing its values together. A dictionary needs the distinct values; run-length encoding needs adjacent equal values; bit-packing needs to know the value range. If a chunk spanned the whole table, the writer would have to buffer every value of one column — potentially the entire table — before it could emit anything. With row groups, the writer accumulates one slice's worth of rows, encodes all of that slice's column chunks, flushes them, releases the memory and starts the next slice. Peak memory becomes a function of slice size and table width, not of table size. ## Reason 2: an independently readable unit for parallelism A row group is self-contained: given its metadata, a reader can decode any column chunk in it without touching any other slice. That makes the row group the natural unit of work distribution. A distributed engine assigns slices to workers; a single-node engine assigns them to threads. There is no coordination, no shared decode state, and no ordering requirement between slices. It also makes partial reads sane — a reader can stream slice by slice with bounded memory and produce results before the file has been fully read. ## Reason 3: metadata at a useful granularity Each chunk carries its own descriptive metadata: byte offset and length, value count, and typically per-chunk summary statistics. Offsets are what make projection possible; counts are what make cheap row counting possible; summaries are what make it possible to decide a chunk cannot contain matching rows and skip it entirely without decoding. One chunk per column per table would give you exactly one summary per column — a statistic so coarse it decides almost nothing. Chunk-level granularity is what turns metadata into a filter. (The details of how those summaries drive skipping, and how sort order determines whether they are tight or useless, belong to the pruning discussion, not here.) ## Reason 4: alignment for reassembly All the column chunks of one row group describe the *same set of rows, in the same order*. The k-th value of every chunk is the k-th row of the slice. That invariant is what allows a row to be rebuilt without storing a row identifier next to each value, and it only holds because the boundaries are shared across columns. Independent per-column chunking with different boundaries would force the format to carry positional information explicitly. ## The sizing trade-off Row group size is one of the few knobs that genuinely cuts both ways. **Too small** and you get: many chunks, so metadata grows and can rival the data for narrow columns; short chunks, so encodings have less material to work with and compression ratios drop (a dictionary rebuilt every 5,000 rows is far less effective than one over 500,000); and many small I/O requests, which is especially punishing on object storage where each request has fixed latency and a per-request cost. Vectorized operators also lose efficiency when a chunk yields fewer values than a batch. **Too large** and you get: higher writer and reader memory, since a slice must be buffered or decoded as a unit; coarser skipping, because a summary covering millions of rows is more likely to overlap any predicate range; and lumpier parallelism, because fewer slices means fewer independent work units — a table with three enormous slices cannot use thirty workers. The practical sizing rule most engines follow is to target a slice large enough that a sequential read of a single column chunk is a worthwhile I/O — comfortably more than the storage layer's minimum efficient read — while small enough that a worker can hold one in memory. Wide tables push the row count down for a fixed byte budget, because a slice holds one chunk per column. ## Where this shows up in practice It explains why an ingestion pattern that writes tiny batches produces files whose row groups are far below target: the writer flushes at the end of each batch regardless of how few rows it has seen, so chunks are short, compression is poor and metadata overhead is high. It also explains why a query over such data spends its time in request overhead rather than decoding. The fix — merging small writes into larger ones — is the compaction story, which is a separate concern; the reason it matters is the chunk sizing described here.

  • What goes wrong when an ingest job writes one row group of 2,000 rows every minute?
    Every column chunk is tiny. Dictionaries and run-length encodings have almost no material to exploit, so compression ratios collapse; metadata grows relative to data; and reads degenerate into many small requests dominated by per-request latency rather than throughput. The remedy is to batch writes larger or merge the output afterwards.
  • Does a larger row group always compress better?
    Usually better, but with diminishing returns and a real cost. Compression improves while the chunk is still short enough that encodings are starved, then flattens once the dictionary or run structure is saturated. Meanwhile writer and reader memory rise linearly and per-chunk statistics get looser, so skipping degrades. Sizing is a compromise, not a maximization.
  • Why can two engines reading the same data choose different degrees of parallelism?
    Because parallelism is bounded by the number of independently readable units. If the data sits in a few very large row groups, an engine that schedules a worker per row group cannot exceed that count; an engine that can split a row group further, or that reads column chunks independently, gets more. Slice granularity is a scheduling property, not just a storage one.

saying these in an interview costs you the question

  • Thinks a column is stored as one unbroken run across the whole table
  • Says bigger row groups are strictly better because compression improves
  • Cannot explain why chunk boundaries must align across columns
  • Confuses a row group with a table partition or a file
  • Assumes metadata is free regardless of how many chunks exist

context