Lakestream
Ursa for Kafka

Deployment

Run Ursa for Kafka on Kubernetes with the Strimzi cluster operator and a separate Ursa compactor Deployment.

Ursa for Apache Kafka (UFK) ships a second image, lakestream/kafka-strimzi, that the Strimzi cluster operator can run in place of quay.io/strimzi/kafka. It is Strimzi's own Kafka image for the same Kafka minor with the bundled Apache Kafka distribution swapped for the UFK release tarball, so everything the operator relies on stays in place: kafka_run.sh and the readiness/liveness scripts, kafka-agent and tracing-agent, Strimzi's third-party libs, Cruise Control and the exporters, UID 1001 and KAFKA_HOME=/opt/kafka. On top of that it adds UFK's libs/, bin/ and config/ plus the isolated Ursa runtime under /opt/kafka/ursa-storage/, which the default ursa.storage.class.path resolves unchanged. Jars are merged by artifact name so nothing appears twice on the classpath.

Use lakestream/kafka-strimzi:latest unless you need to pin a release; versioned tags follow Strimzi's convention, lakestream/kafka-strimzi:<strimzi>-kafka-<kafka>, and every v* release tag publishes this image next to lakestream/kafka. To build it yourself instead, from a checkout of the 4.3-ursa branch:

git clone --branch 4.3-ursa https://github.com/openlakestream/kafka.git
cd kafka
docker/strimzi/build-image.sh lakestream/kafka-strimzi:latest   # ./gradlew releaseTarGz, then the image
docker/strimzi/verify-image.sh lakestream/kafka-strimzi:latest

STRIMZI_VERSION (default 1.2.0) and STRIMZI_KAFKA_VERSION (default 4.3.1) pick the base image; the base Kafka version must share its minor with the tarball. --tarball reuses an existing release tarball, --push --platform linux/amd64,linux/arm64 builds multi-arch and pushes. Without an image name the result is lakestream/kafka-strimzi:<STRIMZI_VERSION>-kafka-<tarball version>.

Not a Strimzi feature

Strimzi has no diskless awareness. It manages a UFK cluster the same way it manages any Kafka cluster; the ursa.* settings are ordinary broker config to it. The one place this shows is broker scale-down.

Prerequisites

  • A Kubernetes cluster with the Strimzi cluster operator installed (1.2.0 or later).
  • An Oxia service reachable from the brokers, for the Lakestream catalog and Ursa storage metadata.
  • An S3-compatible bucket (or GCS/Azure Blob) for the WAL and compacted objects.

Deploy the cluster

docker/strimzi/examples/kafka-diskless.yaml in the Kafka repository is a complete starting point: a Secret with the object-storage credentials, controller and broker KafkaNodePools, a Kafka resource with the Ursa broker settings, and a diskless KafkaTopic. The essential parts:

apiVersion: kafka.strimzi.io/v1
kind: Kafka
metadata:
  name: diskless
spec:
  kafka:
    version: 4.3.1                                        # a version Strimzi ships
    metadataVersion: 4.3-IV0
    image: lakestream/kafka-strimzi:latest                # ours
    config:
      config.providers: env
      config.providers.env.class: org.apache.kafka.common.config.provider.EnvVarConfigProvider
      ursa.storage.enable: true
      ursa.catalog.oxia.service.url: oxia://oxia.oxia.svc:6648/default
      ursa.oxia.service.url: oxia://oxia.oxia.svc:6648/default
      ursa.storage.backend.type: S3
      ursa.storage.path: ursa/wal                         # object prefix for remote backends
      ursa.storage.s3.endpoint: http://minio.minio.svc:9000
      ursa.storage.s3.access.key: ${env:URSA_S3_ACCESS_KEY}
      ursa.storage.s3.secret.key: ${env:URSA_S3_SECRET_KEY}
      ursa.storage.s3.bucket: kafka-ursa
      ursa.storage.s3.region: us-east-1
      ursa.storage.s3.path.style.access: true
      ursa.storage.compaction.bucket: kafka-ursa
      ursa.storage.compaction.prefix: ursa/compacted
    template:
      kafkaContainer:
        env:
          - name: URSA_S3_ACCESS_KEY
            valueFrom:
              secretKeyRef: { name: ursa-s3-credentials, key: accessKey }
          - name: URSA_S3_SECRET_KEY
            valueFrom:
              secretKeyRef: { name: ursa-s3-credentials, key: secretKey }
  entityOperator:
    topicOperator: {}
    userOperator: {}
  • spec.kafka.version must be a Kafka version the installed Strimzi release knows; the operator validates it against its compiled-in list (Strimzi 1.2.0 knows 4.2.0, 4.2.1, 4.3.0 and 4.3.1). Use the one matching the image's base, then spec.kafka.image decides which image runs.
  • Broker settings go in spec.kafka.config. The property names are the ones on the Configuration page. Use the env config provider and template.kafkaContainer.env to pull credentials from a Secret, as above, or drop the key settings and let the SDK's default credentials chain find them.
  • Brokers still need volumes for the KRaft metadata log and for classic topics such as __consumer_offsets; the example gives every node pool a persistent-claim volume with kraftMetadata: shared.

Apply the manifest, wait for the Kafka resource to reach Ready, and use <cluster>-kafka-bootstrap:9092 from clients inside the cluster.

Create diskless topics

A topic becomes diskless through the topic config ursa.storage.enable=true, whether it is created with a KafkaTopic resource or the Kafka admin API. Use replicas: 1; Ursa provides durability and Kafka-level replication is bypassed for diskless topics.

apiVersion: kafka.strimzi.io/v1
kind: KafkaTopic
metadata:
  name: diskless-events
  labels:
    strimzi.io/cluster: diskless
spec:
  partitions: 6
  replicas: 1
  config:
    ursa.storage.enable: "true"
    retention.ms: 604800000

Run the compactor

Compaction is what makes diskless storage reclaimable: retention on a diskless topic only hides records, and WAL objects are deleted behind a watermark that compaction alone advances. The compactor is not a broker, so Strimzi does not run it, and it is not in the Strimzi image. Run it as a plain Deployment of the lakestream/kafka:latest image, which bundles it at /opt/kafka/bin/ursa-compactor.sh, against the same Oxia service and bucket as the brokers. It never connects to Kafka. Apply this in the cluster's namespace; it reads the ursa-s3-credentials Secret the cluster manifest created:

apiVersion: v1
kind: ConfigMap
metadata:
  name: ursa-compactor
data:
  # Everything except credentials; those are appended from the Secret at start.
  ursa-storage.properties: |
    # Lakestream catalog, WAL metadata and compaction task coordination.
    metadataStoreUrl=oxia://oxia.oxia.svc:6648/default
    oxiaStorageUrl=oxia://oxia.oxia.svc:6648/default

    # Ursa WAL written by the brokers.
    backendStorageType=S3
    bucket=kafka-ursa
    prefix=ursa/wal
    cloudStorageEndpoint=http://minio.minio.svc:9000
    region=us-east-1
    s3PathStyleAccess=true

    # Compacted objects that Kafka fetch reads after compaction.
    compactionBackendStorageType=S3
    compactionBucket=kafka-ursa
    compactionPrefix=ursa/compacted
    compactionBucketRegion=us-east-1
    hadoop.fs.s3a.endpoint=http://minio.minio.svc:9000
    hadoop.fs.s3a.path.style.access=true
    hadoop.fs.s3a.connection.ssl.enabled=false

    # Diskless compaction-task ownership stays inside Ursa.
    internalCompactionTaskPublisherEnabled=true
    metastoreRequestRateLimitPerSecond=500

    # Storage-only compaction until lakehouseType and catalog settings turn on stream materialization.
    materializationEnabled=true
  log4j2.properties: |
    status = warn
    name = UrsaCompactor
    appender.console.type = Console
    appender.console.name = STDOUT
    appender.console.target = SYSTEM_OUT
    appender.console.layout.type = PatternLayout
    appender.console.layout.pattern = %d{ISO8601} %-5p [%t] %c{1} - %m%n
    rootLogger.level = info
    rootLogger.appenderRef.stdout.ref = STDOUT
---
apiVersion: apps/v1
kind: Deployment
metadata:
  name: ursa-compactor
  labels:
    app: ursa-compactor
spec:
  replicas: 1
  selector:
    matchLabels:
      app: ursa-compactor
  template:
    metadata:
      labels:
        app: ursa-compactor
    spec:
      containers:
        - name: compactor
          image: lakestream/kafka:latest
          command:
            - /bin/bash
            - -ec
            - |
              cp /etc/ursa/ursa-storage.properties /tmp/ursa-storage.properties
              printf 's3AccessKeyId=%s\ns3SecretAccessKey=%s\n' \
                "${AWS_ACCESS_KEY_ID}" "${AWS_SECRET_ACCESS_KEY}" >> /tmp/ursa-storage.properties
              exec /opt/kafka/bin/ursa-compactor.sh --conf /tmp/ursa-storage.properties
          env:
            - name: URSA_JAVA_OPTS
              value: "-Xmx2G -XX:+UseZGC -Dlog4j.configurationFile=/etc/ursa/log4j2.properties"
            - name: OTEL_SDK_DISABLED
              value: "true"
            - name: AWS_REGION
              value: us-east-1
            - name: AWS_EC2_METADATA_DISABLED
              value: "true"
            - name: AWS_ACCESS_KEY_ID
              valueFrom:
                secretKeyRef:
                  name: ursa-s3-credentials
                  key: accessKey
            - name: AWS_SECRET_ACCESS_KEY
              valueFrom:
                secretKeyRef:
                  name: ursa-s3-credentials
                  key: secretKey
          resources:
            requests:
              cpu: "1"
              memory: 3Gi
            limits:
              memory: 3Gi
          volumeMounts:
            - name: config
              mountPath: /etc/ursa
              readOnly: true
      volumes:
        - name: config
          configMap:
            name: ursa-compactor
  • Match metadataStoreUrl to the brokers' ursa.catalog.oxia.service.url and oxiaStorageUrl to ursa.oxia.service.url, and keep the WAL bucket/prefix and compaction bucket/prefix identical to the brokers' ursa.* settings. Otherwise brokers and compactor look at different data. The camelCase keys are documented under Ursa configuration.
  • Instances elect a leader through Oxia: the leader publishes tasks and runs the WAL cleaner, every instance runs workers. One replica is enough to start; more add throughput.
  • Keep materializationEnabled=true. With no lakehouseType or catalog settings the stream materialization framework registers no table catalog and compaction is storage-only: compacted objects under ursa/compacted, no table; Stream materialization below adds one. Do not set it to false: that selects a legacy compaction path that fails on Ursa 1.0.0 with Unsupported lakehouse type: NONE, quarantines every partition and never reclaims WAL objects.

Stream materialization

Stream materialization turns a diskless topic into an Iceberg table. The table is a second output of the same compaction pass, not a copy of the compacted objects Kafka reads; see the stream materialization framework and Ursa lakehouse tables. Turning it on is compactor configuration only — the brokers and the Strimzi resources do not change. Add to ursa-storage.properties in the ConfigMap above:

materializationDefaultNamespace=default
clusterSdtEnabled=true
lakehouseType=ICEBERG
catalog.name=polaris
iceberg.catalog.polaris.type=rest
iceberg.catalog.polaris.uri=http://polaris.lakehouse.svc:8181/api/catalog
iceberg.catalog.polaris.warehouse=ursa
iceberg.catalog.polaris.credential=${POLARIS_CLIENT_ID}:${POLARIS_CLIENT_SECRET}
iceberg.catalog.polaris.scope=PRINCIPAL_ROLE:ALL
iceberg.catalog.polaris.oauth2-server-uri=http://polaris.lakehouse.svc:8181/api/catalog/v1/oauth/tokens
iceberg.catalog.polaris.catalog-backend=POLARIS
iceberg.catalog.polaris.io-impl=org.apache.iceberg.aws.s3.S3FileIO
iceberg.catalog.polaris.client.region=us-east-1
iceberg.catalog.polaris.s3.endpoint=http://minio.minio.svc:9000
iceberg.catalog.polaris.s3.path-style-access=true
iceberg.catalog.polaris.s3.access-key-id=${AWS_ACCESS_KEY_ID}
iceberg.catalog.polaris.s3.secret-access-key=${AWS_SECRET_ACCESS_KEY}

These are the keys the Compose stack's lakehouse/compactor-entrypoint.sh writes for its Polaris catalog, minus clusterSbtEnabled and streamTableMode, which Ursa 1.0.0 no longer defines. catalog.name picks the iceberg.catalog.<name>.* group; warehouse is the catalog name inside the REST service, which must exist before the compactor starts (polaris-setup.sh shows the management-API call that creates it). Ship the credentials the same way the example ships the S3 keys — appended from a Secret at start, not in the ConfigMap. At startup the compactor logs default-policy bridge: namespace default → catalog polaris (ICEBERG), and the first compaction of each diskless topic creates default.<topic> in the catalog under <warehouse base location>/default/<topic>/.

The materialized table is named after the logical topic, so recreating a topic with the same name appends to the existing table. Running a REST catalog and its object store is outside this page; neither Strimzi nor UFK ships one.

Schema registry

Strimzi does not ship a schema registry, and diskless topics do not need one: brokers serve compacted objects without it, and storage-only compaction never reads a schema. Only stream materialization uses it. Without one, every table has a single binary payload column. To get typed columns, run a registry that serves the Confluent-compatible REST API (the Compose stack uses Karapace, storing schemas in a classic _schemas topic on the same brokers), register Avro, JSON Schema or Protobuf under <topic>-value before the first record reaches the compactor, and point the compactor at it:

schemaRegistryUrl=http://schema-registry.kafka.svc:8081

The compactor caches a missing subject and treats that topic as raw bytes until it restarts; a table created as payload is not retroactively converted. See the quickstart for the schema-to-table mapping.

Broker scale-down

Strimzi validates default.replication.factor, offsets.topic.replication.factor and transaction.state.log.replication.factor against the broker count first; lower them from the example's 3 before shrinking the pool below three brokers, or the Kafka resource goes NotReady with an InvalidResourceException and the diskless check never runs.

Before removing a broker, Strimzi refuses if the broker still holds partition replicas. Diskless partitions report their current owner as their single replica, so a broker that owns any diskless partition looks in use: the operator logs Cannot scale down brokers [n] because [n] have assigned partition-replicas, sets a ScaleDownPreventionCheck warning on the Kafka resource and reverts the node pool to its previous replica count.

To scale down anyway, annotate the Kafka resource with strimzi.io/skip-broker-scaledown-check: "true" and lower the pool's replicas again. Diskless partitions owned by the removed broker re-home to the remaining brokers on their own and their data stays readable and writable, because it lives in Ursa, not on the broker.

The annotation disables the check for classic topics too

Before scaling down with the annotation, make sure no classic partition, including __consumer_offsets, has its last replica on the brokers being removed.

Next steps