Lakestream
API

Catalog

StreamCatalog is the Lakestream API for creating and resolving streams, and for their metadata, namespaces, layouts and lifecycle state.

StreamCatalog (io.lakestream.api, extends AutoCloseable) separates metadata management from data-plane resources. Creating or loading a stream returns a StreamMetadata snapshot, never a reader, writer, or log handle. It has 34 methods — 31 abstract, 3 default — and implementations must be safe for concurrent use.

Types documented on this page (12): StreamCatalog, StreamIdentifier, Namespace, StreamMetadata, StreamCatalogEntry, LifecycleState, StreamConfig, SchemaConfig, Partitioning, PartitioningStrategy, SourceMetadataProperties, CatalogPaths.

Identity and initialization

String name();
CompletableFuture<Void> initialize(String name, Map<String, String> properties);

name() is synchronous and returns the catalog's own name. initialize sets up the catalog with a name and initialization properties; its Javadoc declares no exceptions.

Table catalog and namespace operations

CompletableFuture<Void> registerTableCatalog(TableCatalog catalog);
CompletableFuture<Boolean> unregisterTableCatalog(String name);
CompletableFuture<TableCatalog> getTableCatalog(String name);
CompletableFuture<List<TableCatalog>> listTableCatalogs();

CompletableFuture<Void> createNamespace(Namespace namespace);
CompletableFuture<List<Namespace>> listNamespaces();
CompletableFuture<Namespace> loadNamespaceMetadata(String namespaceName);
CompletableFuture<Boolean> dropNamespace(String namespaceName);
CompletableFuture<Boolean> namespaceExists(String namespaceName);
CompletableFuture<Void> setNamespaceProperties(String name, Map<String, String> props);
CompletableFuture<Void> removeNamespaceProperties(String name, List<String> keys);

registerTableCatalog fails with AlreadyExistsException if a catalog with the same name is already registered; Ursa 1.0 deviates here — see Implementation status. createNamespace fails with AlreadyExistsException if the namespace already exists. loadNamespaceMetadata fails with NoSuchNamespaceException if the namespace isn't found. dropNamespace fails with NamespaceNotEmptyException if the namespace still contains streams. setNamespaceProperties has merge semantics: it adds or overwrites the given keys without touching properties it doesn't mention.

Stream operations

CompletableFuture<List<StreamIdentifier>> listStreams(String namespaceName);
CompletableFuture<List<StreamCatalogEntry>> listStreamEntries(String namespaceName);

CompletableFuture<StreamMetadata> createStream(
    StreamIdentifier id, StreamConfig config, Partitioning partitioning,
    SchemaConfig schema, Map<String, String> properties);
CompletableFuture<StreamMetadata> createStream(
    StreamIdentifier id, StreamConfig config, Partitioning partitioning,
    SchemaConfig schema, Map<String, String> properties,
    Optional<TableMaterializationPolicy> materialization);

CompletableFuture<StreamMetadata> loadStream(StreamIdentifier identifier);
CompletableFuture<StreamMetadata> increasePartitions(StreamIdentifier identifier, int targetPartitionCount);
CompletableFuture<StreamMetadata> replaceStreamProperties(
    StreamIdentifier identifier, Map<String, String> properties, long sourceRevision);
CompletableFuture<Boolean> dropStream(StreamIdentifier identifier, boolean purge);
CompletableFuture<Boolean> streamExists(StreamIdentifier identifier);
  • createStream fails with AlreadyExistsException if the stream already exists, or StreamPermanentlyDeletedException if the identifier carries a permanent-deletion tombstone. Both overloads can also fail with UnsupportedOperationException if recovery finds retired keyed allocations that the catalog storage cannot durably fence.
  • listStreamEntries returns non-terminal lifecycle records only — streams being created, active, or being deleted. Completed deletion tombstones are never included.
  • loadStream fails with NoSuchStreamException if the stream isn't found.
  • increasePartitions is idempotent. Existing partitions stay available while new ones provision; the larger committed layout is published atomically only once every new partition has durable log and catalog metadata. It fails with NoSuchStreamException if the stream isn't active, or StreamPermanentlyDeletedException if the identifier is tombstoned.
  • replaceStreamProperties replaces the entire property snapshot. A sourceRevision older than or equal to the last applied revision is a no-op. It fails with NoSuchStreamException or StreamPermanentlyDeletedException on the same terms as increasePartitions.
  • dropStream always installs a permanent-deletion tombstone for the identifier, even when no live stream currently exists for it — so a false result means only that no live stream existed, not that catalog metadata was unchanged. The purge flag requests that data be purged as well as metadata removed when true; when false, only metadata is removed. A later createStream for the same identifier then fails with StreamPermanentlyDeletedException. dropStream itself can fail with UnsupportedOperationException if the catalog storage cannot durably fence keyed stream-ID mappings and writes to retired log IDs.

createStream's identity-tombstone check, dropStream's durable tombstone, and increasePartitions's atomic layout publication are this interface's view of rules the Storage Spec's Catalog semantics states normatively for the Catalog role, not for every implementation.

Materialization on the catalog

CompletableFuture<Void> setNamespaceMaterialization(String namespace, TableMaterializationPolicy policy);
CompletableFuture<Void> clearNamespaceMaterialization(String namespace);
CompletableFuture<Void> setStreamMaterialization(StreamIdentifier id, TableMaterializationPolicy policy);
CompletableFuture<Void> clearStreamMaterialization(StreamIdentifier id);
CompletableFuture<Optional<ResolvedMaterialization>> resolveMaterialization(StreamIdentifier identifier);

CompletableFuture<Void> setClusterDefaultMaterialization(TableMaterializationPolicy policy);   // default
Optional<TableMaterializationPolicy> clusterDefaultMaterialization();                          // default

setNamespaceMaterialization sets the namespace-level policy (the active baseline for every stream in the namespace); setStreamMaterialization sets a stream-level override. Both can be cleared independently of setting the other layer.

Data plane

CompletableFuture<StreamLayout> getLayout(StreamIdentifier identifier);
CompletableFuture<Log> openLog(StreamIdentifier identifier, LogId logId);
CompletableFuture<Log> openLog(StreamIdentifier id, int partitionIndex);   // default
CompletableFuture<StreamWriter> openWriter(StreamIdentifier identifier);
CompletableFuture<StreamReader> openReader(StreamIdentifier identifier);

openLog is a pure data-plane open: it must not allocate a log ID, register a partition, grow the stream, or mutate catalog metadata — see the Storage Spec's Catalog semantics for the normative form of this rule. The (identifier, logId) overload fails with NoSuchStreamException if the stream isn't active, or IllegalArgumentException if the log ID is not in the stream's committed layout. The (id, partitionIndex) default additionally fails with StreamPermanentlyDeletedException if the identifier is tombstoned. openReader's returned handle owns its lazily opened child logs and must be closed before the catalog is closed. See Logs and Streams for the handles these methods return.

Defaults

MethodDefault behavior
setClusterDefaultMaterialization(policy)Throws UnsupportedOperationException("cluster-default materialization is not supported by this catalog") directly — not through a failed future
clusterDefaultMaterialization()Returns Optional.empty(), synchronously
openLog(id, int partitionIndex)Calls loadStream(id), then metadata.layout().logIds(); if the index is out of range, returns a failed future carrying IllegalArgumentException; otherwise calls openLog(id, logIds.get(partitionIndex))

A StreamCatalog may hold a cluster-wide default materialization policy, set and read through these two methods. This API does not define how a cluster default combines with a namespace policy or a stream policy — TableMaterializationPolicy.resolve (see Materialization types) takes only a namespace policy and a stream policy as input; it has no cluster-default parameter at all. See Implementation status for how a specific implementation combines the three.

Not in the API

sealStream and truncateStream are not catalog operations. LifecycleState includes SEALED and TRUNCATING as enum members, but their existence in the enum is not itself a claim that the catalog interface exposes commands to reach them, or that every implementation produces them; see Implementation status.

Types

TypeKindShape
StreamIdentifierrecord(String namespace, String name); fullName() returns namespace + "/" + name; of(namespace, name) is a factory
Namespacerecord(String name, Map<String, String> properties, Optional<TableMaterializationPolicy> materialization); the materialization field is the namespace-level baseline that every stream in the namespace inherits
StreamMetadatarecord(StreamIdentifier identifier, StreamConfig config, Partitioning partitioning, SchemaConfig schema, Map<String, String> properties, Optional<TableMaterializationPolicy> materialization, LifecycleState state, StreamLayout layout, long metadataVersion). An immutable snapshot: it owns no data-plane resources and does not need closing. Loading a stream's metadata does not open a reader, writer, log, or cache
StreamCatalogEntryrecord(StreamIdentifier identifier, LifecycleState state, Map<String, String> properties, long metadataVersion) — the metadata-only view listStreamEntries returns
LifecycleStateenumCREATING, ACTIVE, SEALED, TRUNCATING, DELETING
StreamConfigrecord(Map<String, String> properties); the no-argument constructor gives an empty map
SchemaConfigrecord(String schemaType, Map<String, String> properties); the no-argument constructor gives ("NONE", {})
Partitioningrecord(PartitioningStrategy strategy, Map<String, String> config); numPartitions() reads the numPartitions key of config, defaulting to 1
PartitioningStrategyenumINDEXED (fixed number of logs accessed by integer index), RANGE (key-range segments with split/merge). Ursa 1.0 supports only INDEXED; creating a stream with RANGE fails — see Implementation status
SourceMetadataPropertiesfinal classHolds LOGICAL_NAME_PROPERTY = "lakestream.source.logical.name" and the static helper logicalName(streamId, properties)
CatalogPathsinterfaceAn implementation's strategy for constructing the catalog's metadata key paths (10 abstract methods, 3 default)

SourceMetadataProperties.logicalName(streamId, properties) returns the non-blank lakestream.source.logical.name property if present, otherwise the non-blank legacy lakestream.kafka.topic.name property, otherwise streamId.name(). The Javadoc is explicit that this is metadata, not the storage identity: a source integration must set it itself and prevent users from overriding it.

Namespace also has two convenience constructors beyond the three-argument canonical one: Namespace(name, properties), which omits the materialization policy, and Namespace(name), which omits both the properties and the policy.

CatalogPaths

CatalogPaths lets a catalog implementation vary the key layout it uses for stream, namespace, and table-catalog records without changing catalog behavior. It defines one constant, TOMBSTONE_SEGMENT = "_tombstones", and one rule: a stream namespace must never be named after a reserved segment, or its streams would collide with catalog records stored under that segment.

String streamMetadataPath(StreamIdentifier id);
String streamConfigPath(StreamIdentifier id);
String partitionMetadataPath(StreamIdentifier id, int partitionIndex);
String namespacePrefix(String namespace);
String partitionPrefix(StreamIdentifier id);
String namespacePath(String namespace);
String namespacesPrefix();
String streamConfigPrefix(String namespace);
String tableCatalogPath(String name);
String tableCatalogsPrefix();

String streamTombstonePrefix();                                  // default
String streamTombstonePath(StreamIdentifier id);                 // default
String compactedReaderName(StreamIdentifier id, int logIndex);   // default

namespacesPrefix() and tableCatalogsPrefix() take no arguments — each returns a prefix for scanning every namespace, or every registered table catalog, rather than a prefix scoped to one namespace. partitionMetadataPath takes both a StreamIdentifier and a partition index. The three default methods build on the abstract ones:

MethodDefault behavior
streamTombstonePrefix()Returns streamConfigPrefix(TOMBSTONE_SEGMENT)
streamTombstonePath(id)Returns streamTombstonePrefix() + id.namespace() + "/" + id.name()
compactedReaderName(id, logIndex)Returns id.fullName() + "-partition-" + logIndex

The exact key layout each of these methods produces is implementation-defined; see the Storage Spec's Appendix B: Oxia mapping for Ursa's.

Materialization resolution

resolveMaterialization(identifier) resolves a stream's effective materialization target without opening any data-plane resource. It returns a CompletableFuture<Optional<ResolvedMaterialization>>: the Optional is empty when the stream is not materialized (disabled, no catalog reference, or the referenced table catalog is not registered), otherwise it holds the resolved TableCatalog, TableIdentifier, and effective TableMaterializationPolicy.

The catalog side of materialization is this small set of methods:

  • registerTableCatalog / unregisterTableCatalog / getTableCatalog / listTableCatalogs manage the named TableCatalog registry that policies reference by name.
  • setNamespaceMaterialization / clearNamespaceMaterialization manage the namespace-level baseline policy.
  • setStreamMaterialization / clearStreamMaterialization manage the stream-level override.
  • setClusterDefaultMaterialization / clusterDefaultMaterialization manage an optional, cluster-wide default policy; this API does not define how it combines with the namespace and stream layers — see Implementation status.

The merge and fallback rules that turn a namespace policy, a stream policy, and a registered catalog into one resolved result are normative and live on the Materialization Spec, not here: see Materialization Spec: Resolution. The Java shape of the policy and result types is on Materialization types.