Lakestream Materialization Spec
The Materialization Spec defines the policy model of the stream materialization framework: policies, resolution, table naming and table catalogs.
This is a specification for the policy model of the Lakestream stream materialization framework, which declares how a stream is materialized into a lakehouse table.
The key words "MUST", "MUST NOT", "REQUIRED", "SHALL", "SHALL NOT", "SHOULD", "SHOULD NOT", "RECOMMENDED", "NOT RECOMMENDED", "MAY", and "OPTIONAL" in this document are to be interpreted as described in BCP 14 [RFC 2119] [RFC 8174] when, and only when, they appear in all capitals, as shown here.
Versioning
This specification describes the policy model in lakestream-api 1.0.0 (Maven org.openlakestream:lakestream-api:1.0.0, package io.lakestream.api.materialization). It has no version number of its own. Changes are made through Lakestream Improvement Proposals (LIPs) and ship in a lakestream-api release.
The serialized form of a policy is implementation-defined. It is not part of this specification, and neither is the serialized form of a table catalog. The format version of the Lakestream Storage Spec does not apply to this specification.
Goals
- Declarative -- A policy states which stream becomes which table, and how. Resolving a policy performs no I/O beyond looking up the named table catalog.
- Layered policies -- A namespace policy sets a baseline for the streams in its namespace; a stream policy overrides it by fixed rules.
- Destination-neutral -- A policy names a registered table catalog by reference, so one policy model covers Iceberg, Delta, ClickHouse and storage-only output.
Overview
A policy declares whether a stream is materialized into a table, which table, and with which settings. Policies attach at two layers. A namespace policy is the baseline for every stream in a namespace, and a stream policy overrides it for one stream. Either layer can be missing.
Resolution combines the two layers for one stream. It looks up the table catalog that the policies name, derives the table identifier, and merges the stream policy's settings over the namespace policy's. It ends in one of three outcomes: a resolved materialization, which holds the table catalog, the table identifier and the effective policy; not materialized; or a failure.
namespace policy --+
stream policy --+--> resolution <-- table catalogs
stream properties --+ |
+--> not materialized
+--> failure
+--> resolved materialization
|
v
materializer --> tableA materializer writes a resolved stream into its table. This specification defines what each policy field means to a materializer that applies it. It does not define how a materializer reads the stream, converts its records, or commits them to the destination. Each materializer states which fields it applies, and Implementation status lists them.
Materialization is separate from storage compaction. The Lakestream Storage Spec defines compacted objects, the files that hold a stream's records after compaction. These internal compacted objects are never registered in a table catalog. A table catalog of type NONE selects storage-only compaction: internal compacted objects, and no external table.
Non-goals
- A resolved policy is not a commit. Resolution reports a destination and the effective settings. It does not show that a materializer applies every setting, that files were written, or that a table snapshot exists.
- An external table is a separate output. An external table has its own files and lifecycle, apart from the internal compacted objects, even when one compaction pass writes both. It is not zero-copy, and it shares no files with them.
- Execution is not specified. How a materializer reads a stream, decodes its records and commits them is implementation-defined. Ursa's execution is described in Materialization SPI.
- Compaction is not key-based. "Compaction" in this specification means rewriting a stream's records into compacted objects, as the Storage Spec defines. It is not Kafka's key-based log compaction, which keeps only the latest value for each key.
Specification
Terms
- Stream -- A stream in a stream catalog, identified by its namespace and its name.
- Stream catalog -- The catalog that holds namespaces and streams, their policies, and the registered table catalogs.
- Stream properties -- The map from string to string that a stream carries. Resolution reads two reserved keys that source integrations set, described in Source logical name, and any key that a naming template names with
${stream.property.<key>}. - Policy -- The structure defined in Policies.
- Namespace policy -- The policy attached to a namespace. It is the baseline for every stream in the namespace.
- Stream policy -- The policy attached to one stream. Its present fields override the namespace policy, except
tableNaming, which is read from the namespace policy only. - Present, absent -- A field that can be absent is either present, with a value, or absent. An empty list is present, not absent.
- Blank -- A string that is empty, or whose every code point is whitespace: U+0009 to U+000D, U+001C to U+001F, or a Unicode space, line or paragraph separator (general category Zs, Zl or Zp) other than U+00A0, U+2007 and U+202F.
- Effective policy -- The policy that resolution builds by merging the stream policy over the namespace policy.
- Table catalog -- A named destination registered with a stream catalog, such as an Iceberg catalog or a ClickHouse server.
- Catalog reference -- The name of a table catalog, held in a policy's
catalogReffield. - Table identifier -- A namespace and a name that identify a table within a table catalog.
- Naming template -- A rule, set in a namespace policy, that derives a table identifier from a stream.
- Source integration -- A system that creates streams for another protocol, such as a Kafka broker, and sets their reserved properties.
- Source logical name -- The name a source integration gives a stream for user-facing defaults, such as the default table name. It is metadata, not the stream's storage identity.
- Resolution -- The procedure in Resolution. It ends in a resolved materialization, not materialized, or a failure.
- Resolved materialization -- The outcome of resolution for a materialized stream: a table catalog, a table identifier, and an effective policy.
- Not materialized -- The outcome of resolution for a stream whose stream policy disables it, that has no effective catalog reference, or whose effective catalog reference names no registered table catalog. It is not an error.
- Failure -- The outcome of resolution when a policy cannot be applied to a stream, such as a naming template that references an unknown variable. It is an error.
- Apply -- A materializer applies a policy field when the field's effective value changes what the materializer writes, or how it writes it.
Conformance
This specification defines two roles. A system conforms for the roles it implements.
- Resolver -- A system that resolves policies, such as a stream catalog.
- Materializer -- A system that writes resolved streams into their destinations.
Each requirement starts with the role it binds, in brackets. Requirements tagged [Source] bind source integrations, which set the reserved stream properties that resolution reads. A source integration is not a conformance role of this specification.
[Resolver] A resolver MUST implement every requirement tagged [Resolver] in Resolution and Table naming.
Most [Materializer] requirements are conditional. A conditional requirement starts "A materializer that applies" and names a field, and it binds only a materializer that applies that field.
[Materializer] A materializer MUST meet every unconditional [Materializer] requirement, and every conditional one for each field it applies.
[Materializer] A materializer MUST state, for each policy field, including the fields of nested structs, whether it applies that field.
This specification does not define how a materializer treats a policy that sets a field the materializer does not apply.
Implementation status lists each implementation's roles and the fields that each of its materializers applies.
Where policies attach
A stream catalog holds at most one namespace policy for each namespace, and at most one stream policy for each stream.
- Namespace policy -- Held with the namespace. It can be given when the namespace is created. It applies to every stream in the namespace, including streams created after it is set.
- Stream policy -- Held with the stream. It can be given when the stream is created.
Both can be set and cleared later. Appendix A lists the operations. When a change reaches resolution depends on whether the resolver caches policies, which this specification does not define; Appendix B describes Ursa's behavior.
A stream with neither policy is not materialized, because no layer names a table catalog.
lakestream-api 1.0.0 also lets a stream catalog hold a cluster-wide default policy. This specification does not define it. Appendix B describes how Ursa uses it.
Policies
A policy has the fields below. Each nested struct is defined in the section its description links to. In every field table on this page, the Requirement column marks a field required when it is always present and optional when it can be absent; see present and absent in Terms.
| Requirement | Field | Type | Description |
|---|---|---|---|
| optional | catalogRef | string | Name of the table catalog to write to. |
| optional | tableNaming | struct | Naming template that derives the table identifier. Read from the namespace policy only. See Table naming. |
| optional | tableIdentifier | struct | Explicit table identifier. Read from the stream policy only. See Table naming. |
| optional | enabled | boolean | Whether materialization is enabled. Only false in a stream policy affects resolution: it makes the stream not materialized. |
| optional | framework | struct | Engine-agnostic settings. See Framework configuration. |
| optional | evolution | struct | Schema changes the destination table permits. See Schema evolution. |
| optional | primaryKey | list of string | Primary-key column names. See Table configuration. |
| optional | baseSchemaVersion | int64 | Base schema version for compatibility checks. See Schema evolution. |
| optional | table | struct | Engine-specific table settings. See Table configuration. |
| required | connectionOverrides | map from string to string | Per-stream overrides of the table catalog's connection settings. Can be empty. Read from the stream policy only. See Table catalogs. |
Resolution
Resolution decides, for one stream, whether the stream is materialized and, if it is, into which table and with which settings.
Inputs
Resolution has five inputs:
- the namespace policy of the stream's namespace, if there is one;
- the stream policy, if there is one;
- the stream's namespace and name;
- the stream properties;
- a lookup that returns the table catalog registered under a given name, if there is one.
[Resolver] A resolver MUST treat a missing namespace policy or stream policy as a policy with an empty connectionOverrides and every other field absent.
A resolver that resolves a stream held in a stream catalog takes the policies and the stream properties from that catalog. This specification does not define whether they can come from a cache.
[Resolver] A resolver that resolves a stream held in a stream catalog MUST NOT open the stream's logs, readers or writers to gather its inputs.
[Resolver] Apart from the table catalog lookup, a resolver MUST NOT perform I/O while it evaluates the steps below.
Evaluation order
[Resolver] A resolver MUST produce the outcome that these steps produce when they are taken in order. Resolution ends at the first step that yields an outcome.
- Disabled check. [Resolver] If the stream policy's
enabledis present andfalse, resolution MUST end with not materialized. [Resolver] The namespace policy'senabledMUST NOT make a stream not materialized. - Catalog reference. [Resolver] The effective catalog reference MUST be the stream policy's
catalogRefif it is present, and the namespace policy's otherwise. [Resolver] If neither is present, resolution MUST end with not materialized. - Catalog lookup. [Resolver] If no table catalog is registered under the effective catalog reference, resolution MUST end with not materialized.
- Merge of
table. The stream policy'stableis merged over the namespace policy's, as Merge rules describe. This step always succeeds. - Table identifier. [Resolver] The table identifier MUST be the stream policy's
tableIdentifierif it is present; otherwise the namespace policy'stableNamingapplied to the stream, as Table naming describes, if that is present; and otherwise the default identifier, whose namespace is the stream's namespace and whose name is the source logical name. [Resolver] If applying the naming template fails, resolution MUST fail, and MUST NOT fall back to the default identifier. [Resolver] If the table identifier has an empty namespace or an empty name, resolution MUST fail. - Result. [Resolver] Resolution MUST end with a resolved materialization that holds the table catalog found in step 3, the table identifier from step 5, and the effective policy that Merge rules build.
Because of this order, a stream that is disabled, whose policies name no table catalog, or whose table catalog is not registered is not materialized, even when its naming template would fail.
Merge rules
In this section, a field chosen by stream wins takes the stream policy's value if that is present, even when it is an empty list, and the namespace policy's value otherwise. A struct merged field by field takes each of its fields by stream wins, unless the table below says otherwise, when both layers have the struct. When only one layer has the struct, the effective struct is that layer's.
[Resolver] A resolver MUST build each field of the effective policy as this table states.
| Field | Effective value |
|---|---|
catalogRef | The effective catalog reference from step 2 (stream wins). |
tableNaming | The namespace policy's, present or absent. The stream policy's is ignored. |
tableIdentifier | The table identifier from step 5. The namespace policy's is ignored. |
enabled | Stream wins. |
framework | Merged field by field. Within it, commit is merged field by field, and errorHandling is chosen whole by stream wins: its mode and dlqTopic are never merged separately. |
evolution | Merged field by field. |
primaryKey | Stream wins. The list is replaced whole, never concatenated. |
baseSchemaVersion | Stream wins. |
table | Merged field by field. Within it, retention is merged field by field, and the lists partitionBy and sortBy are replaced whole, never concatenated. |
connectionOverrides | The stream policy's map, even when it is empty. The namespace policy's is ignored. |
[Resolver] Resolution MUST NOT supply default values: apart from tableIdentifier, which step 5 always sets, a field that is absent from both layers MUST be absent from the effective policy. What a materializer does in place of an absent field is its own choice.
The effective tableNaming is kept for inspection only; the table identifier is already derived.
Using a resolved materialization
The effective enabled is chosen by stream wins like other fields, so a resolved materialization can carry enabled set to false by the namespace policy.
[Materializer] A materializer MUST NOT read meaning into the effective policy's enabled field. A resolved materialization means that the stream is materialized, whatever that field holds.
[Materializer] A materializer MUST write a materialized stream only to the table that its resolved materialization names: the table identifier, in the table catalog.
How a materializer maps a table identifier onto a destination's names, for example when the destination does not accept a character that the identifier contains, is materializer-defined. Appendix B describes Ursa's mapping for ClickHouse.
[Materializer] A materializer MUST NOT write an external table for a stream that resolves to not materialized.
This rule does not affect storage compaction: a stream's internal compacted objects are not an external table.
Table naming
A naming template derives the table identifier when the stream policy has no tableIdentifier. It is read from the namespace policy only.
Table identifier
| Requirement | Field | Type | Description |
|---|---|---|---|
| required | namespace | string | Table namespace. Non-empty. |
| required | name | string | Table name within the namespace. Non-empty. |
The table catalog is not part of the identifier. It comes from the policy's catalogRef.
Naming template
| Requirement | Field | Type | Description |
|---|---|---|---|
| optional | tableNamespacePrefix | string | Table namespace of every identifier the template derives. Used verbatim. |
| required | tableNameTemplate | string | Template for the table name. Non-empty. |
A variable reference in a template is ${, then one or more characters other than }, then }. The characters between the braces are the variable name.
variable-reference = "${" variable-name "}"
variable-name = 1*name-char
name-char = %x00-7C / %x7E-10FFFF ; any character except "}"[Resolver] A resolver MUST scan a template once, from left to right, replacing each variable reference with its variable's value and resuming after the reference's closing }, and MUST copy every other character unchanged.
[Resolver] A resolver MUST insert each value verbatim, and MUST NOT scan inserted text for variable references.
[Resolver] A resolver MUST recognize exactly these variable names, compared case-sensitively:
| Variable name | Value |
|---|---|
stream.namespace | The stream's namespace. |
stream.name | The stream's name. |
stream.logicalName | The source logical name. |
stream.property.<key> | The value of the stream property <key>. The key is the rest of the variable name after stream.property., and can contain .. |
[Resolver] Applying a template MUST fail when any of these holds:
- a variable name is not one of those above;
- a
stream.property.<key>variable names a property that is absent from the stream properties, or whose value is blank; - the resulting table name is blank.
[Resolver] The table namespace MUST be tableNamespacePrefix, verbatim, if it is present, and the stream's namespace otherwise. The prefix replaces the stream namespace; it is not prepended to it, and it is not scanned for variables. A present but empty prefix gives an empty table namespace, which makes resolution fail in step 5.
There is no escape sequence. Template text that has the form of a variable reference is always one, so template text cannot express a literal ${stream.name}; only an inserted value, which is never scanned, can contain one. Text such as ${}, or a ${ with no closing }, is not a variable reference and is copied unchanged.
Resolution does not escape or rewrite characters. The stream's namespace, name and properties appear in the table identifier exactly as they are.
| Stream (namespace, name) | Stream properties | Prefix | Template | Table identifier (namespace, name) |
|---|---|---|---|---|
sales, orders | none | absent | ${stream.name} | sales, orders |
sales, orders | none | analytics | ${stream.name} | analytics, orders |
sales, orders | none | absent | ${stream.namespace}_${stream.name} | sales, sales_orders |
default, orders-topic-id-abc | lakestream.source.logical.name = orders | analytics | ${stream.name}_archive | analytics, orders-topic-id-abc_archive |
default, orders-topic-id-abc | lakestream.source.logical.name = orders | absent | ${stream.logicalName} | default, orders |
sales, orders | none | ${stream.name} | literal_name | ${stream.name}, literal_name |
ns, s | none | absent | ${stream.property.missing} | Fails: the property is absent. |
sales, orders | none | absent | ${foo.bar} | Fails: the variable is unknown. |
Source logical name
[Resolver] The source logical name MUST be the value of the stream property lakestream.source.logical.name if it is present and not blank; otherwise the value of lakestream.kafka.topic.name if it is present and not blank; otherwise the stream's name.
lakestream.kafka.topic.name is a compatibility key. Kafka integrations wrote it before lakestream.source.logical.name existed.
| Stream properties | Source logical name |
|---|---|
lakestream.source.logical.name = logical-stream, lakestream.kafka.topic.name = legacy-kafka-name | logical-stream |
lakestream.kafka.topic.name = orders | orders |
| neither key | The stream's name. |
[Source] A source integration MUST set lakestream.source.logical.name itself, and MUST prevent users from overriding it.
The default identifier depends only on the stream's namespace and its source logical name. Two streams in one namespace with the same logical name, such as two incarnations of a source topic that was deleted and created again, resolve to the same default table.
Table catalogs
A table catalog is a named destination registered with a stream catalog. A policy refers to it by name, through catalogRef.
| Requirement | Field | Type | Description |
|---|---|---|---|
| required | name | string | The name that policies reference. Non-empty, and unique among the table catalogs registered with a stream catalog. |
| required | type | TableCatalogType | The kind of destination. |
| required | connection | map from string to string | Connection settings, such as a URI, a warehouse, a DSN, or references to credentials. |
| required | properties | map from string to string | Catalog-level tuning defaults, such as a target file size or feature flags. |
This specification defines no keys for connection or properties. Each destination's materializer defines its own.
Registering, listing and unregistering table catalogs are stream catalog operations; see Appendix A. A policy holds only the table catalog's name, so a stream whose effective catalog reference names no registered table catalog resolves to not materialized (step 3). Unregistering a table catalog therefore leaves every stream that references it not materialized, from the first resolution whose lookup no longer finds it.
TableCatalogType
| Value | Destination |
|---|---|
ICEBERG | An Apache Iceberg catalog. |
DELTA | Storage-only Delta Lake tables, without Unity Catalog. |
DELTA_UC | Delta Lake tables in a Unity Catalog. |
CLICKHOUSE | ClickHouse. |
NONE | No external table: storage-only compaction. The stream's only compacted output is its internal compacted objects. |
[Materializer] A materializer MUST NOT write a stream whose resolved table catalog has type NONE to any external table.
Internal compacted objects are the Storage Spec's compacted objects. They are never registered in a table catalog.
TableCatalogType is a closed enum. Adding a type takes a LIP and a lakestream-api release.
Connection overrides
[Materializer] A materializer that applies connectionOverrides MUST use, for each connection key, the effective connectionOverrides value if the key is present there, and the table catalog's connection value otherwise.
Framework configuration
framework holds engine-agnostic settings.
| Requirement | Field | Type | Description |
|---|---|---|---|
| optional | writeMode | WriteMode | How records are applied to the table. |
| optional | startPosition | StartPosition | Where a newly enabled materialization begins reading the stream. |
| optional | paused | boolean | Whether materialization is paused. |
| optional | errorHandling | struct | How failures are handled. Merged whole. See Error handling. |
| optional | commit | struct | Commit and retry settings. See Commit settings. |
Write mode
WriteMode value | Meaning |
|---|---|
APPEND | Every record becomes a new row. |
UPSERT | A record replaces the row that has the same primary key. |
CDC | Each record inserts, updates or deletes a row, as its upstream operation code says. |
[Materializer] A materializer that applies writeMode MUST, for APPEND, write every record as a new row; for UPSERT, replace the row that has the record's primary key, or add a new row if there is none; and for CDC, insert, update or delete a row as the record's upstream operation code says.
UPSERT and CDC identify rows by the effective primaryKey. This specification does not define either mode when primaryKey is absent or empty, or how a record carries its operation code.
Start position
StartPosition value | Meaning |
|---|---|
EARLIEST | Begin at the oldest record the stream retains. |
LATEST | Begin at the newest position, skipping the records already in the stream. |
OFFSET | Begin at a given log offset. |
TIMESTAMP | Begin at a given timestamp. |
[Materializer] A materializer that applies startPosition MUST begin a newly enabled materialization at the oldest record the stream retains for EARLIEST, and at the end of the stream for LATEST, so that the records already in the stream are not written.
lakestream-api 1.0.0 has no field that carries the offset for OFFSET or the timestamp for TIMESTAMP.
Pausing
[Materializer] A materializer that applies paused MUST NOT write to the table while the effective paused is true.
paused does not affect resolution. A paused stream still resolves to a resolved materialization.
Error handling
| Requirement | Field | Type | Description |
|---|---|---|---|
| required | mode | ErrorMode | How to handle a failure. |
| optional | dlqTopic | string | Dead-letter topic for failed records. |
ErrorMode value | Meaning |
|---|---|
SUSPEND | Suspend the materialization until an operator intervenes. |
SKIP | Skip the failed records and continue. |
LOG | Log the failure and continue. |
[Materializer] A materializer that applies errorHandling MUST handle a record or batch that fails to be written as the effective mode says: for SUSPEND, stop writing the stream until an operator intervenes; for SKIP, skip the failed records and continue; and for LOG, log the failure and continue.
[Materializer] A materializer that applies dlqTopic MUST use the topic that dlqTopic names for every failed record it sends to a dead-letter topic.
This specification does not define which system holds that topic, which failed records go to it, or how dlqTopic combines with each mode.
Commit settings
| Requirement | Field | Type | Description |
|---|---|---|---|
| optional | maxRetries | int32 | Maximum retry attempts for a failed commit. |
| optional | retryDelayMs | int64 | Base retry delay, in milliseconds. |
| optional | batchSize | int32 | Records per commit batch. |
[Materializer] A materializer that applies maxRetries MUST NOT retry a failed commit more than maxRetries times.
[Materializer] A materializer that applies retryDelayMs MUST wait at least retryDelayMs milliseconds before each retry of a failed commit.
[Materializer] A materializer that applies batchSize MUST put at most batchSize records in each commit batch it writes to the destination.
Table configuration
primaryKey and table configure the destination table. The settings in table are engine-specific: their effect depends on the destination.
Primary key
[Materializer] A materializer that applies primaryKey MUST use the listed columns as the destination table's primary key.
Table settings
| Requirement | Field | Type | Description |
|---|---|---|---|
| optional | partitionBy | list of struct | Partition fields. See Partitioning. |
| optional | sortBy | list of struct | Sort fields. See Sort order. |
| optional | retention | struct | Retention settings. See Retention. |
| optional | targetFileSizeBytes | int64 | Target size of the data files, in bytes. |
| optional | compression | Compression | Codec for the data files. |
There is no table ownership setting. This specification does not define which system creates or drops a destination table.
Partitioning
A partition field has these fields:
| Requirement | Field | Type | Description |
|---|---|---|---|
| required | column | string | Source column. Non-empty. |
| required | transform | PartitionTransform | How the partition value is derived from the column. |
| optional | parameter | string | Argument of the transform. |
Apart from EXPRESSION, the PartitionTransform values mirror the names of Apache Iceberg's partition transforms.
PartitionTransform value | parameter | Partition value |
|---|---|---|
IDENTITY | none | The column value. |
BUCKET | The bucket count N | The column hashed into N buckets. |
TRUNCATE | The width W | The column truncated to width W. |
YEAR | none | The year of a timestamp or date column. |
MONTH | none | The month of a timestamp or date column. |
DAY | none | The day of a timestamp or date column. |
HOUR | none | The hour of a timestamp column. |
EXPRESSION | Engine-specific | A custom, engine-specific expression. |
[Materializer] A materializer that applies partitionBy MUST partition the table by each listed field whose transform it supports, deriving the partition value from column as transform describes.
[Materializer] A materializer that applies partitionBy MUST reject a BUCKET, TRUNCATE or EXPRESSION field that has no parameter. Resolution does not check transforms or parameters.
[Materializer] A materializer that applies partitionBy SHOULD reject a field whose transform it does not support.
[Materializer] A materializer that applies partitionBy MUST treat the column name __partition as the stream partition that holds each record, not as a record field. This specification does not define how that partition is represented in the table.
Sort order
A sort field has these fields:
| Requirement | Field | Type | Description |
|---|---|---|---|
| required | column | string | Source column. Non-empty. |
| required | direction | SortDirection | ASC for ascending, or DESC for descending. |
| required | nullsFirst | boolean | Whether nulls sort before other values. |
[Materializer] A materializer that applies sortBy MUST use the listed fields, in list order, as the table's sort order: each field sorts ascending for ASC and descending for DESC, with nulls first if nullsFirst is true and last if it is false.
Retention
| Requirement | Field | Type | Description |
|---|---|---|---|
| optional | snapshotRetentionMs | int64 | How long to keep table snapshots, in milliseconds. |
| optional | maxSnapshots | int32 | Maximum number of snapshots to keep. |
| optional | rowRetentionMs | int64 | How long to keep rows, in milliseconds. |
[Materializer] A materializer that applies maxSnapshots MUST NOT keep more than maxSnapshots snapshots of the table.
snapshotRetentionMs and rowRetentionMs say how long snapshots and rows are kept. lakestream-api 1.0.0 does not say whether each is a minimum or a maximum, from which time a row's age is measured, or how snapshotRetentionMs combines with maxSnapshots.
Data files
[Materializer] A materializer that applies targetFileSizeBytes MUST use it as the target size of the data files it writes.
[Materializer] A materializer that applies compression MUST write its data files with the codec that compression names.
Compression value | Codec |
|---|---|
ZSTD | Zstandard |
SNAPPY | Snappy |
GZIP | gzip |
LZ4 | LZ4 |
UNCOMPRESSED | No compression |
Schema evolution
evolution states which schema changes the destination table permits. Each flag can be absent. true permits the change, and false forbids it.
| Requirement | Field | Type | Description |
|---|---|---|---|
| optional | addColumn | boolean | Adding a new non-null column. |
| optional | addNullableColumn | boolean | Adding a new nullable column. |
| optional | dropColumn | boolean | Dropping a column. |
| optional | widenType | boolean | Widening a column's type, such as int to long. |
| optional | narrowType | boolean | Narrowing a column's type. |
| optional | renameColumn | boolean | Renaming a column. |
| optional | reorderColumns | boolean | Reordering columns. |
| optional | nullabilityRelax | boolean | Changing a non-null column to nullable. |
| optional | nullabilityTighten | boolean | Changing a nullable column to non-null. |
[Materializer] A materializer that applies evolution MUST NOT make a schema change to the table when the effective flag for that change is false.
When a flag is absent, whether the change is permitted is the materializer's choice.
baseSchemaVersion names the base schema version for compatibility checks. lakestream-api 1.0.0 does not define which schema that version refers to, or how a check uses it.
Materialization state
MaterializationState names the runtime state of a materialization.
MaterializationState value | Meaning |
|---|---|
PENDING | Configured, but not yet running. |
RUNNING | Running and writing to the destination table. |
DEGRADED | Running with reduced function, such as elevated retries or limited throughput. |
SUSPENDED | Suspended after repeated failures, and waiting for an operator. |
PAUSED | Paused explicitly by a user or an operator. |
[Materializer] If an implementation reports materialization state, it MUST use these values with these meanings.
lakestream-api 1.0.0 has no operation that reports or queries materialization state.
Appendix A: Java binding
The Java binding is the package io.lakestream.api.materialization in lakestream-api 1.0.0, with a few types and operations in io.lakestream.api. The Java API reference documents every type.
| This specification | Java |
|---|---|
| Policy | TableMaterializationPolicy |
| Resolved materialization | ResolvedMaterialization |
| Table catalog | TableCatalog |
| Table identifier | TableIdentifier |
| Naming template | TableNaming |
framework | FrameworkConf |
framework.errorHandling | ErrorHandling |
framework.commit | CommitConfig |
table | TableConf |
| Partition field | PartitionSpec |
| Sort field | SortColumn |
table.retention | RetentionConfig |
evolution | EvolutionPolicy |
| The enums | TableCatalogType, WriteMode, StartPosition, ErrorMode, PartitionTransform, SortDirection, Compression, MaterializationState |
| Source logical name | SourceMetadataProperties.logicalName and SourceMetadataProperties.LOGICAL_NAME_PROPERTY, in io.lakestream.api |
| Namespace policy, stream policy | Namespace.materialization() and StreamMetadata.materialization(), in io.lakestream.api |
- Presence. A field that can be absent is an
Optional, and absent isOptional.empty().connectionOverridesis a plainMap. - Nulls. Every record constructor rejects
null; passOptional.empty()for an absent field. Lists and maps are copied, and the copies reject null elements. - Empty strings. The constructors reject an empty naming template, table namespace, table name, table catalog name, and partition or sort column.
- Missing layer.
TableMaterializationPolicy.empty()returns the policy that resolution uses for a missing layer: an emptyconnectionOverridesand every other field absent. - Resolution.
TableMaterializationPolicy.resolveimplements Resolution. It returnsOptional.empty()for not materialized, and throwsIllegalArgumentExceptionfor a failure. Its only I/O is thecatalogLookupfunction it is given. The five-argument overload takes the stream properties. The four-argument overload passes none, so a${stream.property.<key>}variable fails, and the source logical name is the stream's name.TableNaming.toTableIdentifierhas the same two forms. - Stream catalog operations.
StreamCatalogsets and clears the namespace policy (createNamespacewith a policy,setNamespaceMaterialization,clearNamespaceMaterialization) and the stream policy (createStreamwith a policy,setStreamMaterialization,clearStreamMaterialization). It registers, unregisters, gets and lists table catalogs (registerTableCatalog,unregisterTableCatalog,getTableCatalog,listTableCatalogs).resolveMaterializationresolves a stream without opening data-plane resources, and returns aCompletableFuture<Optional<ResolvedMaterialization>>. See Catalog. - Cluster default.
setClusterDefaultMaterializationandclusterDefaultMaterializationhave default implementations: the setter throwsUnsupportedOperationExceptionsynchronously, and the getter returnsOptional.empty(). - Presets.
EvolutionPolicy.forIceberg()andEvolutionPolicy.forDelta()permitaddColumn,addNullableColumnandwidenType, and forbid the other six changes.EvolutionPolicy.forClickHouse()permits onlyaddColumnandaddNullableColumn. Each sets all nine flags. They are values to put in a policy, not defaults: resolution applies none. - Stream partition.
PartitionSpec.STREAM_PARTITION_COLUMNis__partition, andPartitionSpec.streamPartition()returns the partition field with column__partition, transformIDENTITYand no parameter. - No ownership mode.
lakestream-api1.0.0 removed theTableModeenum andTableConf.mode. - Blank. A blank string, in Terms, is one for which
String.isBlank()returnstrue; its whitespace test isCharacter.isWhitespace.
Appendix B: Implementation Notes
These notes describe Ursa 1.0. They are not part of the specification. Implementation status lists which policy fields each Ursa materializer applies.
- Cluster default. Ursa's stream catalog,
IndexedStreamCatalog, supports a cluster default policy. It uses the namespace policy as the baseline when the namespace has one, and the cluster default otherwise, and resolves the stream policy over that baseline. The cluster default replaces the namespace layer whole; it is never merged with a namespace policy field by field. It is held in the memory of the catalog instance that sets it, and is not persisted. Ursa's compactor sets it at every start, when its configuration asks for a cluster-wide default. - Fresh inputs. At every resolution, Ursa's stream catalog reads the stream's record, its namespace's record and the referenced table catalog from its metadata store. A change to a policy, to the stream's properties or to a table catalog therefore applies to the next resolution.
- Missing namespace record. If the stream's namespace has no record in the catalog, Ursa resolves the stream as if the namespace had no policy.
- Lookup before the disabled check. Ursa's stream catalog reads the table catalog that the effective catalog reference names before it runs the steps in Evaluation order. The outcome is the same, except that an error while reading the table catalog fails resolution even for a disabled stream.
- Fallback when not materialized. When resolution ends with not materialized, Ursa's compactor resolves the stream again from its own configuration and the compaction task's properties, a compatibility path for deployments that configure materialization there. If those configure no external destination, it uses a
NONEtable catalog namedinternal-compaction, which keeps internal compaction running. If they configure one, Ursa writes an external table for the stream; Implementation status lists this as a deviation from Using a resolved materialization. - Pinned table identifier. After it writes a task's output, Ursa stores the resolved table identifier and table catalog name on the compaction task. The asynchronous group commit reads them from the task, so it commits to the table the writer used rather than resolving again.
- ClickHouse names. ClickHouse has databases and tables, and no further namespace level. Ursa's ClickHouse materializer uses the table namespace as the database, and replaces each
/in the table name with., keeping the result as one quoted identifier. When Ursa builds a ClickHouse table catalog and policy from its flat lakehouse configuration and no template is configured, it sets the naming template${stream.namespace}.${stream.logicalName}, with the configured ClickHouse database (by default,default) as the prefix, so that streams with the same name in different namespaces get different tables. A stream in namespacesaleswith the logical nameordersthen resolves to the table identifier (default,sales.orders), and is written to the tablesales.orders, one identifier, in databasedefault. - Table catalogs from configuration. At every start, Ursa's compactor registers table catalogs from configuration keys of the forms
iceberg.catalog.<name>.<key>(typeICEBERG),delta.catalog.<name>.<key>(typeDELTA) andclickhouse.catalog.<name>.<key>(typeCLICKHOUSE). The<key>entries become the catalog'sconnection. Adelta.catalog.group becomesDELTA_UCwhen its name isunity, or when one of its keys starts withunity-catalog-orunity_catalog_or iscatalog-impl, ignoring case. The flatunityCatalog…keys, such asunityCatalogUri, become oneDELTA_UCcatalog namedunity, with keys such asunity-catalog-uri. - Serialized form. Ursa stores policies and table catalogs as JSON in its metadata store. A stored table catalog whose type is not a
TableCatalogTypevalue fails to load.