Materialization types
The materialization policy model is expressed as Java records and enums in the io.lakestream.api.materialization package.
io.lakestream.api.materialization holds the Java types for the stream materialization framework's policy model: what streams become which tables, in what write mode, with what partitioning, evolution, and commit behavior. This package is not limited to data types — TableMaterializationPolicy.resolve implements the layer-merge algorithm itself, and TableNaming.toTableIdentifier implements template interpolation, both as plain, synchronous Java code in this module.
This page describes the Java shape: records, enums, constructor validation, and what resolve() actually does. What each field means, and which fields a conforming materializer must apply, is normative on the Lakestream Materialization Spec; this page links to it throughout rather than repeating it.
Types documented on this page (21): TableMaterializationPolicy, ResolvedMaterialization, TableCatalog, TableCatalogType, TableIdentifier, TableNaming, FrameworkConf, WriteMode, StartPosition, ErrorHandling, ErrorMode, CommitConfig, TableConf, PartitionSpec, PartitionTransform, SortColumn, SortDirection, RetentionConfig, Compression, EvolutionPolicy, MaterializationState.
Every record's canonical constructor rejects a null Optional field with a message telling the caller to use Optional.empty() instead, and defensively copies any list or map field with List.copyOf/Map.copyOf — which themselves reject a null element, key, or value.
TableMaterializationPolicy
record TableMaterializationPolicy(
Optional<String> catalogRef,
Optional<TableNaming> tableNaming,
Optional<TableIdentifier> tableIdentifier,
Optional<Boolean> enabled,
Optional<FrameworkConf> framework,
Optional<EvolutionPolicy> evolution,
Optional<List<String>> primaryKey,
Optional<Long> baseSchemaVersion,
Optional<TableConf> table,
Map<String, String> connectionOverrides) {
static TableMaterializationPolicy empty();
static Optional<ResolvedMaterialization> resolve(
Optional<TableMaterializationPolicy> namespacePolicy,
Optional<TableMaterializationPolicy> streamPolicy,
StreamIdentifier streamId,
Function<String, Optional<TableCatalog>> catalogLookup);
static Optional<ResolvedMaterialization> resolve(
Optional<TableMaterializationPolicy> namespacePolicy,
Optional<TableMaterializationPolicy> streamPolicy,
StreamIdentifier streamId,
Function<String, Optional<TableCatalog>> catalogLookup,
Map<String, String> properties);
}This one record is applied at two layers — a namespace policy (the baseline) and a stream policy (a sparse override) — and resolve merges the two into an effective policy. empty() returns an all-Optional.empty() policy with an empty connectionOverrides map: applying it as an override changes nothing.
| Field | Type | Holds |
|---|---|---|
catalogRef | Optional<String> | The name of a registered TableCatalog |
tableNaming | Optional<TableNaming> | A naming template — only read from the namespace layer (see below) |
tableIdentifier | Optional<TableIdentifier> | An explicit table identifier — only read from the stream layer (see below) |
enabled | Optional<Boolean> | Whether materialization is enabled |
framework | Optional<FrameworkConf> | Engine-agnostic configuration |
evolution | Optional<EvolutionPolicy> | Schema-evolution permissions |
primaryKey | Optional<List<String>> | Primary-key column names |
baseSchemaVersion | Optional<Long> | Base schema version for compatibility checks |
table | Optional<TableConf> | Engine-specific table configuration |
connectionOverrides | Map<String, String> | Per-stream overrides for the catalog's connection settings — always read from the stream layer only |
See Policies for what each field means to a materializer.
resolve() is a plain synchronous method
Both resolve overloads are static and return a plain Optional<ResolvedMaterialization> — not a CompletableFuture. This is different from StreamCatalog.resolveMaterialization(id) (see Catalog), which is asynchronous. All I/O — the one catalog lookup — is delegated through the catalogLookup callback the caller supplies; resolve itself performs none. The four-argument overload is equivalent to the five-argument one called with an empty properties map. Either overload can throw IllegalArgumentException synchronously if the namespace's TableNaming template fails to interpolate against the stream (see below) — since resolve is synchronous, this is a direct thrown exception, not a failed future.
resolve() merges layers; it does not substitute a default value for a field left absent in both. A field absent on the effective policy is a field that neither layer set, nothing more. The merge itself:
- If the stream policy's
enabledisOptional.of(false), resolution returnsOptional.empty()immediately. A namespaceenabledvalue never disables resolution by itself. - The effective
catalogRefis the stream's value if present, else the namespace's; if neither is present, or the referenced catalog isn't registered, resolution returnsOptional.empty(). - The effective
tableIdentifieris the stream's explicit value if present; otherwise the namespace'stableNamingapplied to the stream and its properties, if the namespace has one; otherwise the stream's source logical name (see Catalog:SourceMetadataProperties) paired with the stream's own namespace. framework,evolution, andtablemerge field by field, stream-over-namespace, recursing into their own nested records (commitinsideframework;retentioninsidetable).errorHandling(nested insideframework) is not merged field by field — the stream's wholeErrorHandlingvalue wins if present, else the namespace's whole value.primaryKey,partitionBy, andsortByare replaced as whole lists when the stream sets them, never concatenated with the namespace's list.connectionOverridescomes from the stream policy only: a namespace policy'sconnectionOverridesis a field of the same record and can be non-empty, butresolve()ignores it. The effective value is always the stream policy's map, or empty if the stream policy has none.tableIdentifieris read from the stream layer only: a namespace policy'stableIdentifieris likewise a real, settable field, butresolve()never reads it.tableNamingon the effective policy is always the namespace policy's own value, carried through for inspection. A stream policy can also carry atableNamingvalue — it is the same record type at both layers — butresolve()ignores it; only the namespace's value is ever consulted or carried onto the effective policy.
resolve() can also throw IllegalArgumentException synchronously when a derived table identifier component would be empty — for example, a namespace tableNaming with tableNamespacePrefix = Optional.of("") produces an empty effective namespace, which TableIdentifier's canonical constructor rejects.
The effective policy's enabled is picked the same way as other top-level fields — stream wins, namespace as fallback — so a non-empty result can still carry enabled set to false when the stream policy left it absent and the namespace policy set it false; only a stream policy's own explicit false short-circuits to not materialized. The Materialization Spec says a materializer must not read meaning into the effective policy's enabled field: a resolved materialization means the stream is materialized, regardless of what that field holds.
See Resolution for the full normative rule set this implements, including the case ordering between an empty result and a thrown exception.
TableNaming
record TableNaming(Optional<String> tableNamespacePrefix, String tableNameTemplate) {
TableIdentifier toTableIdentifier(StreamIdentifier streamId);
TableIdentifier toTableIdentifier(StreamIdentifier streamId, Map<String, String> properties);
}tableNameTemplate must be non-empty. Supported variables, case-sensitive, written ${name}:
| Variable | Value |
|---|---|
stream.namespace | The stream's namespace |
stream.name | The stream's storage name |
stream.logicalName | The source-owned logical name (see Catalog: SourceMetadataProperties) |
stream.property.<key> | The stream property named <key> — only resolved by the two-argument overload; the single-argument overload has no properties available and rejects a template using this variable |
Substitution is one regex pass over the template (Matcher.appendReplacement/quoteReplacement); a replaced value is not re-scanned for further variables. An unrecognized ${...} name throws IllegalArgumentException, as does a stream.property.<key> reference to a property that is absent or blank, or an interpolated result that is blank. There is no escape sequence for a literal ${ — text that never matches the variable pattern is passed through unchanged, but no syntax exists to keep ${ literal if it does happen to match one.
tableNamespacePrefix, when present, replaces the stream's namespace outright — it does not prepend to it. The effective table namespace is tableNamespacePrefix taken as a literal value (never interpolated) when present, otherwise streamId.namespace(). See Table naming for the full grammar and normative rules.
TableCatalog and TableCatalogType
record TableCatalog(String name, TableCatalogType type, Map<String, String> connection, Map<String, String> properties) {}
enum TableCatalogType { ICEBERG, DELTA, DELTA_UC, CLICKHOUSE, NONE }name must be non-null and non-empty. connection and properties are deliberately split: connection carries per-catalog connection settings (URI, warehouse, DSN, auth references, and similar); properties carries catalog-level tuning defaults. Both maps are defensively copied.
TableCatalogType is a closed enum — adding a member needs a lakestream-api release. Its five values:
| Value | Meaning |
|---|---|
ICEBERG | An Apache Iceberg catalog |
DELTA | A Delta Lake table stored on object storage, without a separate catalog service |
DELTA_UC | A Delta Lake table managed by a Unity Catalog |
CLICKHOUSE | A ClickHouse table store |
NONE | No external catalog — the stream compacts to storage only, with no external table sink |
See Table catalogs for the normative rules, and Implementation status for which types a given implementation actually writes to.
TableIdentifier and ResolvedMaterialization
record TableIdentifier(String namespace, String name) {}
record ResolvedMaterialization(
TableCatalog catalog,
TableIdentifier tableIdentifier,
TableMaterializationPolicy effectivePolicy) {}Both of TableIdentifier's components must be non-null and non-empty. ResolvedMaterialization is the result resolve() returns on success: the looked-up TableCatalog (not just its name — the caller doesn't need a second lookup), the resolved TableIdentifier, and the merged effectivePolicy. All three components are non-null.
Framework configuration
record FrameworkConf(
Optional<WriteMode> writeMode,
Optional<StartPosition> startPosition,
Optional<Boolean> paused,
Optional<ErrorHandling> errorHandling,
Optional<CommitConfig> commit) {}
enum WriteMode { APPEND, UPSERT, CDC }
enum StartPosition { EARLIEST, LATEST, OFFSET, TIMESTAMP }
record ErrorHandling(ErrorMode mode, Optional<String> dlqTopic) {}
enum ErrorMode { SUSPEND, SKIP, LOG }
record CommitConfig(
Optional<Integer> maxRetries,
Optional<Long> retryDelayMs,
Optional<Integer> batchSize) {}ErrorHandling.mode is required whenever an ErrorHandling value is present — only dlqTopic is optional. StartPosition.OFFSET and StartPosition.TIMESTAMP name a starting point but carry no offset or timestamp value themselves; nothing in this record or its neighbors supplies one. See Framework configuration for what each field means and which write mode requires which upstream operation codes.
Table configuration
record TableConf(
Optional<List<PartitionSpec>> partitionBy,
Optional<List<SortColumn>> sortBy,
Optional<RetentionConfig> retention,
Optional<Long> targetFileSizeBytes,
Optional<Compression> compression) {}
record PartitionSpec(String column, PartitionTransform transform, Optional<String> parameter) {
static final String STREAM_PARTITION_COLUMN = "__partition";
static PartitionSpec streamPartition();
}
record SortColumn(String column, SortDirection direction, boolean nullsFirst) {}
enum SortDirection { ASC, DESC }
record RetentionConfig(
Optional<Long> snapshotRetentionMs,
Optional<Integer> maxSnapshots,
Optional<Long> rowRetentionMs) {}
enum Compression { ZSTD, SNAPPY, GZIP, LZ4, UNCOMPRESSED }TableConf has no table-mode or ownership field — Ursa 1.0 removed TableMode from the framework entirely; this record was never given a replacement for it.
PartitionSpec.transform selects a PartitionTransform:
| Value | Parameter |
|---|---|
IDENTITY | None — partition by the raw column value |
BUCKET | Required: the bucket count |
TRUNCATE | Required: the truncation width |
YEAR, MONTH, DAY | None — a timestamp or date column |
HOUR | None — a timestamp column |
EXPRESSION | Required, engine-specific — not portable across destinations |
PartitionSpec performs no cross-field validation itself: nothing in the record enforces that BUCKET/TRUNCATE/EXPRESSION actually carry a parameter, or rejects one that isn't a valid engine expression — that is left to the materializer. PartitionSpec.streamPartition() returns the sentinel spec (STREAM_PARTITION_COLUMN, IDENTITY, empty), where STREAM_PARTITION_COLUMN is the literal "__partition", denoting the stream's own intrinsic partition rather than a column from its records.
See Table configuration for the normative requirements these types carry, including which transforms a conforming materializer must reject rather than silently ignore.
EvolutionPolicy
record EvolutionPolicy(
Optional<Boolean> addColumn,
Optional<Boolean> addNullableColumn,
Optional<Boolean> dropColumn,
Optional<Boolean> widenType,
Optional<Boolean> narrowType,
Optional<Boolean> renameColumn,
Optional<Boolean> reorderColumns,
Optional<Boolean> nullabilityRelax,
Optional<Boolean> nullabilityTighten) {
static EvolutionPolicy forIceberg();
static EvolutionPolicy forDelta();
static EvolutionPolicy forClickHouse();
}Nine independent Optional<Boolean> flags, each stating whether one kind of schema change is permitted. forIceberg() and forDelta() return identical, permissive presets (addColumn, addNullableColumn, and widenType set true; every other flag set false); forClickHouse() returns a stricter preset (only the two "add" flags true). These factories are ordinary static constructors a caller can use to build a namespace-layer policy — they are not defaults that resolve() applies automatically, and calling none of them leaves every flag Optional.empty() on the resulting policy. See Schema evolution for what each flag means and how a materializer is expected to enforce it.
MaterializationState
enum MaterializationState { PENDING, RUNNING, DEGRADED, SUSPENDED, PAUSED }No method anywhere in this API accepts or returns MaterializationState — there is no way to query it through this interface. See Materialization state for what an implementation that does report state is expected to mean by each value.