Lakehouse – Phase 4 Scanning and Query

Phase 3 showed how data enters the lakehouse. This phase shows how data comes back out. A SELECT query goes through a pipeline with a critical optimization in the middle: statistics-based file pruning. The scan engine never blindly reads all data files. It first queries the metadata catalog to find which files belong to the table at the current snapshot, compares the WHERE clause against each file’s column statistics, and skips files that cannot possibly contain matching rows. This is where the database-backed catalog shines brightest — file pruning is a SQL query.

Here is the roadmap for the phases to come:

  • 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.

The Full SELECT Flow (lake/scan.go)

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
SELECT name, amount FROM orders WHERE price > 100
    │
    ▼
┌──────────────────────────────────────────────────┐
│  1. RESOLVE TABLE + SNAPSHOT                     │
│  Look up TableDef, determine snapshot_id         │
└──────────────────────┬───────────────────────────┘
                       │
                       ▼
┌──────────────────────────────────────────────────┐
│  2. GET DATA FILES (snapshot-filtered)           │
│  SELECT path, file_id FROM ducklake_data_file    │
│  WHERE table_id = $1                             │
│    AND begin_snapshot <= $snap                   │
│    AND (end_snapshot IS NULL OR end_snapshot > $) │
└──────────────────────┬───────────────────────────┘
                       │
                       ▼
┌──────────────────────────────────────────────────┐
│  3. FOR EACH FILE: PRUNE WITH STATS              │
│  Load column stats, check predicates             │
│  Skip files where price max <= 100               │
└──────────────────────┬───────────────────────────┘
                       │
                       ▼
┌──────────────────────────────────────────────────┐
│  4. FOR SURVIVING FILES:                         │
│  a. Load delete positions (skip deleted rows)    │
│  b. Read Parquet with column projection          │
│  c. Apply row-level predicate (price > 100)      │
│  d. Project output columns (name, amount)        │
└──────────────────────┬───────────────────────────┘
                       │
                       ▼
│  5. RETURN combined results                      │

In DuckLake’s C++ source, the scan flow spans DuckLakeScan (src/storage/ducklake_scan.cpp), DuckLakeMultiFileList (src/storage/ducklake_multi_file_list.cpp), and DuckLakeMultiFileReader (src/storage/ducklake_multi_file_reader.cpp).

Statistics-Based File Pruning

The pruning step compares each predicate against the column statistics stored in ducklake_file_column_stats:

1
2
3
4
5
6
7
8
9
10
11
12
Predicate: price > 100

For each file, check its price column stats:

  File 1: min=10,  max=50
    max(50) > 100? NO → SKIP (all values are <= 50)

  File 2: min=80,  max=300
    max(300) > 100? YES → MUST SCAN (some values might match)

  File 3: min=150, max=500
    min(150) > 100? YES → ALL MATCH (every value is > 100)

Pruning rules by operator

Predicate Skip file when All match when
col = val val < min OR val > max min == max == val
col != val min == max == val val < min OR val > max
col > val max <= val min > val
col >= val max < val min >= val
col < val min >= val max < val
col <= val min > val max <= val

The Go implementation checks all predicates against a file’s stats. If ANY predicate can definitively exclude the file, the file is pruned:

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
42
43
44
func canPruneFile(stats []catalog.ColumnStats, predicates []Predicate, columns []catalog.ColumnDef) bool {
    if len(predicates) == 0 {
        return false
    }

    statsByFieldID := make(map[catalog.FieldID]catalog.ColumnStats, len(stats))
    for _, s := range stats {
        statsByFieldID[s.FieldID] = s
    }

    colByName := make(map[string]catalog.ColumnDef, len(columns))
    for _, col := range columns {
        colByName[col.Name] = col
    }

    for _, pred := range predicates {
        col, ok := colByName[pred.Column]
        if !ok {
            continue
        }
        cs, ok := statsByFieldID[col.FieldID]
        if !ok {
            continue
        }

        predVal := formatForComparison(col.Type, pred.Value)
        minVal := cs.MinValue
        maxVal := cs.MaxValue

        switch pred.Op {
        case "=":
            if compareStatStrings(col.Type, predVal, minVal) < 0 ||
                compareStatStrings(col.Type, predVal, maxVal) > 0 {
                return true
            }
        case ">":
            if compareStatStrings(col.Type, maxVal, predVal) <= 0 {
                return true
            }
        // ... similar for >=, <, <=, !=
        }
    }
    return false
}

Because all stat values are stored as strings, the comparison must be type-aware. compareStatStrings parses values according to the column’s logical type:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
func compareStatStrings(lt catalog.LogicalType, a, b string) int {
    switch lt {
    case catalog.TypeInteger:
        ai, _ := strconv.ParseInt(a, 10, 32)
        bi, _ := strconv.ParseInt(b, 10, 32)
        return cmpOrdered(int32(ai), int32(bi))
    case catalog.TypeDouble:
        af, _ := strconv.ParseFloat(a, 64)
        bf, _ := strconv.ParseFloat(b, 64)
        return cmpOrdered(af, bf)
    case catalog.TypeVarchar:
        return strings.Compare(a, b)
    // ...
    }
}

The Scan Implementation

The full scan method ties together file listing, pruning, delete filtering, and predicate evaluation:

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 (l *Lake) Scan(ctx context.Context, schemaName, tableName string, opts ScanOptions) ([]datafile.Row, error) {
    var snapID int64
    if opts.SnapshotID != nil {
        snapID = *opts.SnapshotID
    } else {
        snapID, _ = l.latestSnap(ctx)
    }

    table, _ := l.getTableAt(ctx, schemaName, tableName, snapID)
    files, _ := l.store.GetDataFiles(ctx, snapID, table.TableID)

    var allRows []datafile.Row

    for _, file := range files {
        stats, _ := l.store.GetColumnStats(ctx, file.FileID)

        if canPruneFile(stats, opts.Predicates, table.Columns) {
            continue
        }

        deletedPositions, _ := l.loadDeletedPositions(ctx, file.FileID, snapID)

        localPath := l.localPath(file.Path)
        rows, _ := datafile.ReadFile(localPath, opts.Columns)

        for i, row := range rows {
            if deletedPositions != nil {
                if _, deleted := deletedPositions[int64(i)]; deleted {
                    continue
                }
            }
            if matchesPredicates(row, opts.Predicates, colByName) {
                allRows = append(allRows, row)
            }
        }
    }
    return allRows, nil
}

The pipeline for each file is: load stats → prune → load deletes → read Parquet → filter deletes → filter predicates → project columns.

Delete Filtering

When rows are deleted from a data file (Phase 7), the data file is not modified. Instead, delete files list which row positions to skip. The scan engine loads these positions into a map for O(1) lookup:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
func (l *Lake) loadDeletedPositions(ctx context.Context, dataFileID, snapID int64) (map[int64]struct{}, error) {
    deleteFiles, _ := l.store.GetDeleteFiles(ctx, snapID, dataFileID)
    if len(deleteFiles) == 0 {
        return nil, nil
    }

    deleted := make(map[int64]struct{})
    for _, df := range deleteFiles {
        localPath := l.localPath(df.Path)
        rows, _ := datafile.ReadFile(localPath, nil)
        for _, row := range rows {
            if pos, ok := row["position"]; ok {
                if p, ok := toInt64Any(pos); ok {
                    deleted[p] = struct{}{}
                }
            }
        }
    }
    return deleted, nil
}

Column Projection

Column projection ensures we only read the columns the query actually uses from each Parquet file. This happens at two levels:

Parquet-level projection: When opening a Parquet file, we tell the reader which columns to read. Unrequested column chunks are never read from disk.

Output projection: A column might be needed for filtering (price in WHERE price > 100) but not in the output (SELECT name, amount). After filtering, extra columns are dropped.

1
2
3
4
5
6
7
8
9
Query: SELECT name, amount FROM orders WHERE price > 100

Columns in file: [id, name, amount, price, region, status]

Parquet projection: read [name, amount, price]
  → Skip: id, region, status (never read from disk)

After filtering on price > 100:
  → Output: [name, amount] only

The Predicate Model

Predicates are simple column comparisons with AND semantics:

1
2
3
4
5
6
7
8
9
10
11
type Predicate struct {
    Column string
    Op     string // "=", "<", ">", "<=", ">=", "!="
    Value  any
}

type ScanOptions struct {
    Columns    []string    // nil or empty = all columns
    Predicates []Predicate // AND-combined WHERE filters
    SnapshotID *int64      // nil = latest snapshot
}

Row-level evaluation checks each predicate using type-aware comparison:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
func matchesPredicates(row datafile.Row, predicates []Predicate, colByName map[string]catalog.ColumnDef) bool {
    for _, pred := range predicates {
        col, ok := colByName[pred.Column]
        if !ok {
            return false
        }
        val, exists := row[pred.Column]
        if !exists || val == nil {
            return false
        }
        if !evaluatePredicate(col.Type, val, pred.Op, pred.Value) {
            return false
        }
    }
    return true
}

Schema Evolution: Filling Missing Columns

When reading old Parquet files that lack columns added later, the scan fills missing columns with nil. This is how ADD COLUMN works without rewriting data:

1
2
3
4
5
6
7
8
func fillMissingColumns(row datafile.Row, columns []catalog.ColumnDef) datafile.Row {
    for _, col := range columns {
        if _, exists := row[col.Name]; !exists {
            row[col.Name] = nil
        }
    }
    return row
}

This simple function handles the entire schema evolution reconciliation during reads — old files produce nil for new columns, and the caller sees a consistent schema regardless of which file a row came from.