KubeMQ
LearnEvents StoreTutorials

Stream Publishing

Publish persistent events at high throughput using bidirectional streaming with acknowledgment.

This tutorial uses Events Store — persistent, replayable event storage with a per-event acknowledgment. For the ephemeral, fire-and-forget version of stream publishing, see Events stream publishing.

Stream publishing uses a bidirectional gRPC stream to send persistent events at high throughput. Unlike single-event publishing, the stream keeps a persistent connection open and returns an acknowledgment for every stored event, making it ideal for bulk ingestion scenarios.

One persistent stream carries many events to the store; each is acknowledged back to the publisher as it is persisted.

Stream vs Single Send

AspectSingle SendStream Send
ConnectionNew request per eventPersistent bidirectional stream
ThroughputModerateHigh (batch-friendly)
AcknowledgmentPer-call responseAsync ack per event on stream
Use caseOccasional publishesBulk ingestion, high-frequency data

Prerequisites

Step-by-Step

Open a Stream and Publish Events

Open a persistent stream connection and send order events at high throughput.

stream_publisher.go
package main

import (
    "context"
    "fmt"
    "log"

    "github.com/kubemq-io/kubemq-go/v2"
)

func main() {
    ctx := context.Background()
    client, err := kubemq.NewClient(ctx,
        kubemq.WithAddress("localhost", 50000),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer client.Close()

    stream, err := client.SendEventsStoreStream(ctx,
        kubemq.WithOnResult(func(result *kubemq.EventResult) {
            log.Printf("Ack: ID=%s Sent=%v", result.EventID, result.Sent)
            if result.Error != "" {
                log.Printf("Error: %s", result.Error)
            }
        }),
        kubemq.WithOnStreamError(func(err error) {
            log.Printf("Stream error: %v", err)
        }),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer stream.Close()

    for i := 1; i <= 1000; i++ {
        body := fmt.Sprintf(`{"orderId":"ORD-%05d","item":"widget","qty":%d}`, i, i%10+1)
        stream.Send(kubemq.NewEvent().
            SetChannel("orders.ingest").
            SetBody([]byte(body)),
        )
    }

    log.Println("Sent 1000 events via stream")
}
stream_publisher.py
import json
from kubemq import PubSubClient, EventStoreMessage

def on_result(result):
    if result.error:
        print(f"Error storing: {result.error}")
    else:
        print(f"Ack: ID={result.id} Sent={result.sent}")

with PubSubClient(address="localhost:50000") as client:
    stream = client.open_events_store_stream(
        on_result_callback=on_result,
        on_error_callback=lambda e: print(f"Stream error: {e}"),
    )

    for i in range(1, 1001):
        body = json.dumps({"orderId": f"ORD-{i:05d}", "item": "widget", "qty": i % 10 + 1})
        stream.send(EventStoreMessage(
            channel="orders.ingest",
            body=body.encode("utf-8"),
        ))

    print("Sent 1000 events via stream")
    stream.close()
stream_publisher.ts
import { KubeMQClient, createEventStoreMessage } from 'kubemq-js';

const client = await KubeMQClient.create({ address: 'localhost:50000' });

const stream = client.createEventStoreStream();
stream.onError((err) => console.error('Stream error:', err.message));

for (let i = 1; i <= 1000; i++) {
  // send() resolves once the server confirms persistence, rejects on failure
  await stream.send(
    createEventStoreMessage({
      channel: 'orders.ingest',
      body: JSON.stringify({ orderId: `ORD-${String(i).padStart(5, '0')}`, item: 'widget', qty: (i % 10) + 1 }),
    })
  );
}

console.log('Sent 1000 events via stream');
stream.close();
await client.close();
StreamPublisher.java
PubSubClient client = PubSubClient.builder()
    .address("localhost:50000")
    .clientId("stream-publisher")
    .build();

EventStoreStream stream = client.openEventsStoreStream(
    result -> {
        if (result.getError() != null && !result.getError().isEmpty()) {
            System.err.printf("Error storing: %s%n", result.getError());
        } else {
            System.out.printf("Ack: ID=%s%n", result.getId());
        }
    },
    err -> System.err.printf("Stream error: %s%n", err.getMessage())
);

for (int i = 1; i <= 1000; i++) {
    String body = String.format(
        "{\"orderId\":\"ORD-%05d\",\"item\":\"widget\",\"qty\":%d}", i, i % 10 + 1);
    stream.send(EventStoreMessage.builder()
        .channel("orders.ingest")
        .body(body.getBytes())
        .build());
}

System.out.println("Sent 1000 events via stream");
stream.close();
client.close();
StreamPublisher.cs
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

var stream = client.OpenEventsStoreStream(
    onResult: result =>
    {
        if (!string.IsNullOrEmpty(result.Error))
            Console.Error.WriteLine($"Error: {result.Error}");
        else
            Console.WriteLine($"Ack: ID={result.Id}");
    },
    onError: err => Console.Error.WriteLine($"Stream error: {err.Message}")
);

for (var i = 1; i <= 1000; i++)
{
    var body = $"{{\"orderId\":\"ORD-{i:D5}\",\"item\":\"widget\",\"qty\":{i % 10 + 1}}}";
    stream.Send(new EventStoreMessage
    {
        Channel = "orders.ingest",
        Body = Encoding.UTF8.GetBytes(body),
    });
}

Console.WriteLine("Sent 1000 events via stream");
stream.Close();
StreamPublisher.kt
val client = KubeMQClient.pubSub {
    address = "localhost:50000"
    clientId = "stream-publisher"
}

client.use {
    val stream = client.openEventsStoreStream(
        onResult = { result ->
            if (result.error.isNotEmpty()) {
                System.err.println("Error: ${result.error}")
            } else {
                println("Ack: ID=${result.id}")
            }
        },
        onError = { err -> System.err.println("Stream error: $err") }
    )

    for (i in 1..1000) {
        stream.send(eventStoreMessage {
            channel = "orders.ingest"
            body = """{"orderId":"ORD-${"%05d".format(i)}","item":"widget","qty":${i % 10 + 1}}""".toByteArray()
        })
    }

    println("Sent 1000 events via stream")
    stream.close()
}
stream_publisher.cc
kubemq::ClientOptions options;
options.set_address("localhost", 50000);
options.set_client_id("stream-publisher");
auto client = kubemq::Client::Create(options).value();

auto stream = client->OpenEventsStoreStream(
    [](const kubemq::EventResult& result) {
        if (!result.error().empty()) {
            std::cerr << "Error: " << result.error() << std::endl;
        } else {
            std::cout << "Ack: ID=" << result.id() << std::endl;
        }
    },
    [](const std::string& err) {
        std::cerr << "Stream error: " << err << std::endl;
    });

for (int i = 1; i <= 1000; ++i) {
    kubemq::EventStoreMessage msg;
    msg.set_channel("orders.ingest");
    msg.set_body("{\"orderId\":\"ORD-" + std::to_string(i) + "\",\"item\":\"widget\"}");
    stream->Send(msg);
}

std::cout << "Sent 1000 events via stream" << std::endl;
stream->Close();
stream_publisher.rs
use kubemq::prelude::*;
use kubemq::EventStoreBuilder;

#[tokio::main]
async fn main() -> kubemq::Result<()> {
    let client = KubemqClient::builder()
        .host("localhost")
        .port(50000)
        .build()
        .await?;

    let mut stream = client.send_event_store_stream().await?;

    for i in 1..=1000 {
        let event = EventStoreBuilder::new()
            .channel("orders.ingest")
            .body(format!(r#"{{"orderId":"ORD-{:05}","item":"widget","qty":{}}}"#, i, i % 10 + 1).into_bytes())
            .build();
        stream.send(event).await?;
    }

    println!("Sent 1000 events via stream");

    // Drain per-event results; `sent == false` indicates a failed store
    while let Ok(result) = stream.results().try_recv() {
        if !result.sent {
            eprintln!("Error: id={}, error={}", result.event_id, result.error);
        }
    }

    stream.close();
    client.close().await?;
    Ok(())
}
stream_publisher.rb
require 'kubemq'

client = KubeMQ::PubSubClient.new(address: 'localhost:50000', client_id: 'stream-publisher')
sender = client.create_events_store_sender

(1..1000).each do |i|
  msg = KubeMQ::PubSub::EventStoreMessage.new(
    channel: 'orders.ingest',
    body: %({"orderId":"ORD-#{format('%05d', i)}","item":"widget","qty":#{i % 10 + 1}})
  )
  result = sender.publish(msg)
  warn "Error: id=#{result.id}" unless result.sent
end

puts 'Sent 1000 events via stream'
sender.close
client.close
stream_publisher.exs
# The Elixir SDK does not expose a dedicated stream sender; send_event_store/2
# reuses the underlying connection and returns a confirmation per event.
{:ok, client} = KubeMQ.Client.start_link(address: "localhost:50000", client_id: "stream-publisher")

for i <- 1..1000 do
  event =
    KubeMQ.EventStore.new(
      channel: "orders.ingest",
      body: ~s({"orderId":"ORD-#{:io_lib.format("~5..0B", [i]) |> to_string()}","item":"widget","qty":#{rem(i, 10) + 1}})
    )

  case KubeMQ.Client.send_event_store(client, event) do
    {:ok, %{sent: false} = result} -> IO.puts(:stderr, "Error storing: #{result.error}")
    {:error, err} -> IO.puts(:stderr, "Send error: #{err.message}")
    _ok -> :ok
  end
end

IO.puts("Sent 1000 events via stream")
KubeMQ.Client.close(client)

Handle Acknowledgments

Each event sent through the stream receives an asynchronous acknowledgment. Monitor the onResult callback for confirmation or errors. Failed events can be retried.

Ack: ID=abc123 Sent=true
Ack: ID=def456 Sent=true
Error storing: storage has reached to 96.5% utilization and is not allowed

Best Practices

PracticeRecommendation
Batch sizeSend events as fast as the stream allows; backpressure is handled automatically
Error handlingLog failed acks and retry with exponential backoff
Stream lifecycleReuse a single stream for the lifetime of your publisher process
Channel separationUse separate channels for different event types to enable independent replay

Stream publishing is available only over gRPC (port 50000). The REST API does not support streaming.

Next Steps

Was this page helpful?

On this page