Lakestream
Specification

Implementation status

Which roles, format versions, optional capabilities and policy fields each implementation supports.

This page states which Storage Spec and Materialization Spec roles each Lakestream implementation fills today, which format versions and optional API capabilities it provides, and which policy fields each of its materializers applies. It complements the specifications and the API reference; it is not a conformance suite, and it does not report a certification result.

Implementations

  • Ursa 1.0 — a storage engine that implements the Lakestream API and the Storage Spec. Published on Maven Central as org.openlakestream:*:1.0.0 (for example org.openlakestream:lakestream-api:1.0.0 and org.openlakestream:ursa-storage-core:1.0.0).
  • Ursa for Apache Kafka (UFK) 4.3.1 — a Kafka distribution built on the Lakestream API and specification. It uses Ursa through the Kafka binding rather than reimplementing storage or catalog logic; below, "UFK" refers to it.

Version labels on this page are Ursa 1.0 and UFK 4.3.1. The Kafka binding section below cites the literal source tag its facts were checked against.

Storage Spec roles

Ursa 1.0 fills all four Storage Spec roles:

RoleUrsa 1.0
ReaderYes, with the deviations listed below
WriterYes, with the session-loss note below
CompactorYes, with the session-loss note below
CatalogYes, except that Ursa 1.0 does not reject a namespace or a stream name containing /

Reading entries stored as compacted objects — index records with file_type = PARQUET — requires a compacted-object reader factory configured at catalog startup. When none is configured, Ursa installs a no-op reader factory whose reads fail; only WAL (RAW) reads succeed until a reader factory is configured.

Both reader factories default to floor lookup by secondary-index key, which is the row-lookup requirement the Reader role has to meet: an offset that isn't a row's exact start still resolves to the row at or before it. The standalone Kafka reader factory — the one UFK always installs — defaults its approximate-matching setting to true. The general lakehouse reader factory's underlying configuration setting defaults to false, but the factory itself sets it to true during initialization unless a deployer has already configured it explicitly, so its effective default is also floor lookup. Configuring either factory's setting to exact match switches it away from the required behavior.

Two more Reader gaps against Reading's offset-clamping rules: reading a PARQUET range passes the caller's offset straight to the compacted-object reader without first raising it to the trim marker. The read then either fails, when the offset is below the compacted record's own start, or succeeds and returns entries that end at or below the trim marker — which the Reading rules forbid — when the offset is below the trim marker but still inside the record's range; a RAW read avoids both outcomes by clamping first. And the first offset Ursa reports for a log is the covering record's own start offset, not that offset raised to the trim marker, so a reported first offset can itself be below the marker.

Ursa 1.0 follows neither session-loss recommendation in Fences and log deletion, for writers or for compactors. Its per-log write leases have no session-loss handling: when the metadata-store session that created a lease record ends, a writer goes on publishing index records under that lease, and a compactor goes on with the compaction index update and hard trims under it, until it closes the lease.

UFK does not fill these roles independently — its diskless topics read and write through Ursa, so its storage-role conformance follows Ursa's. Its own entry payloads meet the Kafka binding's requirements.

Format versions

Ursa 1.0 writes format version 3 by default. When configured with a format version below 3, it writes version 2 index values even when configured for version 1. It can also read version 1 and version 2 WAL objects and offset-index values — the legacy formats — when configured to do so:

Format versionUrsa 1.0 writesUrsa 1.0 reads
v1NoOnly if configured
v2Only if configuredOnly if configured
v3Yes (default)Yes (default)

Version selection is entirely a matter of configuration, not something Ursa detects from bytes on read. The v3 WAL object layout carries its own version marker, and the v3 offset-index value carries its own version field, but Ursa writes both without checking either back on read. A deployment's readers, writers and compactors must be configured to agree on one format version; nothing in the running system enforces that agreement. UFK does not set a format-version property when it configures Ursa, so its brokers always run Ursa at format version 3, the default.

Optional API capabilities

The API leaves several capabilities optional: some as Log, LogStorage, LogCursor or StreamCatalog methods with a default implementation an override can replace, others as enum values or behaviors an implementation may or may not support. This is what Ursa 1.0 provides for each one.

CapabilityUrsa 1.0
Cluster-wide default materialization (set and get)Yes, held in memory; not persisted across restarts
openLog by partition indexYes
Durable cursors (openCursor, loadCursor)No — both return the API's default failed future
Ephemeral cursors (openEphemeralCursor)Yes
Retention trim offset (computeRetentionTrimOffset)Yes
Binary search by entry header (binarySearchOffset)Yes
Non-blocking closeAsyncYes, on catalog-opened logs: returns before the close finishes and retries under supervision until the log and its write lease are released
The optional LogCursor methods (acknowledgement, filtered reads, state persistence, backlog, cursor deletion)Implemented, but reachable only through ephemeral cursors; persisting an ephemeral cursor's state is a no-op, since its state is kept in memory only
readEntriesByIndex, preFetchEntries, getFirstOffset with includeTrimmedYes, all three
INDEXED partitioningYes
RANGE partitioningNo — stream creation fails
Lifecycle states producedOnly CREATING, ACTIVE and DELETING. SEALED and TRUNCATING are never produced
deleteLog through a catalog-opened logNo — fails; delete a stream's logs through dropStream instead
Durable dropStream fencingYes — see the Storage Spec's write leases and fences

Deviations from the API contract

Ursa 1.0 differs from the documented contract in Logs and Catalog in the following ways:

API contractUrsa 1.0
LogStorage.getFirstOffset and LogStorage.getLastOffset on an empty log complete the returned future exceptionally.Both complete normally, carrying LogOffset.NOT_FOUND.
registerTableCatalog fails with AlreadyExistsException when a catalog of that name is already registered.Overwrites the existing registration, with a logged warning, and succeeds.
Log.readEntry resolves to the entry at the given offset; the contract does not address a missing one.Ursa's implementation completes with null.
Log.getMessageCount(startOffset, endOffset) counts the messages in that range.Ignores endOffset and returns the message count of the single index entry found at startOffset.
Log.readEntries(..., includeTrimmed) includes soft-trimmed entries when includeTrimmed is true.Ignores the flag; a RAW read never returns entries below the trim marker regardless of its value.
Log.softTrim and LogStorage.softTrim resolve to the first entry's offset after trimming.Retries a trim-marker write that conflicts with another trim up to three times. If the write still fails, or fails for another reason, both resolve to the boundary that soft trim computed anyway, so the returned value does not show that the marker advanced.

Materialization

Two roles: Resolver and Materializer. A materializer declares which policy fields it applies; in Ursa, fields it doesn't apply are not read at all.

Resolver

Ursa 1.0 resolves every stream's materialization policy with lakestream-api's own resolution code (TableMaterializationPolicy.resolve) — the code that ships in the API artifact and that the Materialization Spec's resolution rules bind. Ursa does not carry a separate resolver of its own.

Ursa adds one layer the resolution rules don't define: an optional, cluster-wide default policy, which Ursa's compactor sets from configuration at every start rather than persisting it. It stands in for a missing namespace policy as a whole; it is not merged field by field with one. Ursa also diverges from the evaluation order: its stream catalog looks up the effective catalog reference's table catalog before running the disabled check, so an error reading that table catalog fails resolution even for a disabled stream, ahead of where the specified order would already have ended with not materialized. See the Materialization Spec's implementation notes for both.

Materializers

Ursa's materializers are separate modules behind a service-provider interface; see Materialization SPI for how Ursa selects and runs one. This matrix states which policy fields each one applies.

Policy fieldIcebergDeltaDelta/UCClickHouseNONE
catalogRefnot appliednot appliednot appliednot applied—
tableNamingnot appliednot appliednot appliednot applied—
tableIdentifierappliedappliedappliedapplied—
enablednot appliednot appliednot appliednot applied—
connectionOverridesappliedappliedappliedapplied—
primaryKeynot appliednot appliednot appliedapplied—
writeModenot appliednot appliednot appliedapplied, with deviations (see below)—
commit.batchSizenot appliednot appliednot appliedapplied—
table.* (partitionBy, sortBy, retention, targetFileSizeBytes, compression)not appliednot appliednot appliednot applied—
evolutionnot appliednot appliednot appliednot applied—
startPositionnot appliednot appliednot appliednot applied—
pausednot appliednot appliednot appliednot applied—
errorHandlingnot appliednot appliednot appliednot applied—
commit.maxRetriesnot appliednot appliednot appliednot applied—
commit.retryDelayMsnot appliednot appliednot appliednot applied—
baseSchemaVersionnot appliednot appliednot appliednot applied—

No materializer applies catalogRef: resolution fully consumes it before a materializer runs, to decide which table catalog — and so which materializer — handles the stream. The same is true of tableNaming: resolution already derives tableIdentifier from it, and a materializer writes to that resolved identifier directly. None applies enabled either; the Materialization Spec in fact requires a materializer not to read meaning into the effective enabled field at all.

A stream whose effective catalog has type NONE registers no external table catalog — compacted objects are its only output — so no materializer runs and none of these fields apply.

writeMode chooses ClickHouse's destination table engine: UPSERT or CDC selects ReplacingMergeTree, relying on the engine to deduplicate by its ORDER BY key — and so does a non-empty primaryKey alone, even under APPEND, with no writeMode set, or with no policy at all. MergeTree, the append-only family, is what's left: APPEND, an absent writeMode, or an absent policy, each with no primary key set. A stream configured with APPEND and a primary key therefore does not write every record as a new row, which is what the Materialization Spec's APPEND requirement describes, so the matrix marks writeMode "applied, with deviations" for ClickHouse rather than plain conformance. CDC selects the same engine as UPSERT but does not otherwise apply insert, update and delete per operation code the way the spec's CDC requirement describes: its schema-aware row encoder rejects a null value — one common way a source signals a delete — as an error rather than deleting the row; a fallback path that decodes Kafka records without a schema registry configured skips a null value silently instead. Either way, this engine choice is only what commit metadata records as assumed; the materializer itself never issues DDL to create or alter a ClickHouse table, which is the schema service's separate responsibility.

Where a materializer needs a setting the policy doesn't supply to it, Ursa takes it from writer configuration instead — task properties carried alongside the resolved policy, separately from the policy record and applied after it, so they win over it. Iceberg and Delta take a partition key, an identifier-fields list, an upsert-mode flag and a base-schema-version number this way, rather than from the policy's table.partitionBy, primaryKey, writeMode or baseSchemaVersion. A task property can even override the catalog name that the resolved materialization otherwise supplies, when that property is present and non-blank. That diverges from the Materialization Spec's rule that a materializer write a materialized stream only to the table that its resolved materialization names: the table identifier, in the resolved table catalog.

Ursa's compactor also diverges from the Materialization Spec's rule that a materializer not write an external table for a stream that resolves to not materialized. When resolution ends with not materialized, the compactor resolves the stream again from its deployment configuration merged with the compaction task's properties, a compatibility path for deployments that configure materialization there rather than in policies. If those configure an external destination, Ursa writes an external table for the stream, even when its stream policy sets enabled to false. If they configure none, it uses a NONE table catalog named internal-compaction, and writes no external table.

Dead-letter handling is likewise not driven by the policy's errorHandling.dlqTopic. Iceberg always writes a dead-letter table. Delta and Delta/UC write one only when a separate writer-configuration flag enables it. ClickHouse has no dead-letter path at all: a row it can't write fails the materialization task. The compaction scheduler every materializer runs under also reports no materialization state — its metrics handler is a no-op, and there is no API to query a stream's materialization state.

Kafka binding

These facts describe UFK 4.3.1, which depends on Ursa 1.0. They were checked against the source of openlakestream/kafka at tag v4.3.1.1.

UFK identifies each Kafka topic incarnation as the stream default/{topic}-topic-id-{topic ID}, with one log per partition. Including the topic ID keeps a deleted-and-recreated topic of the same name from attaching to the previous incarnation's log when best-effort cleanup of the old one hasn't finished yet.

Creating or reconciling a topic's stream writes source metadata onto it: lakestream.source.logical.name and lakestream.kafka.topic.name (both the topic name), lakestream.kafka.topic.id, lakestream.kafka.source.revision, and an internal ownership marker. Topic-config keys that collide with these — including any already prefixed lakestream.kafka., and three more reserved for materialization bookkeeping — are filtered out rather than copied onto the stream.

Transactional produce is rejected with INVALID_REQUEST before it reaches storage. Idempotent producers — a batch carrying a producer ID without the transactional flag — are accepted normally. Independently of that check, payload validation also fails control batches and transactional batches outright.

UFK depends on several of the optional capabilities above:

  • ephemeral cursors, opened per partition reader;
  • filtered reads bounded by an exclusive maxOffset, for fetch requests;
  • binarySearchOffset, to seek close to a target timestamp before scanning forward at most 256 entries;
  • computeRetentionTrimOffset followed by softTrim, to turn retention.ms and retention.bytes into a soft trim;
  • non-blocking closeAsync, so closing a partition's log doesn't block its caller until the write lease is released.

Deleting a topic calls dropStream with purge set to true, which also erases the stream's data rather than only its catalog metadata. The call itself relies on the durable dropStream fencing capability above, which the API requires of any catalog storage the operation runs against, independent of the purge argument. See Limitations for UFK's operational boundaries beyond this page's scope.

Mixed-writer deployments

The Storage Spec requires only that WAL object names be unique within a deployment and never reused; naming convention and the mechanism for reclaiming them are both implementation-defined.

Ursa 1.0 names WAL objects with a date-and-time prefix, yyyy/MM/dd/HH/mm/ss/{UUID}, in the writer's local time, and its WAL garbage collector assumes every object it might delete follows that scheme. Each cleanup pass derives a deletion boundary from the oldest uncompacted object name across every stream; a name that doesn't parse in that format stops the entire pass rather than skipping just that one object.

A deployment where more than one implementation writes WAL objects into the same WAL root must therefore either have every writer use Ursa's naming scheme with the same time zone as Ursa's own writers, or disable Ursa's WAL garbage collector and reclaim WAL objects some other way. The name carries no time zone of its own, and the garbage collector compares prefixes as plain date-and-time values, so writers in different zones can make it reclaim an object a reader still needs, or delay reclaiming one longer than necessary.

Not established by implementing the API

Implementing the Lakestream API does not, by itself, establish:

  • cross-log transactions, or any transaction spanning more than one log;
  • generic retry deduplication or exactly-once delivery;
  • distributed single-writer ownership — Log.fence() blocks further appends on a log, but is not a leader-election or ownership protocol;
  • RANGE partition split or merge support, from the existence of the RANGE enum value;
  • atomic publication across the metadata store and every table sink a stream materializes into;
  • Kafka transaction support, from using Ursa as Kafka's diskless storage layer.

There is no cross-implementation certification program behind this page, and no published, passing conformance suite. See Verification for the separate, non-certifying verification work: protocol model checking and black-box runs against a running system.