Lakestream
Ursa

Quickstart

Write to a stream and read it back through the Lakestream API, with Ursa embedded in a small Java program.

Ursa's public API is the Lakestream API. This quickstart embeds Ursa 1.0.0 from Maven Central in a Java program, creates a stream, appends one entry, and reads it back by offset. It uses the local-filesystem storage backend and one Oxia container, and no Kafka. Treat this as a way to see the API move, not a deployment recipe.

1. Prerequisites

  • JDK 17 or later
  • Maven 3.6.3 or later
  • Docker, for Oxia

2. Start Oxia

Oxia holds the stream catalog, the offset indexes, and cursor state. Run it standalone, the same image and version the Kafka Compose stack pins:

docker run -d --name oxia -p 6648:6648 oxia/oxia:0.16.7 /oxia/bin/oxia standalone

Standalone mode is a single-node development server. Its data lives inside the container.

3. Create the project

Ursa's artifacts are published under the org.openlakestream group; the Java packages stay io.lakestream.api and io.lakestream.ursa. A pom.xml for the example:

<project xmlns="http://maven.apache.org/POM/4.0.0">
  <modelVersion>4.0.0</modelVersion>
  <groupId>example</groupId>
  <artifactId>lakestream-quickstart</artifactId>
  <version>0.1</version>

  <dependencies>
    <dependency>
      <groupId>org.openlakestream</groupId>
      <artifactId>ursa-storage-lakestream</artifactId>
      <version>1.0.0</version>
    </dependency>
    <!-- ursa-storage-lakestream declares this as provided, so add it yourself -->
    <dependency>
      <groupId>org.openlakestream</groupId>
      <artifactId>ursa-storage-core</artifactId>
      <version>1.0.0</version>
      <scope>runtime</scope>
    </dependency>
  </dependencies>

  <build>
    <plugins>
      <plugin>
        <groupId>org.apache.maven.plugins</groupId>
        <artifactId>maven-compiler-plugin</artifactId>
        <version>3.13.0</version>
        <configuration><release>17</release></configuration>
      </plugin>
      <plugin>
        <groupId>org.codehaus.mojo</groupId>
        <artifactId>exec-maven-plugin</artifactId>
        <version>3.5.0</version>
        <configuration><mainClass>example.QuickStart</mainClass></configuration>
      </plugin>
    </plugins>
  </build>
</project>

ursa-storage-lakestream brings lakestream-api with it, but its published POM lists ursa-storage-core with provided scope. Without the second dependency the program compiles and then fails at startup with NoClassDefFoundError: io/lakestream/ursa/storage/impl/StorageConfig.

4. Write and read

Save this as src/main/java/example/QuickStart.java:

package example;

import io.lakestream.api.LogEntry;
import io.lakestream.api.Namespace;
import io.lakestream.api.Partitioning;
import io.lakestream.api.PartitioningStrategy;
import io.lakestream.api.RoutingKey;
import io.lakestream.api.SchemaConfig;
import io.lakestream.api.StreamCatalog;
import io.lakestream.api.StreamConfig;
import io.lakestream.api.StreamIdentifier;
import io.lakestream.api.StreamReader;
import io.lakestream.api.StreamWriter;
import io.lakestream.ursa.lakestream.impl.StreamCatalogService;
import io.netty.buffer.Unpooled;

import java.nio.charset.StandardCharsets;
import java.util.Map;
import java.util.Properties;

public class QuickStart {
    public static void main(String[] args) throws Exception {
        Properties props = new Properties();
        props.setProperty("backendStorageType", "LOCAL");
        props.setProperty("storagePath", "/tmp/lakestream-quickstart");
        props.setProperty("oxiaStorageUrl", "oxia://localhost:6648/default");

        StreamIdentifier id = StreamIdentifier.of("quickstart", "events");

        try (StreamCatalog catalog =
                 new StreamCatalogService().open("oxia://localhost:6648/default", props)) {

            if (!catalog.namespaceExists("quickstart").join()) {
                catalog.createNamespace(new Namespace("quickstart")).join();
            }
            if (!catalog.streamExists(id).join()) {
                catalog.createStream(
                    id,
                    new StreamConfig(),
                    new Partitioning(PartitioningStrategy.INDEXED, Map.of("numPartitions", "1")),
                    new SchemaConfig(),
                    Map.of()).join();
            }

            StreamWriter.WriteResult written;
            try (StreamWriter writer = catalog.openWriter(id).join()) {
                written = writer.write(
                    RoutingKey.roundRobin(), 1,
                    Unpooled.copiedBuffer("hello, lakestream", StandardCharsets.UTF_8)).join();
            }
            System.out.printf("wrote to log %d at offset %d%n", written.logId().id(), written.offset());

            try (StreamReader reader = catalog.openReader(id).join()) {
                StreamReader.ReadResult result =
                    reader.read(written.logId(), written.offset(), 10, 1024 * 1024).join();
                for (LogEntry entry : result.entries()) {
                    try (entry) {
                        System.out.printf("offset %d: %s%n",
                            entry.offset(), entry.payload().toString(StandardCharsets.UTF_8));
                    }
                }
                System.out.printf("next offset %d%n", result.nextOffset());
            }
        }
    }
}

What each part does:

  • Two Oxia addresses. The URI passed to open is the stream catalog's metadata store. oxiaStorageUrl is where the storage engine keeps offset indexes and log state. They can be the same Oxia namespace, as here, but both must be set; leaving out oxiaStorageUrl fails with Invalid metadata URL. Must start with 'oxia://'. In Kafka's configuration these are ursa.catalog.oxia.service.url and ursa.oxia.service.url.
  • backendStorageType=LOCAL writes WAL objects under storagePath on the local disk. It is a development backend; it does not give you shared, diskless durability. The object storage settings page lists the S3, GCS, and AZUREBLOB settings.
  • Namespace before stream. createStream requires an existing namespace, so the program creates quickstart when it is missing. Both checks make the program safe to rerun.
  • One partition, one log. PartitioningStrategy.INDEXED with numPartitions=1 gives the stream a single log. RoutingKey.roundRobin() lets the layout pick it. With more partitions, WriteResult.logId() tells you which log a write landed in.
  • Entries, not records. write takes the number of records the payload contains and a Netty ByteBuf. read addresses one log by record offset and returns ReadResult(entries, nextOffset). Every LogEntry must be closed, and payload() is a borrowed view that is only valid until then.
  • Close order. Writer and reader are closed before the catalog. StreamCatalogService is Ursa's embedding entry point; it returns an IndexedStreamCatalog, which implements StreamCatalog.

5. Run

mvn -q compile exec:java

Expected output:

wrote to log 2 at offset 0
offset 0: hello, lakestream
next offset 1

The log ID is allocated by the catalog, so yours may differ. Running the program again appends at offset 1 and reads that entry back. The WAL object and its .crc32c sidecar appear under /tmp/lakestream-quickstart/, in a date-based directory.

The process does not exit on its own

After catalog.close() returns, Ursa 1.0.0 leaves one non-daemon thread running, storage-simple-wal-callback-processor, an executor that ObjectWalStorageImpl.close() does not shut down. Stop the program with Ctrl+C, or end main with System.exit(0) in a throwaway program. A long-lived service is not affected.

What this did not do

  • No compaction. Nothing here runs the compactor, so the WAL is never rewritten into compacted objects and never reclaimed. See architecture for why a real deployment runs it.
  • No table. No materialization policy is set and no table catalog is registered, so nothing reaches Iceberg or Delta. Lakehouse tables covers that path.
  • No ServiceLoader. lakestream-api also has StreamCatalogLoader, which finds a StreamCatalogProvider through ServiceLoader. The only provider Ursa 1.0.0 registers is UrsaKafkaStreamCatalogProvider in ursa-storage-kafka-runtime, the Kafka-flavored wiring UFK loads. A plain embedding constructs StreamCatalogService directly, as above.

Cleanup

docker rm -f oxia
rm -rf /tmp/lakestream-quickstart

Where next