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 standaloneStandalone 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
openis the stream catalog's metadata store.oxiaStorageUrlis 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 outoxiaStorageUrlfails withInvalid metadata URL. Must start with 'oxia://'. In Kafka's configuration these areursa.catalog.oxia.service.urlandursa.oxia.service.url. backendStorageType=LOCALwrites WAL objects understoragePathon the local disk. It is a development backend; it does not give you shared, diskless durability. The object storage settings page lists theS3,GCS, andAZUREBLOBsettings.- Namespace before stream.
createStreamrequires an existing namespace, so the program createsquickstartwhen it is missing. Both checks make the program safe to rerun. - One partition, one log.
PartitioningStrategy.INDEXEDwithnumPartitions=1gives 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.
writetakes the number of records the payload contains and a NettyByteBuf.readaddresses one log by record offset and returnsReadResult(entries, nextOffset). EveryLogEntrymust be closed, andpayload()is a borrowed view that is only valid until then. - Close order. Writer and reader are closed before the catalog.
StreamCatalogServiceis Ursa's embedding entry point; it returns anIndexedStreamCatalog, which implementsStreamCatalog.
5. Run
mvn -q compile exec:javaExpected output:
wrote to log 2 at offset 0
offset 0: hello, lakestream
next offset 1The 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-apialso hasStreamCatalogLoader, which finds aStreamCatalogProviderthroughServiceLoader. The only provider Ursa1.0.0registers isUrsaKafkaStreamCatalogProviderinursa-storage-kafka-runtime, the Kafka-flavored wiring UFK loads. A plain embedding constructsStreamCatalogServicedirectly, as above.
Cleanup
docker rm -f oxia
rm -rf /tmp/lakestream-quickstartWhere next
- Architecture — Ursa's module boundaries and storage components.
- Catalog and Logs — the contracts behind
createStream,write, andread. - Ursa for Apache Kafka (UFK) quickstart — the same storage behind a Kafka cluster, with the compactor and an Iceberg demo.
- Build from source — only needed to change Ursa itself.