Skip to main content

litelink

CI license Python Iceberg

Durable append-only capture into Iceberg tables

Embedded and local-first.

Introduction

litelink is a Python library for the thing every capture pipeline hand-rolls badly: getting a stream of observations onto disk durably, into well-sized Parquet, and eventually into object storage — without a daemon, a broker, or a catalog service. append() returns once the row is durable, and a query a moment later sees it.

SQLite buffer          durable on commit. unsealed rows only.
      │  seal at target_seal_size
      ▼
local Iceberg table    a rolling window. reads land here.
      │  sync: upload data files, register into the archive
      ▼
remote Iceberg table   full history, on S3.

Reads span all three tiers, and the catalog is a SQLite file rather than a service, so no read on the hot path touches the network. Every other machine reads the archive instead, with any Iceberg engine and nothing from litelink.

It exists because doing this by hand goes wrong the same way every time: one production capture system had 125,884 objects, 62.5% of them under 16 KiB, Parquet files at 2 rows each, a compaction routine nothing ever scheduled, and an in-memory buffer a SIGKILL emptied. Durable capture, file sizing and tiering are each easy alone and nobody's job together.

Status: early. All three tiers work, and a log survives losing its machine. Read what it is not and not implemented yet first.

Install

pip install litelink        # or: uv add litelink

That is the whole of it. The wheel carries what the library shells out to — a checksum-verified litestream, and the DuckDB iceberg, avro and httpfs extensions built for the DuckDB it pins — so a machine with no egress, nothing on PATH and no DuckDB extension cache still reads, writes and restores. Verified that way in CI, not assumed.

That costs about 124 MB per wheel, and it buys the failure mode you do not want: a missing binary discovered during a failover, or an extension that autoinstalls fine on your laptop and cannot on the box that matters.

python -m litelink          # PASS/FAIL per requirement, non-zero exit if not ready

Platform wheels are published for Linux and macOS on x86-64 and arm64. Anywhere else, pip builds from the sdist — which produces a working pure-Python wheel with no binaries — and you supply litestream and the DuckDB extensions yourself; python -m litelink tells you which are missing and how to get them.

Quick start

git clone https://github.com/nhobin219/litelink && cd litelink
just bootstrap         # uv sync + git hooks + DuckDB extensions
just demo-websocket    # capture a live public feed, one process, ~30 seconds

For working ON litelink you need uv and just; for using it you need neither. Either way the demo needs nothing else: no producer, no credentials, no maintainer, no container. That demo is examples/websocket.py, and this is its shape:

import litelink
import pyarrow as pa

# A trade feed: durable the moment it arrives, queryable a moment later.
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()),
    pa.field("side", pa.int64()),        # 0 buy, 1 sell
])

# new() takes the shape, fixed at creation. open() takes none of it — schema,
# sort order, config and archive all come from the log itself.
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, "side": 0})  # durable on return

# extend() commits the whole group in ONE transaction — one fsync for the batch,
# not one per row. That call size is the write throughput lever, and a call-site
# choice: no LogConfig setting tunes it.
log.extend(group_of_rows)                          # append(row) is extend([row])

log = litelink.open("data", "trades")
reader = litelink.open("data", "trades", read_only=True)         # alongside a live writer, no write surface

recent = log.scan(where="event_ts > 1787772776000000").read_all()
log.maintain()                                     # compact, evict, expire

That is the whole API surface for local capture. Everything below is optional, and every public call is in docs/API.md — one page, forty-one of them, and most deployments use six.

More demos

A synthetic feed you can drive as hard as you like, with one process per storage role:

just demo-capture      # append continuously — the hot path, and nothing else
just demo-maintain     # in another terminal: seal, compact, evict, expire
just demo-tail         # in a third: watch where the rows are

To add the archive tier, against a local S3-compatible store or a real bucket:

just rustfs            # object storage in one container
just demo-archive      # capture, with an archive configured
just demo-maintain     # also pushes to it, and evicts what it has pushed

cp .env.example .env   # or: set LITELINK_DEMO_ARCHIVE=s3://your-bucket/prefix

Credentials are not in that file unless you put them there — the library reads them from the environment through the ordinary AWS chain, so a profile, instance metadata or SSO all work untouched, and a log directory never carries a key with it.

To survive losing the machine, ship the SQLite WAL alongside:

just demo-replicate    # generates litestream.yml from the log, runs it

litelink.restore(root, name, archive=...) then rebuilds the log on another box, reserving an offset window so nothing the dead machine served is reissued. Verified against a local S3-compatible store and against AWS. See examples/.

Run the sidecar as its own process

litelink emits the config; your supervisor runs the binary. replication_config() writes a litestream.yml describing which databases carry the log's state, and systemd, Kubernetes or anything else runs litestream replicate -config … beside the writer. litelink never starts it and never supervises it, and that is a design commitment rather than an omission:

  • It keeps the network out of the write path. A sidecar reads the WAL from outside the process. If litelink owned the replicator, a stalled upload could push back on append, and "durable when it returns, with no network in the write path" is the property the whole design rests on.
  • The lifecycles are the wrong way round otherwise. A writer that owns its replicator kills it at exactly the moment you need the last frames shipped. A separate process keeps draining what is already on disk.
  • §1 allows one writer per log. Two Log handles each spawning a replicator would run two litestream instances against one database, which litestream forbids.

litelink does run litestream in one place — restore, which shells out to it once and waits. That is a batch call, not a supervised daemon, and it is why the binary ships in the wheel rather than being left to PATH: the alternative is discovering it is missing during a failover.

Check the machine before you need it:

python -m litelink                             # extensions, litestream
python -m litelink s3://bucket/prefix trades   # ...and the archive is reachable

Run that as the user and from the process manager that will own the log. A systemd user unit does not inherit a login shell's PATH, so which litestream succeeding in your terminal says nothing about the unit that will actually perform the restore.

Reading it from another machine

The demos above are the writer's read. Everywhere else reads the archive, which is an ordinary Iceberg table publishing version-hint.text at every commit — so an engine pointed at the prefix resolves the current metadata itself, with no catalog service, no archive.db, no local root and no litelink install:

SELECT count(*), max(litelink_offset)
FROM iceberg_scan('s3://bucket/prefix/litelink/trades',
                  version_name_format = '%s%s.metadata.json');

litelink_offset is monotonic and never reused, so a reader keeps the highest one it has seen and asks for what came after — which is how you follow an archive that sync is publishing into. The extensions, the credential shapes, why version_name_format is not optional, and the polling pattern in full are in docs/API.md.

That read is only as fresh as the last sync. When you have a WAL sidecar running (wal_replication), litelink.follow does better: it restores the writer's buffer alongside the archive and merges them, so the reader sees down to the replication lag instead.

with litelink.follow("trades", archive="s3://bucket/prefix", s3=opts) as reader:
    reader.coverage()   # Coverage(archive=(1, 1928), buffered=(1929, 2100), gap=None, ...)
    reader.scan(where="side = 0", columns=["event_ts", "price"])

It never writes anything the primary shares and cannot append — a LogHandle has no write surface at all, rather than one that raises. It is a snapshot, not a subscription: refreshing means assembling another one, and the root it builds is a temporary directory removed on close. coverage() is how it stays honest about what it can and cannot serve. See §3b for why a gap in that report is not necessarily loss.

How it works

  • Iceberg is used, not reimplemented. Manifests, per-file column statistics, schema with field IDs, and atomic snapshot commits all come from it.
  • The library owns exactly one column, litelink_offset — monotonic, never reused. It is the boundary mechanism between tiers. Everything else is the caller's schema.
  • Sync is a watermark, not CDC. There are no updates or deletes to replicate.
  • Parts are sealed once and never rewritten. Rewriting a growing partition costs ~144x write amplification and buys nothing, because the local WAL already made the row durable.
  • Read boundaries are derived from committed table state, never from a stored flag — so no seal window can double-count or drop.
  • The seal cut is chosen by the appender, in the transaction that crosses target_seal_size, and queued. A sealer that falls behind therefore writes several correctly-sized files rather than one oversized one.
  • Sizing is two targets, not one. A seal wants to be small, because the buffer is what a hot read scans; a file wants to be large, because per-file overhead dominates scans and uploads. Compaction bridges them, on local disk, at 8× the seal size by default.

Read performance is the cost of reading Parquet, plus ~4 ms of fixed overhead. The numbers behind all of this are docs/SPEC.md §7 and §12; how the pieces run is docs/RUNTIME.md; just bench is the same measurement on your hardware.

What it is not

An OLTP or key-value store. A point lookup is ~1,600x slower than an indexed row store, and no configuration closes that gap. It is a local, in-process, real-time analytics store: data is durable at commit and queryable immediately, so freshness is sub-second with durability — but "real-time" means fresh, not point-lookup fast.

Nor is it an unbounded local archive. Keeping everything on one machine — archive=None with no local_retention — degrades as the table grows. A seal's cost tracks what the table's metadata holds, and a residue grows with the file count: compaction never revisits a file already at the target size, and only eviction removes one. With no retention set nothing evicts, and the seal is on the write path, so the cost lands on appends. Configure a retention, or an archive to evict into — docs/SPEC.md §13.7.

Not implemented yet

Schema evolution (docs/SPEC.md §9) and blob fields (§15) are specified and unbuilt, and are what the code lacks against its own design. add_column, rename_column and drop_column exist and raise; binary columns are refused outright, because §15 has large payloads bypass the buffer rather than travel through it.

Still open: payload encoding, local-disk backpressure, bulk ingest, and extension provisioning for embedders. All four in docs/SPEC.md §13.

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
just check              # lint + format-check + typecheck + tests, same as CI
just --list             # the rest

just bootstrap provisions the iceberg, avro and httpfs DuckDB extensions, which are downloaded rather than bundled — see docs/SPEC.md §7. Tooling is uv + ruff + ty + pytest. Commits follow Conventional Commits, enforced by a commit-msg hook; CONTRIBUTING.md has the types, scopes and style gates.

License

Apache License 2.0 — see LICENSE and NOTICE.

Metadata

Release files for litelink 0.1.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for litelink 0.1.0
File Size Uploaded
litelink-0.1.0.tar.gz 270.2 kB Details

Built distributions (wheels)

Table of built distributions (wheels) for litelink 0.1.0
File
litelink-0.1.0-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl Python 3 none Linux glibc 2.17+ x86-64 Details
litelink-0.1.0-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl Python 3 none Linux glibc 2.17+ ARM64 Details
litelink-0.1.0-py3-none-macosx_11_0_x86_64.whl Python 3 none macOS 11.0+ x86-64 Details
litelink-0.1.0-py3-none-macosx_11_0_arm64.whl Python 3 none macOS 11.0+ ARM64 Details

Total release size: 149.0 MB

Release history Release notifications | RSS feed

0.5.1

5 release files

0.5.0

5 release files

0.4.1

5 release files

0.4.0

5 release files

0.3.1

5 release files

0.3.0

5 release files

0.2.3

5 release files

0.2.2

5 release files

0.2.1

5 release files

0.2.0

5 release files

This release

0.1.0 This release

5 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page