Lakestream
API

API

The Lakestream Java API is published as lakestream-api 1.0.0; this reference documents its three packages, bootstrapping, threading and errors.

This section is the Java API reference for Lakestream: every public type in lakestream-api, its contract as declared by the interface and its Javadoc, and the default behavior of every default method. It describes the API's contract, not any one implementation's behavior. Where Ursa 1.0 differs from the contract, or leaves an optional capability unimplemented, that is noted with a link to Implementation status rather than documented here as the contract.

Artifact and packages

The API is published as org.openlakestream:lakestream-api:1.0.0, built for Java 17, with one compile-scoped dependency, netty-buffer. It has three packages:

PackageContents
io.lakestream.apiThe catalog, log, and stream interfaces, their records and enums, and the bootstrap helpers
io.lakestream.api.materializationThe stream materialization framework's policy model — records and enums, plus the resolution and naming-template code that operates on them
io.lakestream.api.exceptionThe exception hierarchy

Older design documents and examples using io.streamnative.lakestream.api, or a resource-owning Stream handle, do not describe this API.

The package documentation calls LogStorage "Level 1" and StreamCatalog "Level 2".

Types documented on this page (9): the bootstrap helpers StreamCatalogLoader and StreamCatalogProvider, and the seven types in io.lakestream.api.exception.

Bootstrapping

final class StreamCatalogLoader {
    static StreamCatalog open(String catalogMetadataUri, Properties properties) throws Exception;
    static StreamCatalog open(String catalogMetadataUri, Properties properties, ClassLoader classLoader) throws Exception;
}

interface StreamCatalogProvider {
    StreamCatalog open(String catalogMetadataUri, Properties properties) throws Exception;
}

StreamCatalogLoader.open(uri, properties) uses the current thread's context class loader; if that is absent, it falls back to the class loader that loaded StreamCatalogLoader itself. The three-argument overload takes an explicit class loader instead.

Either overload requires exactly one StreamCatalogProvider to be visible to that class loader through ServiceLoader:

ConditionResult
No provider foundIllegalStateException("No StreamCatalogProvider found")
More than one provider foundIllegalStateException("Multiple StreamCatalogProvider implementations found")
The provider's open returns nullIllegalStateException naming the provider class

properties is copied before it reaches the provider, so the caller's Properties object is never mutated or retained by reference. A StreamCatalogProvider implementation owns every resource it creates and transfers that ownership to the StreamCatalog it returns — the catalog is responsible for closing those resources when it is closed. Ursa registers one provider, UrsaKafkaStreamCatalogProvider, so the "exactly one" rule is satisfied by construction in that packaging.

Thread safety

StreamCatalog, Log, LogCursor, LogStorage, StreamLayout, StreamReader, and StreamWriter all declare that implementations must be safe for concurrent use. LogEntry and LogEntryHeader declare weaker guarantees, worded as should rather than must, and not identical to each other: LogEntryHeader says implementations should be immutable and safe for concurrent reads, while LogEntry says only that implementations should be safe for concurrent reads, with no immutability claim.

LogStateManager's methods are individually thread-safe, but a check followed by an action on its state is not atomic without extra synchronization supplied by the caller — reading FENCED and then acting on it can race with a concurrent setState.

Ephemeral cursors are meant to be opened one per read, not pooled: allocating one is documented as allocating in-memory state only, and the Javadoc on openEphemeralCursor says callers should open one per read rather than pool them.

CompletableFuture conventions

I/O on these interfaces returns CompletableFuture. One default breaks that pattern outright: StreamCatalog.setClusterDefaultMaterialization's default implementation throws UnsupportedOperationException directly, on the calling thread, instead of returning a future at all. Three other defaults conform to the pattern in an easy-to-miss way: Log's default openCursor, openEphemeralCursor, and loadCursor each return an already-failed future carrying UnsupportedOperationException, so the call shape looks exactly like a real asynchronous operation even though nothing asynchronous happens.

Log.closeAsync() never throws — a close failure is reported by completing the returned future exceptionally, never as a thrown exception.

Not every method returns a CompletableFuture: StreamCatalog.name(), clusterDefaultMaterialization(), and several LogCursor accessors (see Cursors) are synchronous by signature. Two of the synchronous methods stand out because their work looks like it could need I/O: Log.getMessageCount returns long directly (-1 if the count is not available from cached data), and LogCursor.getNumberOfEntriesInBacklog also returns long directly. No interface promises a synchronous method won't block.

Several methods that return CompletableFuture also declare @throws clauses in their Javadoc — for example, StreamCatalog.createStream's AlreadyExistsException, or StreamLayout.resolveForWrite's IllegalArgumentException for an out-of-range routing hint. The interfaces do not say whether such an exception is thrown synchronously before any future is returned, or delivered by completing the future exceptionally; treat either as possible unless a specific implementation documents which it does.

Close order

StreamCatalog's Javadoc states directly that every data-plane handle it opens — a log, a reader, or a writer — must be closed before the catalog itself is closed, and that a catalog-opened reader additionally owns its own lazily opened child logs. Cursors are one level further down: they are opened from a Log (see Cursors), not from the catalog directly. The interfaces do not state an order between a cursor and its log as directly as they do for the catalog, but a LogCursor holds a reference to its Log (LogCursor.log()), so the same shape applies one level down: close cursors before the log they read from.

Exceptions

Every exception is an unchecked RuntimeException. Most signal a catalog-operation failure; LogFencedException and its subtype signal a log-fencing failure instead.

ExceptionExtendsSignals
AlreadyExistsExceptionRuntimeExceptionA namespace, stream, or table catalog with that name already exists
NoSuchNamespaceExceptionRuntimeExceptionA referenced namespace does not exist
NamespaceNotEmptyExceptionRuntimeExceptionA namespace drop was requested but the namespace still contains streams
NoSuchStreamExceptionRuntimeExceptionA referenced stream does not exist. Declares protected constructors for the subclass below
StreamPermanentlyDeletedExceptionNoSuchStreamExceptionThe stream identity carries a durable permanent-deletion tombstone; catch this before the broader NoSuchStreamException if the distinction matters
LogFencedExceptionRuntimeExceptionA log can no longer accept mutations because its writer was fenced
PartitionLifecycleFencedExceptionLogFencedExceptionA native partition allocation lifecycle is fenced by retained metadata that cannot be reconciled with the partition's current metadata and ownership generation

Buffer ownership

Payloads on the data plane are Netty ByteBuf, not ByteBuffer. In short: on append, the caller owns its reference until the returned future completes and must then release it exactly once; on read, the caller must close every returned entry, and an entry's payload() is a borrowed, read-only view that must not be released directly. See Logs: Buffer ownership for the full contract and a worked example.

io.lakestream.api.materialization

io.lakestream.api.materialization is not a separate SPI — it is part of this API, covering the stream materialization framework's policy model: the records and enums that state which streams become which tables, in what write mode, with what partitioning, evolution, and commit behavior. See Materialization types for the Java types, and the Lakestream Materialization Spec for what they mean.