Streams
StreamWriter and StreamReader combine a stream's layout with per-log operations, so writing or reading a stream is a single call.
StreamWriter and StreamReader (io.lakestream.api, both extend AutoCloseable) combine a stream's StreamLayout with per-log operations into a single write or read call. StreamWriter.write spares the caller from resolving a routing key to a LogId itself; StreamReader.read still takes a LogId explicitly, since reading targets one specific log, but spares the caller from driving LogStorage or a Log handle directly. Both interfaces declare that implementations must be safe for concurrent use, and both are opened through StreamCatalog.openWriter/openReader.
Types documented on this page (5): StreamWriter (with nested WriteResult), StreamReader (with nested ReadResult), StreamLayout, RoutingKey, StreamPosition.
StreamWriter and StreamReader
interface StreamWriter extends AutoCloseable {
CompletableFuture<WriteResult> write(RoutingKey key, int numberOfRecords, ByteBuf data);
StreamLayout layout();
record WriteResult(LogId logId, long offset) {}
}
interface StreamReader extends AutoCloseable {
CompletableFuture<ReadResult> read(LogId logId, long startOffset, int maxMessageCount, long maxSizeBytes);
CompletableFuture<List<LogId>> logIds();
StreamLayout layout();
record ReadResult(List<LogEntry> entries, long nextOffset) {}
}write resolves the target log from the routing key through the stream's layout, then appends to it, returning the LogId the write landed on and the offset assigned to the entry. Buffer ownership follows the same rule as Log.append and LogStorage.append: the caller retains ownership of data until the returned future completes and must then release it exactly once — see Logs: Buffer ownership.
read reads from one log the caller names explicitly — StreamReader does not merge or order entries across a stream's logs into one sequence. It returns the entries read plus nextOffset, the offset to pass as startOffset on the following call; the caller must close every entry in entries, the same as any other read.
StreamLayout
Partitioning partitioning();
CompletableFuture<List<LogId>> logIds();
CompletableFuture<LogId> resolveForWrite(RoutingKey key);
int logCount();
StreamPosition position(LogId logId, long offset);logIds() returns every log in the stream, in order; for PartitioningStrategy.INDEXED (see Catalog), that order is the partition index order. resolveForWrite picks a log for a given RoutingKey, and fails with IllegalArgumentException if the key's index hint is outside the stream's log count. position(logId, offset) creates a StreamPosition — the layout-specific way of naming a place in the stream — for callers that need to hold onto a position without depending on a specific log/offset pair's meaning.
RoutingKey and StreamPosition
record RoutingKey(OptionalInt indexHint) {
static RoutingKey ofIndex(int index);
static RoutingKey roundRobin();
}
interface StreamPosition {}RoutingKey supports two ways to route a write: ofIndex(i) targets one log by ordinal index, and roundRobin() (an empty indexHint) leaves the choice to the layout's own distribution. There is no key-hash routing on this type.
StreamPosition is a marker interface with no methods of its own — it exists so a layout can hand back an opaque value from position(logId, offset) that only that same layout knows how to interpret, rather than every caller working directly in (LogId, offset) pairs.
Not in the API
There are no stream-level softTrim/hardTrim methods. Trimming is per log, through Log or LogStorage (see Logs); a stream-wide trim means calling it once per log in the stream's layout.
A sketch of the shape
This sketch uses only real members of the API, with ... marking wiring a real program would supply. It is not runnable as written — for a runnable version against Maven Central and a live catalog, see the Ursa quickstart.
StreamCatalog catalog = ...;
// in format version 3, neither a namespace nor a stream name can contain "/"
StreamIdentifier id = StreamIdentifier.of("default", "events");
Partitioning partitioning =
new Partitioning(PartitioningStrategy.INDEXED, Map.of("numPartitions", "1"));
ByteBuf payload = ...; // event bytes
catalog.createStream(id, new StreamConfig(), partitioning, new SchemaConfig(), Map.of())
.thenCompose(metadata -> catalog.openWriter(id))
.thenCompose(writer -> writer.write(RoutingKey.roundRobin(), 1, payload))
.thenCompose(written -> catalog.openReader(id)
.thenCompose(reader -> reader.read(written.logId(), written.offset(), 1, 1_000_000)))
.thenAccept(result -> {
for (LogEntry entry : result.entries()) {
try (entry) {
System.out.println(entry.offset());
}
}
});createStream resolves to a StreamMetadata snapshot, not a resource-owning handle (see Catalog). The stream's namespace is "default", not "public/default": the Storage Spec's Catalog semantics require that, in format version 3, neither a stream namespace nor a stream name contain /. The writer and reader are opened explicitly off the catalog and must be closed before it, and payload must be released once write's future completes — this sketch omits that cleanup and the payload's release for brevity, but does close each entry the read returns, per buffer ownership. WriteResult carries the LogId and offset the write landed at, which the read call above reuses to fetch the same entry back; a real consumer would more often track read position with a LogCursor (see Cursors) than round-trip a single write like this.