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:

  1. The transaction begins by reading the latest snapshot within the Postgres transaction, giving us a consistent base.
  2. The Parquet file is written to disk before any metadata changes.
  3. txn.AddDataFile accumulates the file and stats. They are flushed to Postgres when txn.Commit is called.
  4. 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:

  1. Inserts a new snapshot row with incremented IDs.
  2. Inserts each pending data file with the snapshot’s begin_snapshot.
  3. Inserts column stats for each file.
  4. Inserts any pending delete files.
  5. Records changes in the audit trail.
  6. 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.