Lakestream
Ursa

Materialization SPI

The service-provider interface Ursa uses to run materialization, and where its materializers live.

Lakestream's stream materialization framework defines how a stream becomes a lakehouse table. This page describes the machinery Ursa uses to execute that: the service-provider interface sink back-ends implement, and where the materializers built on it live.

SPI, not spec

The policy model — the vocabulary for saying which stream becomes which table, and under what settings — is part of Lakestream. The SPI on this page, and the materializers behind it, are Ursa's, not part of the Lakestream specification: an implementation could satisfy the policy model with entirely different execution machinery. Implementation status lists, per materializer, which policy fields it applies.

ursa-storage-materialization holds the service-provider interface in io.lakestream.ursa.materialization. Sink back-ends implement it; the compaction orchestrator drives it.

Service and work units

TypeKindContent
MaterializationServiceInterface, AutoCloseableOne instance per deployment. initialize(MaterializationRuntime, MaterializationServiceConfig) runs once at scheduler startup; materialize(MaterializationTask) runs per task; invalidate(StreamIdentifier) drops cached state on failure or stream delete. A default resolveFromTaskProperties(...) returns empty and exists for deployments that drove materialization through task properties rather than catalog policy.
MaterializationServiceProviderFinal utility classload(className) reflectively instantiates a service through its public no-arg constructor, rethrowing any reflective failure as IllegalStateException so startup fails fast. The caller initializes the returned instance.
MaterializationServiceConfigRecordworkerPoolSize, walReadRateLimitWindow (a Duration), walReadRateLimitBytes, and an additionalProperties escape hatch. Sizes must be positive.
MaterializationTaskRecordOne unit of work: streamMetadata, resolvedMaterialization, sourceTopic, a numeric streamId, an inclusive startOffset, an exclusive endOffset, and a nullable sourceTask. The service reads the source entries for that range itself; the orchestrator does not carry entries.
MaterializationRuntimeRecordThe framework services injected into sinks: schemaService, schemaEvolutionManager, materializationExecutor, logger, metrics, failureMessageHandler, a nullable compactTaskManager, a nullable storageApi, and taskProperties. withStorageApi(...) and withTaskProperties(...) return copies. Its SOURCE_TOPIC_PROPERTY constant is ursa.materialization.source.topic.
MaterializationContextRecordPer-record context: stream, offset, timestamp, an optional sourceSchemaVersion, and a sourceMetadata map for source-format metadata such as Kafka headers.

Sink contract

TableMaterializer<R> is the pluggable sink interface, parameterized by the sink-side record type. Its lifecycle is write* → commit → write* → commit → … → close. write(R, MaterializationContext) buffers a record and takes ownership of it, including on rejection; commit() returns a CommitResult; supportedEvolutions() returns the EvolutionPolicy the framework gates incoming evolutions against.

TableMaterializerFactory builds those materializers. Each TableCatalogType has at most one factory, and the framework dispatches on catalogType() across ServiceLoader-discovered factories. create(policy, resolvedCatalog, streamMetadata, runtime) returns a materializer; schemaService(...) returns the sink's schema service, or null when the sink performs no schema evolution.

CommitResult is a record of recordsCommitted, bytesCommitted, and an opaque sinkMetadata map.

Failure handling and observability

FailureMessageHandler forwards unwritable records to a dead-letter sink through sendFailureMessage(FailureRecord), returning a future that completes once the record is durably accepted; noop() supplies a stub. FailureRecord is a record of stream, catalogType, an optional dlqTopic, a reason, and a payload whose buffer ownership transfers to the handler.

MaterializationMetrics is a small interface rather than a registry binding, so sinks need no observability dependency: recordWritten, recordCommitDuration, recordCommitRetry, recordSchemaEvolutionApplied, recordSchemaEvolutionRejected, recordDlqRecord, and setState. noop() supplies a stub.

Decoding and schemas

The serde sub-package converts storage entries into sink records:

  • EntryEncoder<T> converts one storage entry into one or more materialization records; ownership of the entry transfers to the encoder on every path.
  • EntryDecoder<T> runs the other direction, from materialization records back to a generic entry.
  • EntrySerdeFactory is the registry keyed by a SerdeType: KAFKA_ICEBERG, KAFKA_DELTA, KAFKA_PARQUET, KAFKA_BATCHED_RAW_PARQUET, KAFKA_CLICKHOUSE. Integration modules register their providers, which keeps the generic serde framework free of table-format dependencies.
  • MaterializationRecord<T> pairs a sink record with optional entry metadata.
  • SchemaService<T> resolves source-side schemas by topic and version.
  • SchemaEvolutionManager evolves the destination table schema as the source schema advances.

Flow

  1. Policy resolution. The stream catalog resolves a namespace baseline and stream overrides into a ResolvedMaterialization, or into an empty result when the stream opts out, when the effective policy has no catalog reference, or when that reference names a table catalog that is not registered. When the result is empty, Ursa's compactor resolves the stream again from its deployment configuration and the compaction task's properties, and writes an external table if those configure one; this departs from the Materialization Spec's rule for streams that resolve to not materialized, and Implementation status lists it as a deviation.
  2. Task publication. Compaction produces the work. Each WAL task that resolves to a materialization-enabled stream becomes a MaterializationTask covering one [startOffset, endOffset) range of one source log.
  3. Writing. The service builds a TableMaterializer through the factory registered for the resolved catalog type, decodes the entries in the range, and writes records into it.
  4. Commit. commit() publishes the buffered output to the destination table catalog. Until that commit lands, the data is readable from the log but not visible in the table.

See lakehouse tables for how internal compacted objects — defined normatively in the Storage Spec — differ from external delivery in Ursa.

The materializers themselves live outside this module. ursa-storage-lakehouse registers the Iceberg, Delta, and Delta/Unity Catalog factories; ursa-storage-clickhouse registers the ClickHouse factory. Each registers through META-INF/services/io.lakestream.ursa.materialization.TableMaterializerFactory in its own module.