| 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=Truemeans the same everywhere: push everything through this stage now, regardless of thresholds.advance(flush=True)passes it tosealandpublish, 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) andadvance()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.
pip install litelink # or: uv add litelinkNothing 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.
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 fixedThe 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. newtakes the shape;opentakes 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()oradvance(); 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.
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, sweepextend() 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.
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 asksSo 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.
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.
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 machineClone 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.
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.
-
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, withpublish()running, bounds the rest. Numbers and the reasoning are indocs/SPEC.md§13.7.
-
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;
binarycolumns are carried, for ids and other small values rather than payloads (SPEC §15).
docs/API.md— every public call, on one pagedocs/SPEC.md— the design, and in places still ahead of the codedocs/RUNTIME.md— writer and maintainer, threads, processes, what crosses between themexamples/— the websocket capture, and the synthetic feed with one process per rolebenchmarks/— the harness, including what litelink costs over raw SQLiteCONTRIBUTING.md— setup, the gates, and what a good PR here looks likeSECURITY.md— what to report privately, and what is a known limit instead
just bootstrap # uv sync + git hooks + DuckDB extensions + litestream
just check # lint + format-check + typecheck + tests, same as CI
just --list # the restA 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.