Lakehouse – Phase 2 Parquet Data Files

Phase 1 built the metadata catalog — Postgres tables that track what tables exist, what columns they have, and which files belong to them. But metadata alone is useless without actual data. This phase builds the data file layer: writing and reading Apache Parquet files with per-column statistics that the metadata catalog will reference for file pruning. We use the parquet-go library rather than implementing Parquet from scratch — the educational value is in understanding how Parquet’s structure enables file pruning and schema evolution.

Here is the roadmap for the phases to come:

  • Phase 2: Parquet data files
  • Phase 3: Ingest pipeline
  • Phase 4: Scanning and query
  • Phase 5: Snapshots and time travel
  • Phase 6: Schema evolution
  • Phase 7: Deletes and updates
  • Phase 8: Storage and access

Full Source Code

The code referenced in this post can be found in https://gitlab.com/kimserey.lam/lake-learn.

Why Parquet

A lakehouse stores data as Parquet files on storage (local disk, S3, GCS). Each table might have dozens or thousands of Parquet files, each containing a subset of the table’s rows.

Parquet is a columnar format. Within each file, data is organized by column, not by row. A query that only needs 2 of 20 columns can skip reading the other 18 entirely. Combined with per-column statistics (min, max, null count), the scan engine can decide which files and which row groups within files to read.

I/O reduction. For a 20-column table where a query uses 2 columns, you read roughly 10% of the data. This compounds with file pruning: if statistics say 7 of 10 files cannot match the WHERE clause, you read 2 columns from 3 files — 1.5% of the total data.

Compression efficiency. Values in the same column tend to be similar. Columnar layout groups them together, achieving 5-10x compression ratios.

Self-describing files. Each Parquet file contains its schema and statistics in a footer. Any Parquet-compatible tool (DuckDB, pandas, Spark, Polars) can read the data directly, without consulting the metadata catalog.

In DuckLake’s C++ source, Parquet writing happens through DuckDB’s COPY TO mechanism in src/storage/ducklake_insert.cpp. Reading goes through DuckLakeMultiFileReader (src/storage/ducklake_multi_file_reader.cpp), which maps Parquet columns to catalog columns using field IDs.

Parquet File Structure

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
Parquet File Layout:
┌──────────────────────────────────────────────────────────┐
│  Magic: "PAR1" (4 bytes)                                 │
├──────────────────────────────────────────────────────────┤
│  Row Group 0:                                            │
│    ┌──────────────────┐  ┌──────────────┐  ┌──────────┐ │
│    │ Column Chunk: id │  │ Col: name    │  │ Col: amt │ │
│    │  Page 0 (data)   │  │  Page 0      │  │  Page 0  │ │
│    │  Page 1 (data)   │  │  Page 1      │  │  Page 1  │ │
│    │  [statistics]    │  │  [statistics] │  │  [stats] │ │
│    └──────────────────┘  └──────────────┘  └──────────┘ │
├──────────────────────────────────────────────────────────┤
│  Footer:                                                 │
│    File metadata (schema, row group offsets, statistics) │
│    Footer length (4 bytes)                               │
│    Magic: "PAR1" (4 bytes)                               │
└──────────────────────────────────────────────────────────┘

Row groups are horizontal partitions of the file — batches of rows stored together. Each row group contains one column chunk per column. DuckLake typically writes one row group per file, so file-level statistics and row-group-level statistics are the same thing.

Column chunks and pages divide each column’s data within a row group. Pages are the unit of compression — each page is independently compressed. Pages also carry their own min/max statistics.

The footer is the most important part. It contains the schema definition, row group metadata, and column statistics. A reader always reads the footer first, then decides which row groups and column chunks to read.

Column Statistics: Enabling File Pruning (datafile/stats.go)

Every column chunk records three statistics: min_value, max_value, and null_count. These are the foundation of file pruning.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
Query: SELECT * FROM orders WHERE price > 500

File 1 statistics:
  price: min=10,   max=200,  null_count=5
  → max(200) > 500?  NO → SKIP entire file

File 2 statistics:
  price: min=150,  max=800,  null_count=0
  → max(800) > 500?  YES → Must scan

File 3 statistics:
  price: min=600,  max=1200, null_count=2
  → min(600) > 500?  YES → ALL rows match

Result: skip file 1, scan files 2 and 3

In DuckLake’s C++ source, DuckLakeInsert::ParseColumnStats extracts statistics from the Parquet file’s metadata after writing.

In our Go project, we track statistics on the fly during writes with a StatsCollector:

1
2
3
4
5
6
7
8
9
10
11
12
type StatsCollector struct {
    columns map[catalog.FieldID]*columnTracker
}

type columnTracker struct {
    fieldID   catalog.FieldID
    colType   catalog.LogicalType
    nullCount int64
    minValue  any
    maxValue  any
    hasValue  bool
}

The Observe method processes each value as it is written:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
func (sc *StatsCollector) Observe(fieldID catalog.FieldID, value any) {
    ct, ok := sc.columns[fieldID]
    if !ok {
        return
    }

    if value == nil {
        ct.nullCount++
        return
    }

    if !ct.hasValue {
        ct.minValue = value
        ct.maxValue = value
        ct.hasValue = true
        return
    }

    if compareValues(ct.colType, value, ct.minValue) < 0 {
        ct.minValue = value
    }
    if compareValues(ct.colType, value, ct.maxValue) > 0 {
        ct.maxValue = value
    }
}

Min/max values are serialized to strings for storage in the catalog. This is a design choice shared with DuckLake’s C++ implementation — it avoids having typed columns in the stats table.

Type Mapping: Logical Types to Parquet Types

Our catalog’s logical types map to Parquet physical and logical types:

1
2
3
4
5
6
7
8
9
Catalog LogicalType    Parquet Physical    Go type
──────────────────     ────────────────    ────────
TypeInteger            INT32               int32
TypeBigInt             INT64               int64
TypeFloat              FLOAT               float32
TypeDouble             DOUBLE              float64
TypeVarchar            BYTE_ARRAY (STRING) string
TypeBoolean            BOOLEAN             bool
TypeTimestamp          INT64 (MICROS)      time.Time

The schema mapping happens in buildParquetSchema, which constructs a parquet.Schema from our ColumnDef list:

1
2
3
4
5
6
7
8
9
10
11
12
func buildParquetSchema(columns []catalog.ColumnDef) *parquet.Schema {
    group := parquet.Group{}
    for _, col := range columns {
        node := logicalTypeToNode(col.Type)
        node = parquet.FieldID(node, int(col.FieldID))
        if col.Nullable {
            node = parquet.Optional(node)
        }
        group[col.Name] = node
    }
    return parquet.NewSchema("message", group)
}

Each column gets a Parquet field ID matching the catalog’s FieldID. This is the key to schema evolution — when a column is renamed, the field ID in old Parquet files still matches the catalog entry.

The Writer Pipeline (datafile/writer.go)

The writer accepts column definitions and rows, produces a Parquet file, and returns metadata with statistics:

1
2
3
4
5
6
type WriteResult struct {
    Path          string
    RecordCount   int64
    FileSizeBytes int64
    ColumnStats   map[catalog.FieldID]catalog.ColumnStats
}

The WriteFile method orchestrates the full pipeline:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
func (w *Writer) WriteFile(path string, rows []Row) (*WriteResult, error) {
    schema := buildParquetSchema(w.columns)

    f, err := os.Create(path)
    if err != nil {
        return nil, fmt.Errorf("datafile: create %s: %w", path, err)
    }

    pw := parquet.NewWriter(f, schema)
    stats := NewStatsCollector(w.columns)

    for _, row := range rows {
        for _, col := range w.columns {
            stats.Observe(col.FieldID, row[col.Name])
        }
        coerced := coerceRow(w.columns, row)
        if err := pw.Write(coerced); err != nil {
            _ = f.Close()
            return nil, fmt.Errorf("datafile: write row: %w", err)
        }
    }

    if err := pw.Close(); err != nil {
        _ = f.Close()
        return nil, fmt.Errorf("datafile: close parquet writer: %w", err)
    }
    if err := f.Close(); err != nil {
        return nil, fmt.Errorf("datafile: close file: %w", err)
    }

    fi, _ := os.Stat(path)
    return &WriteResult{
        Path:          path,
        RecordCount:   int64(len(rows)),
        FileSizeBytes: fi.Size(),
        ColumnStats:   stats.Results(),
    }, nil
}

Row values are coerced to exact Go types expected by parquet-go. For optional columns whose value is nil, the key is omitted so that parquet-go writes a null.

The Reader Pipeline (datafile/reader.go)

Reading is the inverse: open a Parquet file, project only needed columns, return typed rows.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
func ReadFile(path string, projectedColumns []string) ([]Row, error) {
    f, err := os.Open(path)
    if err != nil {
        return nil, fmt.Errorf("datafile: open %s: %w", path, err)
    }
    defer f.Close()

    fi, _ := f.Stat()
    pf, err := parquet.OpenFile(f, fi.Size())
    if err != nil {
        return nil, fmt.Errorf("datafile: open parquet %s: %w", path, err)
    }

    // Build projection set
    var projSet map[string]bool
    if len(projectedColumns) > 0 {
        projSet = make(map[string]bool, len(projectedColumns))
        for _, name := range projectedColumns {
            projSet[name] = true
        }
    }

    // Read all row groups
    var rows []Row
    for _, rg := range pf.RowGroups() {
        reader := rg.Rows()
        rowBuf := make([]parquet.Row, 128)
        for {
            n, err := reader.ReadRows(rowBuf)
            for i := 0; i < n; i++ {
                row := parquetRowToMap(rowBuf[i], colNames, projSet)
                rows = append(rows, row)
            }
            if err == io.EOF {
                break
            }
        }
        _ = reader.Close()
    }
    return rows, nil
}

Column projection ensures we only read the columns the query actually uses. If the query is SELECT name FROM users, the reader skips all other column chunks — their bytes are never read from disk, never decompressed, never allocated in memory.

Field IDs: The Key to Schema Evolution

Each column in the Parquet schema carries an integer field ID. DuckLake assigns these at column creation time. When reading a Parquet file, the reader matches columns by field ID, not by name:

1
2
3
4
5
6
7
8
Parquet file written before rename:
  field_id=1 name="amount"  → contains price data

After RENAME COLUMN amount → total:
  Catalog says: field_id=1 is now named "total"
  Parquet file still has: field_id=1 name="amount"

Reading: match on field_id=1 → data is correct regardless of name mismatch

In DuckLake’s C++, DuckLakeMultiFileReader performs this matching via TryFindColumnByFieldId. In our Go code, the same mapping happens during reads when we reconcile the file’s schema with the catalog’s current columns.

File Naming: UUID-Based Paths

DuckLake names data files using UUIDs to avoid collisions between concurrent writers:

1
2
3
4
{schema_name}/{table_name}/data/ducklake-{uuid}.parquet

Example:
main/orders/data/ducklake-550e8400e29b41d4a716446655440000.parquet

Each INSERT creates one new file. The UUID ensures uniqueness without coordination.