Skip to content
nhobin219Public

About

An embedded Iceberg storage engine for append-only data

Topics

Resources

Code of conduct

Contributing

Security policy

Stars

5 stars

Watchers

1 watching

Forks

litelink

CI license Python Iceberg

An embedded Iceberg storage engine for append-only data

DuckDB litelink
runs in your process in your process
without a server, daemon or cluster a server, daemon or catalog service
owns the query the durable write path
speaks SQL over Parquet and Arrow Iceberg v2, on disk and in object storage

What DuckDB is to query execution, litelink is to the durable write path. It takes high-throughput transactional appends and turns them into one well-sized Iceberg table per log — on S3, or in a local directory. append() commits to a SQLite buffer and returns once the row is durable. Behind it, the library seals rows into sorted Parquet, compacts small files up to a target size, publishes settled files to that table, and evicts from local disk what it already holds. Through all of it the log stays one queryable unit: a read sees every row exactly once, whichever tier holds it, and maintenance runs beside appends rather than in front of them.

Rows move through three tables: buffer, staging and published. advance() runs steps 1–3 to move them along, then steps 4–10 to clean up what they left behind:

  log.append(row) · log.extend(rows)        durable on return: SQLite, synchronous=FULL
                │
                ▼
  ┌───────────────────────────┐
  │ buffer         buffer.db  │   rows wait until a file's worth has arrived
  └───────────────────────────┘
                │  1. seal       writes what the size trigger cut       flush: everything
                ▼
  ┌───────────────────────────┐ ◄──┐
  │ staging    local Iceberg  │    │  2. compact   merges runs of small files
  └───────────────────────────┘ ───┘
                │  3. publish    what compaction is finished with       flush: everything
                ▼
  ┌───────────────────────────┐
  │ published        Iceberg  │   local by default, or s3://
  └───────────────────────────┘

  then cleanup:
     4. evict("buffer")        rows staging holds (published, with wal_replication)
     5. evict("staging")       files published holds, never before (I4)
     6. reclaim("buffer")      VACUUM buffer.db         only with vacuum_free_ratio
     7. reclaim("staging")     expire old snapshots, then delete files past their grace
     8. sweep("staging")       files a lost or crashed commit left behind
     9. reclaim("published")   the same, on the published table
    10. sweep("published")

One set of verbs for every table, each taking the table as an argument:

evict: drop what the next copy holds reclaim: turn that into free disk sweep
"buffer" rows staging holds, or published with wal_replication VACUUM —
"staging" files published holds (I4) expire snapshots, then delete stranded metadata
"published" — (it is the durable copy) expire snapshots, then delete stranded metadata

reclaim on staging or published does not delete a file in the call that frees it: the file waits out its table's snapshot retention first, so a scan already reading it can finish. A later reclaim deletes it.

  • flush=True means the same everywhere: push everything through this stage now, regardless of thresholds. advance(flush=True) passes it to seal and publish, so one pass leaves nothing buffered and nothing unpublished. Use it at shutdown, not in a loop: it leaves undersized files behind.
  • Nothing runs unless you call it. The library owns no thread. Call seal() often (it is one indexed read when there is nothing to do) and advance() rarely.
  • Each step is a routine of its own for an orchestrator that wants them on different schedules or in different processes; API.md's "Process split" is the recommended layout. advance() is the one-process version, and raises if another process holds a claim it needs.
  • A publish that fails stops only the published steps. Steps 4–8 still run, so a machine cut off from S3 keeps reclaiming local storage; then the error is raised.

ingest() loads data that is already durable, a Parquet corpus say: it skips the buffer and writes straight into staging at target_compact_size, then publishes, flushing its short last file by default only when the log replicates its WAL. Eviction, reclaiming and the sweeps are left to the next advance(), which a log that only ever ingests still needs. SPEC §1 diagrams both paths.

The Iceberg table is the product. The usual shape is a write path in one system and an analytical store in another, with a job copying between them. Here they are one log, and the Parquet a row is sealed into is the Parquet DuckDB, or any other Iceberg engine, reads. litelink reads across every tier, so no read on the hot path touches the network; everything else reads the published table through its version-hint.text, with no catalog and no litelink:

import duckdb
import litelink

# Written, durable on return.
log = litelink.open("data", "trades")
log.append({"trade_id": 624438572, "event_ts": 1787772776240000,
            "price": 78501.62, "amount": 0.0076})

# Read on the same box, across whichever tiers could hold a match. Inside sql(), the log
# is always the relation `log`, whatever it is named.
log.sql("SELECT count(*), max(price) FROM log").read_all()

# Read the published table with any Iceberg engine, and litelink not installed at all —
# from S3, or for a local-only log from data/trades/published/trades.
duckdb.sql("""
    SELECT count(*), max(price)
    FROM iceberg_scan('s3://bucket/prefix/trades',
                      version_name_format = '%s%s.metadata.json')
""")

The published table holds what publish has pushed, which trails the buffer by the publish interval; rows newer than that are readable through litelink on the writer's machine.

Install

pip install litelink        # or: uv add litelink

Nothing else is required — no producer, no credentials, no maintainer process, no container. Object storage and WAL replication are opt-in, and each is one setting; another machine reads the published table with no litelink at all.

Wheels for Linux and macOS on x86-64 and arm64 carry a checksum-verified litestream and the DuckDB extensions litelink loads, so a box with no egress still reads, writes and restores. That costs ~152 MB. Run python -m litelink to check a machine before you rely on it; see docs/RUNTIME.md for anywhere else.

API

litelink.new(root, name, *, schema, sort_by=None, config=None, published=None,
             s3_options=None, start_offset=1)                      -> WriteHandle
litelink.open(root, name, *, s3_options=None)                      -> WriteHandle
litelink.open(root, name, *, read_only=True, ...)                  -> LocalReadHandle
litelink.restore(root, name, *, published, s3_options=None, ...)   -> WriteHandle
litelink.validate_row(schema, row)                                 # raises as append would
litelink.preflight(...)                                            # what python -m litelink runs

# Every handle reads:
    log.scan(*, columns=None, where=None, start_offset=None, end_offset=None, published=True)
    log.sql(query, *, published=True)                 # the log is `log`; both stream Arrow
    log.column_statistics(*, tier=None) · log.coverage(*, published=True)   # tier: staging|published|buffer|None
    log.end_offset() · buffered_rows() · staging_rows() · staging_files() · published_through()
    log.schema · sort_by · config · published

# A WriteHandle also writes:
    log.append(row) -> int                          # durable on return
    log.extend(rows) -> list[int]                   # ONE transaction, one fsync
    log.ingest(table_or_reader)                     # Arrow straight to Parquet, then published
    log.advance(*, flush=False)                     # the whole pipeline above, in order
    log.seal(*, flush=False) · compact() · publish(*, flush=False)
    log.evict(table=None) · reclaim(table=None) · sweep(table=None) # its steps: clean up
    log.retire()                                    # end the log: all published, none local
    log.set_config(...)                             # the policy; schema, sort_by, published are fixed

The deliberate choices:

  • Handles, not logs. A read handle has no write methods at all, rather than ones that raise, and open(..., read_only=True) is typed so misuse is caught before it runs.
  • new takes the shape; open takes none of it. Schema, sort order, config and published table live in the log, so nothing at the call site can disagree with what is on disk.
  • The library owns no thread. Nothing seals unless you call seal() or advance(); your loop is the schedule.
  • Which tiers a query reads is decided per query, from its predicates. A query bounded inside the staging window never touches the network, however much has been evicted.

Full reference in docs/API.md.

Writing

import litelink
import pyarrow as pa

schema = pa.schema([
    pa.field("trade_id", pa.int64()),
    pa.field("event_ts", pa.int64()),    # microseconds, as the exchange sends them
    pa.field("price", pa.float64()),
    pa.field("amount", pa.float64()),
])

log = litelink.new("data", "trades", schema=schema, sort_by=("event_ts",))

log.append({"trade_id": 624438572, "event_ts": 1787772776240000,
            "price": 78501.62, "amount": 0.0076})     # durable on return
log.extend(group_of_rows)                             # the throughput lever
log.advance()                                        # seal, compact, publish, evict, expire, sweep

extend() commits the whole group in one transaction, so it is one fsync for the batch rather than one per row, and that call size is the write-throughput lever. Loading history is ingest(), which writes Arrow straight to Parquet.

Reading

sql exposes the log as log; scan(where=…, columns=…) is the typed equivalent, and both return a pa.RecordBatchReader rather than a table, so materialising is yours to choose. A reader can open the same log alongside a live writer with litelink.open("data", "trades", read_only=True).

litelink decides which tiers a query reads. Every query reads the buffer; the staging table and the published table are read only when they could hold a matching row — for the published table, a file below the staging table. That is decided from per-column bounds for each tier — the staging table's from its own Iceberg manifests, the published table's kept in buffer.db — so the decision itself never touches the network:

log.scan(where="event_ts > 1787772000000000")   # recent: local disk only
log.scan(where="event_ts < 1700000000000000")   # history: reads the published table too
log.scan()                                      # the whole log
log.scan(published=False)                       # local disk only, whatever it asks

So a query's latency follows its predicates. Bound it on a leading column of sort_by and a recent window stays local; leave it unbounded and it reads every tier, because the whole log is the right answer. Anything the decision cannot read — an OR, a subquery, a comparison with something other than a constant — reads the published table rather than risk skipping a row. The buffer keeps no column statistics, so only an offset bound (scan(start_offset=…, end_offset=…)) can skip it. column_statistics(tier=…) gives every column's bounds and counts without opening a data file, per tier ("staging", "published" below it, "buffer") or for the whole log.

Reads are cached in memory, and a reader on another machine can also cache on disk. Every duckdb_connection keeps DuckDB's memory cache on, local reads included (memory_cache=False turns it off). One built with s3_options reads S3, and with disk_cache=True also caches what it reads on disk across restarts, under ~/.cache/litelink/<cache_key> (or $XDG_CACHE_HOME), shared by every process using the key. The disk cache is off by default and never used by a log's own handles: on the host that writes a log, it would put back on disk exactly what eviction removed.

options = litelink.S3Options()
con = litelink.duckdb_connection(s3_options=options, disk_cache=True, cache_key="trades-reader")
metadata = litelink.current_metadata("s3://bucket/prefix/trades", s3_options=options)
con.sql(f"SELECT count(*) FROM iceberg_scan('{metadata}')")

With the disk cache, resolve the table with current_metadata, per read. It reads version-hint.text outside DuckDB and returns the current metadata.json. Scanning the table by its directory instead resolves the hint through the cache, whose file handle outlives a new publish: the scan fails DuckDB's ETag check instead of reading the new snapshot.

Reading from another machine

litelink reads on the primary: every handle is on the host that holds the log's root. Off that host, the published table is the interface. It is an ordinary Iceberg table that publishes version-hint.text at every commit, so an engine pointed at the prefix resolves the current metadata itself. No catalog service, no local root, no litelink install:

import duckdb

con = duckdb.connect()
con.execute("CREATE SECRET (TYPE s3, PROVIDER credential_chain, REGION 'us-east-1');")

table = con.execute("""
    SELECT count(*), max(litelink_offset)
    FROM iceberg_scan('s3://bucket/prefix/trades',
                      version_name_format = '%s%s.metadata.json')
""").arrow().read_all()

Point it at the table DIRECTORY — <published>/<name> — not at a metadata JSON. version_name_format is not optional: DuckDB defaults to the Hadoop v%s%s.metadata.json while pyiceberg names its metadata 00003-<uuid>.metadata.json, so the format has to stop prepending the v. credential_chain is the ordinary AWS resolution — profile, instance metadata, SSO; against another endpoint pass KEY_ID, SECRET, ENDPOINT and URL_STYLE 'path' instead. No INSTALL/LOAD is needed — DuckDB autoloads iceberg, avro and httpfs when a query names them, and just bootstrap provisions them ahead of time so the first read is not a download.

It reads what the published table holds, which on a quiet stream can lag the writer indefinitely rather than by the publish interval, because publish holds back a trailing run under target_compact_size. Rows still in the primary's buffer or staging table are only readable on the primary.

litelink_offset is monotonic and never reused, so a reader keeps the highest it has seen and asks for what came after — which is how you poll the published table as it grows.

Demos and recovery

just demo-websocket    # a live public feed, one process, ~30 seconds
just demo-capture      # a synthetic feed, driven as hard as you like
just demo-maintain     # in another terminal: seal, compact, publish, evict
just rustfs            # object storage in a container, to publish to S3
just demo-replicate    # ship the SQLite WAL, to survive losing the machine

Clone the repo for these; just bootstrap sets up the toolchain. Credentials are never written to the log directory — the library reads them from the environment through the ordinary AWS chain, so a profile, instance metadata or SSO all work untouched. litelink.restore(root, name, published="s3://...") rebuilds a log on another box, from its WAL replica, or from the published table alone when there is none, reserving an offset window so nothing the dead machine served is reissued.

litelink emits the litestream config; your supervisor runs the binary. Full walkthrough in examples/ and docs/RUNTIME.md.

On a KVM guest, run on the kvm-clock clocksource before turning on wal_replication. litestream panics on a single backwards tick of the monotonic clock, and a KVM guest on tsc produces them. The sidecar then crash-loops while logging healthy syncs, and the WAL stops being replicated without anything saying so. python -m litelink warns about it, and docs/API.md has a systemd unit that makes kvm-clock stick across reboots.

On disk

One directory per stream, holding everything that stream owns — and the published prefix mirrors it, so a stream can be copied, replicated or deleted whole in either tier:

data/trades/                     s3://bucket/prefix/trades/
    buffer.db                        _wal/
    catalog.db                           buffer.db/
    published.db                         catalog.db/
    litestream.yml                       published.db/
    data/                            data/
        *.parquet                        *.parquet
        compacted/*.parquet              compacted/*.parquet
        ingested/*.parquet               ingested/*.parquet
    metadata/                        metadata/
        *.metadata.json                  *.metadata.json
        *.avro                           *.avro
                                         version-hint.text

Data files sit under the table's own location, so the path an engine reads (s3://bucket/prefix/trades) is the directory that holds both halves of the table. A log with no S3 published table keeps the same table on local disk, at data/trades/published/trades/, with no _wal/. A log written before 0.6 keeps archive.db where a new one has published.db.

Upgrading a log written by 0.1.0 takes litelink 0.5.1 first: see Migrating from 0.1.

What it is not

  • Not a mutable store. Rows are only appended, never updated or deleted in place, and a log's schema, sort order and published table are fixed when it is created. To change any of them, retire() the log and start a new one where it ended: new(root, "trades-v2", schema=…, sort_by=…, published=…, start_offset=old.end_offset()). Offsets stay dense across the two, and any engine reads both as one sequence.

  • Not an unbounded staging table. A seal's cost tracks what the table's metadata holds, so a log that never runs advance() and never evicts gets slower on the write path over time. advance() arrests the larger factor; a retention, with publish() running, bounds the rest. Numbers and the reasoning are in docs/SPEC.md §13.7.

Not implemented yet

  • Indexes for point lookups. A lookup by key scans the tiers, pruned only by min/max statistics, so finding one row costs a scan rather than a seek — least on sort_by's leading column, where the statistics are tight.

  • Blob fields — large payloads that bypass the buffer — are specified and unbuilt; binary columns are carried, for ids and other small values rather than payloads (SPEC §15).

Documentation

  • docs/API.md — every public call, on one page
  • docs/SPEC.md — the design, and in places still ahead of the code
  • docs/RUNTIME.md — writer and maintainer, threads, processes, what crosses between them
  • examples/ — the websocket capture, and the synthetic feed with one process per role
  • benchmarks/ — the harness, including what litelink costs over raw SQLite
  • CONTRIBUTING.md — setup, the gates, and what a good PR here looks like
  • SECURITY.md — what to report privately, and what is a known limit instead

Development

just bootstrap          # uv sync + git hooks + DuckDB extensions + litestream
just check              # lint + format-check + typecheck + tests, same as CI
just --list             # the rest

A checkout downloads the DuckDB extensions and litestream that an installed wheel carries, so a contributor provisions what a user does not. Tooling is uv + ruff + ty + pytest; commits follow Conventional Commits, enforced by a hook. See CONTRIBUTING.md.

License

Apache License 2.0 — see LICENSE and NOTICE.

About

An embedded Iceberg storage engine for append-only data

Topics

Resources

Code of conduct

Contributing

Security policy

Stars

5 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages