Lakehouse – Phase 3 Ingest Pipeline
Phase 1 gave us a metadata catalog that can track tables, columns, files, and statistics. Phase 2 gave us the ability to write and read Parquet files with column statistics. This phase connects them: an INSERT writes data to a Parquet file and then atomically registers that file (with its statistics) in the metadata catalog. The design has one core principle: data files are written before metadata is updated, and the metadata update is atomic.
Here is the roadmap for the phases to come:
- 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.
The Full INSERT Flow (lake/insert.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
36
37
38
39
40
41
INSERT INTO orders (id, customer, amount) VALUES (1, 'alice', 100.50), ...
│
▼
┌──────────────────────────────────────────────────┐
│ 1. BEGIN TRANSACTION │
│ Read latest snapshot within Postgres txn │
└──────────────────────┬───────────────────────────┘
│
▼
┌──────────────────────────────────────────────────┐
│ 2. RESOLVE TABLE │
│ Look up TableDef from MetadataStore │
│ → table_id, schema, columns │
└──────────────────────┬───────────────────────────┘
│
▼
┌──────────────────────────────────────────────────┐
│ 3. GENERATE FILE PATH │
│ {schema}/{table}/data/ducklake-{uuid}.parquet │
└──────────────────────┬───────────────────────────┘
│
▼
┌──────────────────────────────────────────────────┐
│ 4. WRITE PARQUET FILE + COLLECT STATS │
│ Write rows as one row group │
│ Track min/max/null per column │
└──────────────────────┬───────────────────────────┘
│
▼
┌──────────────────────────────────────────────────┐
│ 5. REGISTER DATA FILE + STATS IN METADATA │
│ INSERT into ducklake_data_file │
│ INSERT into ducklake_file_column_stats │
└──────────────────────┬───────────────────────────┘
│
▼
┌──────────────────────────────────────────────────┐
│ 6. CREATE SNAPSHOT + COMMIT │
│ INSERT into ducklake_snapshot │
│ All metadata now visible │
└──────────────────────────────────────────────────┘
In DuckLake’s C++ source, the insert pipeline spans DuckLakeInsert (src/storage/ducklake_insert.cpp) and DuckLakeTransaction (src/storage/ducklake_transaction.cpp).
Write Then Register: The Failure Model
The ordering of operations is deliberate. The Parquet file is written to storage before any metadata changes occur:
1
2
3
4
5
6
7
8
9
10
11
12
13
Crash during file write:
→ Partial file on disk, no metadata registered
→ System state: unchanged
→ Cleanup: delete orphaned partial file (optional)
Crash during metadata insert (before COMMIT):
→ Complete file on disk, Postgres auto-rollback
→ System state: unchanged
→ Cleanup: delete orphaned file (optional)
Crash after COMMIT:
→ File on disk + metadata committed
→ System state: insert succeeded
The key insight: an orphaned file (file exists, no metadata) is always safe. It wastes storage but never causes incorrect query results. The reverse — metadata pointing to a missing file — would cause scan failures.
In DuckLake’s DuckLakeTransaction::CleanupFiles, orphaned files from rolled-back transactions are cleaned up asynchronously.
The Insert Implementation
The Lake.Insert method ties together the snapshot manager, metadata store, and Parquet writer:
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
45
46
47
48
49
50
51
52
53
54
55
56
func (l *Lake) Insert(ctx context.Context, schemaName, tableName string, rows []datafile.Row) error {
if len(rows) == 0 {
return nil
}
txn, err := l.snaps.Begin(ctx)
if err != nil {
return fmt.Errorf("lake: begin transaction: %w", err)
}
defer txn.Rollback(ctx)
_, table, err := l.resolveTableInTx(ctx, txn, schemaName, tableName)
if err != nil {
return err
}
fileName := fmt.Sprintf("ducklake-%s.parquet", newUUID())
relativePath := storage.DataPath(schemaName, tableName, fileName)
localPath := l.localPath(relativePath)
if err := os.MkdirAll(filepath.Dir(localPath), 0o755); err != nil {
return fmt.Errorf("lake: create data directory: %w", err)
}
writer := datafile.NewWriter(table.Columns)
result, err := writer.WriteFile(localPath, rows)
if err != nil {
return fmt.Errorf("lake: write parquet: %w", err)
}
// Upload to remote storage if needed
if err := l.uploadFile(ctx, localPath, relativePath); err != nil {
os.Remove(localPath)
return fmt.Errorf("lake: upload data file: %w", err)
}
fileMeta := catalog.DataFileMeta{
TableID: table.TableID,
Path: relativePath,
RecordCount: result.RecordCount,
FileSizeBytes: result.FileSizeBytes,
}
var stats []catalog.ColumnStats
for _, col := range table.Columns {
if cs, ok := result.ColumnStats[col.FieldID]; ok {
stats = append(stats, cs)
}
}
txn.AddDataFile(fileMeta, stats)
txn.RecordChange(fmt.Sprintf("insert %d rows into %s.%s", len(rows), schemaName, tableName))
_, err = txn.Commit(ctx)
return err
}
A few things to note:
- The transaction begins by reading the latest snapshot within the Postgres transaction, giving us a consistent base.
- The Parquet file is written to disk before any metadata changes.
txn.AddDataFileaccumulates the file and stats. They are flushed to Postgres whentxn.Commitis called.- The commit creates a new snapshot atomically with the file registration.
Transaction Coordination (snapshot/snapshot.go)
The Transaction type accumulates pending files and delete files, then flushes everything in a single Postgres commit:
1
2
3
4
5
6
7
8
9
type Transaction struct {
manager *Manager
tx pgx.Tx
baseSnap *Snapshot
newFiles []pendingFile
newDeletes []pendingDelete
changes []string
committed bool
}
On Commit, the transaction:
- Inserts a new snapshot row with incremented IDs.
- Inserts each pending data file with the snapshot’s
begin_snapshot. - Inserts column stats for each file.
- Inserts any pending delete files.
- Records changes in the audit trail.
- Calls
tx.Commit()to make everything visible atomically.
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
45
46
47
func (t *Transaction) Commit(ctx context.Context) (*Snapshot, error) {
snap := &Snapshot{
ID: t.baseSnap.ID + 1,
Time: time.Now().UTC(),
SchemaVersion: t.baseSnap.SchemaVersion,
NextCatalogID: t.baseSnap.NextCatalogID,
NextFileID: t.baseSnap.NextFileID + int64(len(t.newFiles)) + int64(len(t.newDeletes)),
}
// Insert snapshot record
_, err := t.tx.Exec(ctx, `
INSERT INTO ducklake_snapshot
(snapshot_id, snapshot_time, schema_version, next_catalog_id, next_file_id)
VALUES ($1, $2, $3, $4, $5)
`, snap.ID, snap.Time, snap.SchemaVersion, snap.NextCatalogID, snap.NextFileID)
// Register each data file
currentFileID := fileIDStart
for _, pf := range t.newFiles {
_, err := t.tx.Exec(ctx, `
INSERT INTO ducklake_data_file
(file_id, table_id, path, record_count, file_size_bytes, begin_snapshot)
VALUES ($1, $2, $3, $4, $5, $6)
`, currentFileID, pf.tableMeta.TableID, pf.tableMeta.Path,
pf.tableMeta.RecordCount, pf.tableMeta.FileSizeBytes, snap.ID)
// Insert column stats for this file
for _, cs := range pf.stats {
_, err := t.tx.Exec(ctx, `
INSERT INTO ducklake_file_column_stats
(file_id, field_id, null_count, min_value, max_value)
VALUES ($1, $2, $3, $4, $5)
`, currentFileID, cs.FieldID, cs.NullCount, cs.MinValue, cs.MaxValue)
}
currentFileID++
}
// Record audit trail
for _, change := range t.changes {
t.tx.Exec(ctx, `
INSERT INTO ducklake_snapshot_changes (snapshot_id, changes_made)
VALUES ($1, $2)
`, snap.ID, change)
}
return snap, t.tx.Commit(ctx)
}
File Path Management
Data files are organized in a directory hierarchy that mirrors the catalog structure:
1
2
3
4
5
6
7
8
9
10
{data_path}/
└── {schema_name}/
└── {table_name}/
├── data/
│ ├── ducklake-{uuid1}.parquet
│ ├── ducklake-{uuid2}.parquet
│ └── ducklake-{uuid3}.parquet
└── deletes/
├── ducklake-{uuid4}-delete.parquet
└── ducklake-{uuid5}-delete.parquet
The storage package provides path construction helpers:
1
2
3
4
5
6
7
func DataPath(schema, table, filename string) string {
return path.Join(schema, table, "data", filename)
}
func DeletePath(schema, table, filename string) string {
return path.Join(schema, table, "deletes", filename)
}
Table Creation (lake/table.go)
Before inserting data, you need a table. CreateTable allocates IDs from the snapshot’s counter and inserts both the table and its columns:
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
func (l *Lake) CreateTable(ctx context.Context, schemaName, tableName string, columns []catalog.ColumnDef) error {
txn, err := l.snaps.Begin(ctx)
if err != nil {
return fmt.Errorf("lake: begin transaction: %w", err)
}
defer txn.Rollback(ctx)
snapID := txn.SnapshotID()
readSnap := txn.ReadSnapshot()
schema, err := l.store.GetSchemaInTx(ctx, txn.Tx(), readSnap, schemaName)
if err != nil {
return fmt.Errorf("lake: schema %q not found: %w", schemaName, err)
}
catalogIDStart := txn.NextCatalogID()
idsNeeded := int64(1 + len(columns)) // 1 for table + 1 per column
_, err = l.store.CreateTable(ctx, txn.Tx(), snapID, schema.SchemaID, tableName, columns, catalogIDStart)
if err != nil {
return err
}
txn.SetNextCatalogID(catalogIDStart + idsNeeded)
txn.RecordChange(fmt.Sprintf("create table %s.%s (%d columns)", schemaName, tableName, len(columns)))
_, err = txn.Commit(ctx)
return err
}
Each column gets a unique field ID allocated from the snapshot’s NextCatalogID counter. The table itself also consumes one ID. This ensures field IDs are globally unique and never reused.
The Lake Entry Point (lake/lake.go)
The Lake struct is the top-level handle. It connects to Postgres, initializes the metadata schema, and ensures an initial snapshot exists:
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
type Lake struct {
store *catalog.MetadataStore
snaps *snapshot.Manager
backend storage.Backend
dataPath string
}
func Open(ctx context.Context, pgDSN string, backend storage.Backend, dataPath string) (*Lake, error) {
pool, err := pgxpool.New(ctx, pgDSN)
if err != nil {
return nil, fmt.Errorf("lake: connect to postgres: %w", err)
}
if err := catalog.InitializeSchema(ctx, pool); err != nil {
pool.Close()
return nil, fmt.Errorf("lake: initialize catalog schema: %w", err)
}
store := catalog.NewMetadataStore(pool)
mgr := snapshot.NewManager(store)
// Ensure there is at least snapshot 0.
latest, _ := mgr.LatestSnapshot(ctx)
if latest == nil {
txn, _ := mgr.Begin(ctx)
txn.RecordChange("initial snapshot")
txn.Commit(ctx)
}
return &Lake{
store: store,
snaps: mgr,
backend: backend,
dataPath: dataPath,
}, nil
}
The backend parameter controls where data files are stored (local filesystem or S3). This abstraction is built out in Phase 8.