Writing a materializer
Implementing a sink back-end against Ursa's materialization SPI, from factory to commit to registration.
A materializer is the sink back-end that turns decoded records into writes against some destination. Ursa ships four — Iceberg, Delta, Delta on Unity Catalog, and ClickHouse — and the interface they implement is open.
This page is the how-to. For the full type inventory, see the materialization SPI; for the policy vocabulary that decides which stream goes where, see the Materialization Spec.
A new sink needs a change to Ursa
The framework dispatches on TableCatalogType, which is a closed enum in lakestream-api. A destination that is not one of the existing values cannot be selected by policy without adding one, so a genuinely new sink type requires a pull request to openlakestream/ursa — and, because the enum is part of the Lakestream API, a LIP.
You can implement a materializer for an existing type without either.
The two interfaces
A sink is a factory and the materializer it builds.
public interface TableMaterializerFactory {
TableCatalogType catalogType();
TableMaterializer<?> create(
TableMaterializationPolicy policy,
TableCatalog resolvedCatalog,
StreamMetadata streamMetadata,
MaterializationRuntime runtime);
@Nullable
TableSchemaService<?, ?> schemaService(
TableMaterializationPolicy policy,
TableCatalog resolvedCatalog,
StreamMetadata streamMetadata);
}public interface TableMaterializer<R> extends AutoCloseable {
void write(R record, MaterializationContext context);
CommitResult commit();
void close();
EvolutionPolicy supportedEvolutions();
}R is your sink-side record type. The framework decodes storage entries into it through the serde registered for the format, so the materializer never parses bytes.
Lifecycle
The framework creates one materializer per task and drives it:
create → write* → commit → write* → commit → … → closeA materializer is not shared between tasks and does not need to be thread-safe across them. close() runs whether or not the task succeeded.
Implementing create
create is handed everything needed to resolve the destination:
policy— the resolvedTableMaterializationPolicy, including table configuration and framework settings.resolvedCatalog— theTableCatalogrecord: its name, type, connection map and properties. Connection details belong here rather than in global configuration.streamMetadata— the source stream, including the properties a per-topic override was set on.runtime— the framework services below.
MaterializationRuntime carries the schema service, the schema evolution manager, an executor, a logger, metrics, the failure handler, and the task properties. Take what you need from it rather than constructing your own; the metrics and failure handler in particular are how your sink becomes observable without depending on a metrics library.
Build the destination connection here, but keep it cheap. create runs per task.
Implementing write
void write(R record, MaterializationContext context);context gives the stream identifier, the record's offset and timestamp, an optional source schema version, and a map of source metadata such as Kafka headers. Carry the offset through if your destination can store it — it is what makes a table row traceable to a log position.
write takes ownership of the record, including when it rejects it. A record that cannot be written must be released, not dropped on the floor. Route it through the failure handler:
runtime.failureMessageHandler().sendFailureMessage(
new FailureRecord(stream, catalogType, dlqTopic, reason, payload));Ownership of the payload buffer transfers to the handler. The returned future completes once the record is durably accepted.
Buffer rather than writing through. commit is where durability is established, and a sink that writes per record cannot make the commit atomic.
Implementing commit
commit() publishes everything buffered since the last commit and returns a CommitResult of recordsCommitted, bytesCommitted and an opaque sinkMetadata map.
Two properties matter:
- Atomic. Either the whole batch becomes visible or none of it does. Compaction records progress on the strength of a successful commit; a partial commit that reports success leaves the table disagreeing with the offset the compactor believes it has reached.
- Idempotent under retry. A commit can be retried after an ambiguous failure. A sink that cannot make a repeated commit a no-op will duplicate rows.
Throw on failure rather than returning a short count. The framework classifies the exception and quarantines the task; see operations.
Declaring supported evolutions
EvolutionPolicy supportedEvolutions();Return what your destination can actually absorb. The framework gates incoming schema changes against it, so a sink that cannot widen a type should say so rather than failing at write time.
EvolutionPolicy has factories for the common shapes: forIceberg() and forDelta() are permissive about added and widened columns, forClickHouse() allows additions alone.
Return a schema service from schemaService(...) if your sink evolves the destination schema, or null if it does not.
Registration
The framework discovers factories with ServiceLoader and indexes them by catalogType(). Add a file to your module:
META-INF/services/io.lakestream.ursa.materialization.TableMaterializerFactorycontaining your factory's fully-qualified class name.
One factory per catalog type. If two are found for the same type, the first wins and the duplicate is logged as a warning — so a shadowed sink fails quietly rather than loudly. Keep your module off the classpath when you do not intend to override.
Worked reference
ClickHouseTableMaterializerEndToEndTest in ursa-storage-clickhouse exercises a complete sink from factory through write to commit, and is the closest thing to a template.
It is test-scoped and therefore not part of any published artifact — read it, do not depend on it.