Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,16 @@ there was nothing to have changed from. Everything above it is ordinary.

### Added

- **litelink's read-cache settings on every reader of the published tables**
(#77): `memory_cache`, `disk_cache`, `cache_key` and
`disk_cache_volume_limit` on `Stream.snapshot`, `scan`, `sql`, `live`, and on
`connect` for a catch-up.
- **litelink's defaults:** memory on, disk off.
- **The key is the caller's:** relative to litelink's cache root, absolute,
or its `default`. Nothing is keyed automatically.
- **Separate databases:** readers with different settings read through
different DuckDB databases.

- **`Subscription.recv_many(limit=500)` and `Subscription.batches(limit=500)`**
(#86): a subscriber's version of group commit.
- **What a batch holds:** at least one row, then every row already on the
Expand Down Expand Up @@ -41,6 +51,13 @@ there was nothing to have changed from. Everything above it is ordinary.

### Changed

- **litelink 0.9** (`>=0.9.0,<0.10`), for `litelink.current_metadata`.
- **A published table's current metadata is resolved outside DuckDB**, with
`litelink.current_metadata`. Through a disk-cached connection,
`version-hint.text` pinned a reader to the first snapshot it saw
(litelink#141, fixed in litelink 0.9 by that function). Every file the hint
names is written once and caches safely.

- **The one-shot OTel demo runs without a maintainer**, as the migration demo
does. Its maintainer processes were still opening the logs when the run
ended, and cost it several seconds: `demo.main` took 3–11 s with them and
Expand Down
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,8 @@ streamcast.connect(uri, *, offset=<unset>, cursor=None, cursor_uri=None,
catch_up=False, ...) -> Subscription
await sub.recv() · async for offset, ts, row in sub
async for batch in sub.batches(limit=500) # what has arrived, never waiting for more
# Every reader of the published tables also takes litelink's cache settings:
# memory_cache=True, disk_cache=False, cache_key=None, disk_cache_volume_limit=0.8
streamcast.publish(uri, ...) -> Publication # server needs publish=True
await producer.send(row) · await producer.send_many(rows)
await producer.submit(row) -> Future # pipelined: up to max_in_flight=64 unanswered
Expand Down
33 changes: 30 additions & 3 deletions docs/API.md
Original file line number Diff line number Diff line change
Expand Up @@ -653,7 +653,8 @@ arithmetic.
```python
streamcast.connect(uri, *, offset=<unset>, cursor=None, cursor_uri=None,
s3_options=None, upload_every=30.0, catch_up=False,
catch_up_retries=3, metadata=None,
catch_up_retries=3, metadata=None, memory_cache=True,
disk_cache=False, cache_key=None, disk_cache_volume_limit=0.8,
**websockets_kwargs) -> Subscription
```

Expand Down Expand Up @@ -1273,7 +1274,9 @@ kdb tickerplant has: the feed handler parses, the plant stores typed rows.
```python
await streamcast.Stream.snapshot(metadata_uri, *, as_of_offset=None, as_of_ts=None,
broker=None, s3_options=None,
max_tail=1_000_000) -> Snapshot
max_tail=1_000_000, memory_cache=True,
disk_cache=False, cache_key=None,
disk_cache_volume_limit=0.8) -> Snapshot
await streamcast.Stream.scan(metadata_uri, *, <the same>, columns=None, where=None,
filters=(), start_offset=None, end_offset=None) -> pa.Table
await streamcast.Stream.sql(metadata_uri, query, *, <the same>, filters=(),
Expand Down Expand Up @@ -1389,6 +1392,28 @@ do not hold yet, which it keeps in memory. That stays small while publishing kee
`max_tail` rows (1,000,000 by default, counted as read) the snapshot is refused, saying
publishing is behind, rather than running the reader out of memory.

**Caching what is read is litelink's, and yours to choose.** The four keywords are
`litelink.duckdb_connection`'s, with its defaults, and the same on `snapshot`, `scan`,
`sql`, `live` and `connect` (for a catch-up):

| keyword | default | what it does |
|---|---|---|
| `memory_cache` | on | DuckDB's external file cache: what was read stays in memory for the process |
| `disk_cache` | off | `s3://` reads kept on disk by `cache_httpfs`, across restarts |
| `cache_key` | `None` | the disk cache's directory: relative to litelink's cache root (`cache_key="<stream id>"` is `~/.cache/litelink/<stream id>`), absolute as given, or `None` for its `default` |
| `disk_cache_volume_limit` | 0.8 | how full the disk cache's volume may get, everything on it counted |

Nothing chooses a key for you, because only you know what deserves a cache of its own: a
stream id, a team, a job. Readers asking for different settings read through different
DuckDB databases. A disk cache earns its keep where reads cross a network to object
storage; against a store on the same machine, it measured no faster than reading it
again.

**A cached reader still sees every publish.** Each read resolves the table's current
metadata with `litelink.current_metadata`, outside DuckDB's caches: through a disk-cached
connection, `version-hint.text` would pin a reader to an old snapshot (litelink#141).
Everything the hint names is written once, and caches safely.

**Where the metadata file is.** With an `s3://` published location, the copy beside the
tables, which reads from anywhere. Without one, every log publishes under its own
directory and the URI is the local `file://` file, which reads only on the server's
Expand All @@ -1398,7 +1423,9 @@ machine. On another machine the read says so and suggests an `s3://` location.

```python
await streamcast.Stream.live(broker, *, s3_options=None, rebase_every=10.0, where=None,
start_offset=None, max_tail=1_000_000) -> Live
start_offset=None, max_tail=1_000_000, memory_cache=True,
disk_cache=False, cache_key=None,
disk_cache_volume_limit=0.8) -> Live

async with await streamcast.Stream.live("ws://localhost:8765/trades") as live:
await live.wait_for(offset) # until that row is visible
Expand Down
5 changes: 3 additions & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -43,12 +43,13 @@ dependencies = [
# 0.5.0 for `column_statistics`, which the manifest is built from, and
# for finite-only floats, which is what makes a float column prunable;
# 0.5.1 for seals sized by what the Arrow table actually weighs (#84);
# 0.8.1 for Python 3.14.
# 0.8.1 for Python 3.14; 0.9.0 for `current_metadata`, a published
# table's current metadata read outside DuckDB's caches (litelink#141).
#
# Capped below the next minor: under 0.x a minor is litelink's breaking
# release, and a published streamcast must not resolve one it was never
# tested against. Each release moves the cap deliberately.
"litelink>=0.8.1,<0.9",
"litelink>=0.9.0,<0.10",
# Serialisation is on the hot path in BOTH directions — every publish
# encodes a row, every replayed row re-encodes one — which is what earns
# a compiled dependency here. Measured against stdlib `json` on a
Expand Down
9 changes: 8 additions & 1 deletion src/streamcast/_catchup.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@

from streamcast import _snapshot
from streamcast._errors import NotReplayable, StreamcastError
from streamcast._published import DEFAULT_CACHE, ReadCache

if TYPE_CHECKING:
from collections.abc import AsyncGenerator, Awaitable, Callable
Expand Down Expand Up @@ -135,6 +136,7 @@ class Catcher:
"""

__slots__ = (
"_cache",
"_first",
"_handshake",
"_name",
Expand All @@ -156,8 +158,10 @@ def __init__(
start: int,
retries: int,
handshake: Callable[[int], Awaitable[tuple[Any, Any]]],
cache: ReadCache = DEFAULT_CACHE,
) -> None:
self._uri = uri
self._cache = cache
self._stream_id = stream_id
self._name = name
self._s3 = s3_options
Expand All @@ -173,7 +177,10 @@ def __init__(
async def _open(self) -> _snapshot.Snapshot:
try:
return await _snapshot.snapshot(
self._uri, s3_options=self._s3, stream_id=self._stream_id
self._uri,
s3_options=self._s3,
stream_id=self._stream_id,
cache=self._cache,
)
except _snapshot.SnapshotUnavailable as exc:
raise CatchUpUnavailable(str(exc)) from exc
Expand Down
15 changes: 15 additions & 0 deletions src/streamcast/_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@
TooSlow,
)
from streamcast._protocol import decode, parse_greeting, parse_refusal
from streamcast._published import ReadCache
from streamcast._remote import UPLOAD_EVERY, RemoteCursor

if TYPE_CHECKING:
Expand Down Expand Up @@ -718,6 +719,7 @@ class connect: # noqa: N801 — `websockets.connect` is lowercase and this mirr
"""

__slots__ = (
"_cache",
"_catch_up",
"_catch_up_retries",
"_cursor",
Expand Down Expand Up @@ -746,10 +748,22 @@ def __init__(
catch_up: bool = False,
catch_up_retries: int = CATCH_UP_RETRIES,
metadata: str | None = None,
memory_cache: bool = True,
disk_cache: bool = False,
cache_key: str | PathLike[str] | None = None,
disk_cache_volume_limit: float = 0.8,
compression: str | None = None,
**kwargs: Any,
) -> None:
self._uri = uri
# How a catch-up's reads of the published tables are cached:
# litelink's settings, passed through. See `Stream.snapshot`.
self._cache = ReadCache.of(
memory_cache=memory_cache,
disk_cache=disk_cache,
cache_key=cache_key,
disk_cache_volume_limit=disk_cache_volume_limit,
)
# `Any`, so the checker resolves `**self._kwargs` against
# `websockets.connect`'s very precise signature rather than against a
# value type inferred from `compression`.
Expand Down Expand Up @@ -851,6 +865,7 @@ async def _recover(self, refused: NotReplayable) -> Subscription:
start,
self._catch_up_retries,
self._handshake,
self._cache,
)
# BEFORE handing anything back, so an unreadable table raises here
# rather than from whatever line first calls `recv`. Entering the
Expand Down
11 changes: 10 additions & 1 deletion src/streamcast/_live.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
from streamcast import _log, _snapshot
from streamcast._errors import NotReplayable, TooSlow
from streamcast._limits import MAX_TAIL
from streamcast._published import DEFAULT_CACHE, ReadCache

if TYPE_CHECKING:
from collections.abc import Callable, Sequence
Expand Down Expand Up @@ -88,8 +89,12 @@ def __init__(
where: dict[str, object] | None = None,
start: int | None = None,
max_tail: int = MAX_TAIL,
cache: ReadCache = DEFAULT_CACHE,
) -> None:
self._broker = broker
# How every read of the published tables is cached — the base, each
# rebase, and a catch-up — as `Stream.live` was asked.
self._cache = cache
# What the view will hold from the broker, unpublished, before it
# stops: a publisher that has stalled must not run this process out
# of memory. See `_limits.MAX_TAIL`.
Expand Down Expand Up @@ -143,6 +148,7 @@ async def _connect(self) -> Any:
catch_up=True,
s3_options=self._s3,
where=self._where,
**self._cache.keywords(), # ty: ignore[invalid-argument-type]
)
try:
subscription = await connecting
Expand Down Expand Up @@ -257,6 +263,7 @@ async def rebase(self) -> None:
_require(self._broker, self._uri),
s3_options=self._s3,
stream_id=self._stream_id,
cache=self._cache,
)
async with self._lock:
# Converted against the old base first, so a pending row is judged
Expand Down Expand Up @@ -490,6 +497,7 @@ async def live(
where: dict[str, object] | None = None,
start_offset: int | None = None,
max_tail: int = MAX_TAIL,
cache: ReadCache = DEFAULT_CACHE,
) -> Live:
"""See `Stream.live`."""
from streamcast import _client # noqa: PLC0415 — the client imports `_snapshot`
Expand All @@ -515,7 +523,7 @@ async def live(
raise ValueError(msg)

base = await _snapshot.snapshot(
uri, s3_options=s3_options, stream_id=greeting.stream_id
uri, s3_options=s3_options, stream_id=greeting.stream_id, cache=cache
)
view = Live(
broker,
Expand All @@ -526,6 +534,7 @@ async def live(
where=where,
start=start,
max_tail=max_tail,
cache=cache,
)
try:
await view._start() # noqa: SLF001
Expand Down
Loading
Loading