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.