Lakestream
Ursa

Architecture

Ursa is organised into modules that separate its storage path, its catalog and its compaction lifecycle.

This page maps the Lakestream architecture onto Ursa's modules and runtime. The conceptual services are responsibilities, not necessarily separately deployed servers.

Module map

ModuleResponsibility
lakestream-apiPublic stream, log, catalog, provider, and materialization interfaces and records
ursa-storage-lakestreamIndexedStreamCatalog, LogImpl, LogCursorImpl, stream layouts, and StreamCatalogService bootstrap
ursa-storage-coreUrsaStorage, WAL storage, cloud/local file backends, caches, and offset indexes
ursa-storage-commonShared utilities and compaction contracts
ursa-storage-materializationMaterialization SPI, schema handling, and Kafka record codecs
ursa-storage-lakehouseLakehouse compaction, Iceberg/Delta writers, catalog integration, and storage bindings
ursa-storage-lakehouse-kafka-readerReconstructs Kafka entries from compacted data
ursa-storage-kafka-runtimeLeaf artifact assembling storage, catalog, compacted reader, and telemetry for Kafka
ursa-storage-clickhouseClickHouse materialization sink
ursa-storage-compactDistributed compaction scheduling, worker execution, commit, and cleanup orchestration
ursa-storage-containers, ursa-storage-test, ursa-storage-toolsTest infrastructure, integration tests, and command-line tools

The root POM no longer builds ursa-storage-ml, ursa-storage-pulsar, or ursa-storage-pulsar-ml. Kafka's current integration is not a ManagedLedger adapter.

Embedding and runtime assembly

Applications target lakestream-api. Ursa's StreamCatalogService constructs an IndexedStreamCatalog, storage, and Oxia clients from Java properties. Closing the returned catalog also closes resources whose ownership was transferred to it.

For Kafka, UrsaKafkaStreamCatalogProvider is discovered through ServiceLoader. It supplies a KafkaLakehouseReaderFactory, prepares storage properties, and transfers its OpenTelemetry SDK lifecycle to the catalog. The runtime depends on the implementation modules; those modules do not depend back on the Kafka runtime. Kafka keeps these implementation dependencies off its main classpath.

Kafka protocol and diskless adapter (kafka repository)
                 │
          lakestream-api
                 │  ServiceLoader
     ursa-storage-kafka-runtime
          ┌──────┴───────────────────┐
          ▼                          ▼
ursa-storage-lakestream     lakehouse-kafka-reader
          │                          │
          ▼                          ▼
ursa-storage-core          compacted-object reads
          │
     WAL objects + Oxia indexes
          │
          ▼
ursa-storage-compact + lakehouse/materialization
          ├── internal compacted objects → streaming reads
          └── optional external table/sink → analytical reads

This is a runtime/data-flow sketch, not a complete Maven dependency graph.

Write and read path

LogImpl.append(numberOfRecords, data) delegates to LogStorage, which composes WAL persistence with Oxia indexing. See appending for the normative contract; in Ursa's current object-WAL path:

  1. The append is admitted to a pending queue and its payload is copied into a write-cache segment.
  2. The segment is persisted through FileStorage as a RAW object, potentially containing multiple logs. FileStorage selects LOCAL, S3, GCS, or AZUREBLOB; local files are useful for development but do not provide shared, diskless durability.
  3. The per-log indexes are published to Oxia; sequence-key deltas assign record offsets and cumulative byte counts.
  4. Append callbacks receive the per-log indexing outcome. Successful appends are acknowledged after persistence and index publication, not merely after queue admission.

Different logs sharing an object have separate indexing outcomes. This is not an atomic multi-log batch, nor a transaction spanning object storage and Oxia.

How appends are batched

Three workers move an append through that sequence: one assembles pending appends into segments, one awaits persistence and triggers indexing, and one runs client callbacks.

The assembling worker flushes when either threshold is reached — writeBufferFlushSize bytes pending, or writeBufferFlushIntervalMs elapsed since the last flush. Under a light append rate the interval dominates, so it sets the floor on how long an append waits before it is written.

Pending appends are grouped by log and packed into a segment until it reaches writeBufferSize bytes or spans writeBufferMaxStreamIds logs. An append larger than a segment is persisted on its own.

Admission control is measured in bytes, not requests: once un-flushed appends reach maxPendingAddRequestsUsedBytes, further appends are rejected rather than queued. Rejection is how backpressure surfaces to a client.

See WAL settings for these properties.

Reads

A read resolves the offset through the index, and the entry's recorded file type decides where it is served from: RAW entries come from the WAL, PARQUET entries from a compacted object.

The WAL path tries three tiers in order:

  1. Write cache — segments flushed recently but still in memory. This is what lets a tailing read avoid the object store entirely.
  2. Read cache — WAL objects fetched back from the object store, bounded by readCacheMemorySize and evicted by size.
  3. The object store.

Cached segments are leased while a reader holds them, so eviction cannot close a segment out from under an in-flight read.

The compacted path needs a format-aware reader to reconstruct entries; Kafka supplies one through its runtime rather than making the core engine depend on Kafka protocol classes. Which reader is installed depends on how the catalog was opened — see embedding Ursa.

Two metadata connections

Ursa opens two connections to the metadata store, and they are configured separately because they carry different traffic:

SettingCarries
oxiaStorageUrlWAL offset indexes — the hot path for every append and read.
metadataStoreUrlCatalog metadata, locks, and compaction leader election.

They can address the same cluster, but the coordination connection is the one that needs notifications enabled, which is why the split exists rather than being a deployment preference.

A shared storage layer removes the need to replicate partition record data between brokers. It does not remove Oxia's own coordination or Kafka's controller quorum, ownership checks, and protocol-visible leader metadata. See Kafka architecture.

Compaction, materialization, and reclamation

The standalone CompactionMain runs CompactionScheduler. The lakehouse storage bindings discover logs, publish work, commit compacted ranges, and run cleanup; LakehouseMaterializationService performs the configured conversion and sink work. CompactionWorker resolves a stream's effective materialization and dispatches a task to it; see compaction for the normative rules that dispatch follows.

LakehouseMaterializationService can create, from one WAL read pass:

  • an internal compacted-object writer for protocol read-back, controlled by compactedObjectEnabled;
  • a catalog sink for a destination such as Iceberg, Delta, or ClickHouse — Ursa's external, stream-delivered table (SDT) path (see lakehouse tables);
  • storage-only output, the internal compacted-object writer alone, when TableCatalogType.NONE is selected.

Each task gets fresh, single-use materializers rather than a cached, shared one: a partitioned topic's partitions are compacted concurrently under one stream identity, so a shared materializer would be handed to multiple threads at once. Materializers write and commit their outputs before task completion is recorded. File-producing tasks are then persisted for an asynchronous group commit, which commits the files, advances the offload cursor, and removes the task; an inline sink with no file results (such as ClickHouse) retires the task directly after its own commit.

Internal compacted objects remain indexed for streaming reads and are not registered in any table catalog. Optional external materialization writes a separate table output from the same WAL read; registering an Iceberg catalog is not required for internal compaction. See the materialization SPI for the sink interface these outputs implement.

Compaction is required for WAL reclamation

Kafka retention issues a soft trim. It changes the first readable offset, but WAL deletion is gated by the oldest un-compacted position across streams. Without compaction, retention does not reclaim WAL bytes; a stalled stream can hold back shared reclamation. Do not run a persistent diskless cluster with only brokers, Oxia, and an object store.

Object reclamation is itself background work: the compacted-data cleaner scans eligible Parquet files before the mark-deleted boundary, resolves each compacted entry index to its file paths through the per-file CompactedObjectFileIndex, and deletes those files through FileStorage. It touches only internal compacted objects; a delivered table's files and snapshots belong to the destination system. See trim, deletion, and reclamation for the normative boundary the cleaner respects.

See compaction settings for scheduling controls, operations for running the service, and lakehouse tables for output ownership.