KubeMQ
LearnEventsTutorials

Stream Publishing

Send events at high throughput using bidirectional streaming for batched delivery.

What You Will Build

This tutorial uses Events — ephemeral, fire-and-forget pub/sub with no persistence and no replay. For the persistent, replayable version of stream publishing, see Events Store stream publishing.

A high-throughput event publisher that uses bidirectional streaming to send large volumes of events efficiently, with backpressure handling and error recovery.

A stream is a single long-lived bidirectional channel: the client opens it once, pushes many events through it, then closes it.

A persistent stream amortizes gRPC overhead across many events instead of paying it per call.

Prerequisites

Stream vs Single Send

AspectSingle SendStream Publishing
ConnectionNew RPC per eventPersistent bidirectional stream
ThroughputModerateHigh (batched I/O)
OverheadPer-call gRPC overheadAmortized over stream lifetime
BackpressureNone (fire-and-forget)Built-in flow control
Use caseLow-to-moderate volumeHigh-frequency telemetry, log shipping

Steps

Open an Event Stream

Create a persistent streaming connection and send events through it.

stream_publisher.go
package main

import (
    "context"
    "fmt"
    "log"
    "time"

    "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()

    streamCh := make(chan *kubemq.Event, 100)
    resultCh := make(chan *kubemq.EventSendResult, 100)

    go client.StreamEvents(ctx, streamCh, resultCh)

    for i := 0; i < 1000; i++ {
        body := fmt.Sprintf(`{"orderId":"ORD-%04d","timestamp":%d}`, i, time.Now().UnixMilli())
        streamCh <- kubemq.NewEvent().
            SetChannel("order-stream").
            SetBody([]byte(body))

        select {
        case result := <-resultCh:
            if !result.Sent {
                log.Printf("Failed to send event %d: %s", i, result.Error)
            }
        case <-time.After(5 * time.Second):
            log.Printf("Timeout waiting for result on event %d", i)
        }
    }

    close(streamCh)
    log.Println("Stream publishing complete: 1000 events sent")
}
stream_publisher.py
import time
from kubemq.pubsub import Client as PubSubClient
from kubemq.pubsub import EventMessage

client = PubSubClient(address="localhost:50000")

stream = client.open_events_stream()

for i in range(1000):
    body = f'{{"orderId":"ORD-{i:04d}","timestamp":{int(time.time() * 1000)}}}'
    result = stream.send(
        EventMessage(channel="order-stream", body=body.encode("utf-8"))
    )
    if result and not result.sent:
        print(f"Failed to send event {i}: {result.error}")

stream.close()
print("Stream publishing complete: 1000 events sent")
client.close()
stream_publisher.js
const { KubeMQClient, createEventMessage } = require("kubemq-js");

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

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

for (let i = 0; i < 1000; i++) {
  const body = JSON.stringify({
    orderId: `ORD-${String(i).padStart(4, "0")}`,
    timestamp: Date.now(),
  });

  stream.send(createEventMessage({
    channel: "order-stream",
    body: Buffer.from(body),
  }));
}

stream.close();
console.log("Stream publishing complete: 1000 events sent");
await client.close();
StreamPublisher.java
PubSubClient client = PubSubClient.builder()
    .address("localhost:50000")
    .clientId("stream-publisher")
    .build();

EventsStream stream = client.openEventsStream();

for (int i = 0; i < 1000; i++) {
    String body = String.format(
        "{\"orderId\":\"ORD-%04d\",\"timestamp\":%d}",
        i, System.currentTimeMillis());

    EventSendResult result = stream.send(EventMessage.builder()
        .channel("order-stream")
        .body(body.getBytes())
        .build());

    if (result != null && !result.isSent()) {
        System.err.printf("Failed to send event %d: %s%n", i, result.getError());
    }
}

stream.close();
System.out.println("Stream publishing complete: 1000 events sent");
client.close();
StreamPublisher.cs
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

var stream = client.OpenEventsStream();

for (var i = 0; i < 1000; i++)
{
    var body = $"{{\"orderId\":\"ORD-{i:D4}\",\"timestamp\":{DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()}}}";

    var result = await stream.SendAsync(new EventMessage
    {
        Channel = "order-stream",
        Body = Encoding.UTF8.GetBytes(body),
    });

    if (result is { Sent: false })
    {
        Console.Error.WriteLine($"Failed to send event {i}: {result.Error}");
    }
}

stream.Close();
Console.WriteLine("Stream publishing complete: 1000 events sent");
StreamPublisher.kt
val client = PubSubClient("localhost:50000")

val stream = client.openEventsStream()

for (i in 0 until 1000) {
    val body = """{"orderId":"ORD-${"%04d".format(i)}","timestamp":${System.currentTimeMillis()}}"""

    val result = stream.send(EventMessage(
        channel = "order-stream",
        body = body.toByteArray()
    ))

    if (result != null && !result.sent) {
        System.err.println("Failed to send event $i: ${result.error}")
    }
}

stream.close()
println("Stream publishing complete: 1000 events sent")
client.close()
stream_publisher.cpp
auto client = kubemq::PubSubClient("localhost:50000");

auto stream = client.openEventsStream();

for (int i = 0; i < 1000; i++) {
    kubemq::EventMessage event;
    event.channel = "order-stream";
    event.body = "{\"orderId\":\"ORD-" + std::to_string(i) +
                 "\",\"timestamp\":" +
                 std::to_string(std::chrono::system_clock::now()
                     .time_since_epoch().count()) + "}";

    auto result = stream.send(event);
    if (result && !result->sent) {
        std::cerr << "Failed to send event " << i
                  << ": " << result->error << std::endl;
    }
}

stream.close();
std::cout << "Stream publishing complete: 1000 events sent" << std::endl;
stream_publisher.rs
use kubemq::prelude::*;
use kubemq::EventBuilder;

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

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

    for i in 0..1000 {
        let body = format!(
            "{{\"orderId\":\"ORD-{:04}\",\"timestamp\":{}}}",
            i,
            chrono::Utc::now().timestamp_millis()
        );
        let event = EventBuilder::new()
            .channel("order-stream")
            .body(body.into_bytes())
            .build();

        stream.send(event).await?;
    }

    // Drain any stream errors reported asynchronously
    while let Ok(err) = stream.errors().try_recv() {
        eprintln!("Stream error: {}", err);
    }

    stream.close();
    client.close().await?;
    println!("Stream publishing complete: 1000 events sent");
    Ok(())
}
stream_publisher.rb
require "kubemq"

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

sender = client.create_events_sender

1000.times do |i|
  body = %({"orderId":"ORD-#{format('%04d', i)}","timestamp":#{(Time.now.to_f * 1000).to_i}})
  sender.publish(
    KubeMQ::PubSub::EventMessage.new(channel: "order-stream", body: body)
  )
end

sender.close
client.close
puts "Stream publishing complete: 1000 events sent"
stream_publisher.exs
{:ok, client} =
  KubeMQ.Client.start_link(address: "localhost:50000", client_id: "stream-publisher")

{:ok, handle} = KubeMQ.Client.send_event_stream(client)

Enum.each(0..999, fn i ->
  body =
    ~s({"orderId":"ORD-#{String.pad_leading(Integer.to_string(i), 4, "0")}",) <>
      ~s("timestamp":#{System.system_time(:millisecond)}})

  event = KubeMQ.Event.new(channel: "order-stream", body: body)
  KubeMQ.EventStreamHandle.send(handle, event)
end)

KubeMQ.Client.close(client)
IO.puts("Stream publishing complete: 1000 events sent")

Create a Subscriber

Subscribe to the stream channel to verify delivery.

stream_subscriber.go
counter := 0
sub, err := client.SubscribeToEvents(ctx, "order-stream", "",
    kubemq.WithOnEvent(func(event *kubemq.Event) {
        counter++
        if counter%100 == 0 {
            fmt.Printf("Received %d events (latest: %s)\n",
                counter, string(event.Body))
        }
    }),
    kubemq.WithOnError(func(err error) {
        log.Println("Error:", err)
    }),
)
stream_subscriber.py
counter = 0

def on_event(event):
    global counter
    counter += 1
    if counter % 100 == 0:
        print(f"Received {counter} events (latest: {event.body.decode()})")

client.subscribe_to_events(
    subscription=EventsSubscription(
        channel="order-stream",
        on_receive_event_callback=on_event,
        on_error_callback=lambda e: print(f"Error: {e}"),
    ),
    cancel=CancellationToken(),
)
stream_subscriber.js
let counter = 0;

client.subscribeToEvents({
  channel: "order-stream",
  onEvent: (msg) => {
    counter++;
    if (counter % 100 === 0) {
      console.log(
        `Received ${counter} events (latest: ${Buffer.from(msg.body).toString()})`
      );
    }
  },
  onError: (err) => console.error("Error:", err.message),
});
StreamSubscriber.java
AtomicInteger counter = new AtomicInteger(0);

client.subscribeToEvents(EventsSubscription.builder()
    .channel("order-stream")
    .onReceiveEventCallback(event -> {
        int count = counter.incrementAndGet();
        if (count % 100 == 0) {
            System.out.printf("Received %d events (latest: %s)%n",
                count, new String(event.getBody()));
        }
    })
    .onErrorCallback(err ->
        System.err.println("Error: " + err.getMessage()))
    .build());
StreamSubscriber.cs
var counter = 0;

await foreach (var msg in client.SubscribeToEventsAsync(
    new EventsSubscription { Channel = "order-stream" }))
{
    counter++;
    if (counter % 100 == 0)
    {
        Console.WriteLine($"Received {counter} events (latest: "
            + $"{Encoding.UTF8.GetString(msg.Body.Span)})");
    }
}
StreamSubscriber.kt
var counter = 0

client.subscribeToEvents(
    channel = "order-stream",
    onEvent = { event ->
        counter++
        if (counter % 100 == 0) {
            println("Received $counter events (latest: ${String(event.body)})")
        }
    },
    onError = { err -> System.err.println("Error: ${err.message}") }
)
stream_subscriber.cpp
int counter = 0;

client.subscribeToEvents("order-stream", "",
    [&counter](const kubemq::Event& event) {
        counter++;
        if (counter % 100 == 0) {
            std::cout << "Received " << counter << " events (latest: "
                      << event.body << ")" << std::endl;
        }
    },
    [](const std::string& err) {
        std::cerr << "Error: " << err << std::endl;
    }
);
stream_subscriber.rs
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;

let counter = Arc::new(AtomicUsize::new(0));
let counter_cb = counter.clone();

let sub = client
    .subscribe_to_events(
        "order-stream",
        "",
        move |event| {
            let counter = counter_cb.clone();
            Box::pin(async move {
                let n = counter.fetch_add(1, Ordering::SeqCst) + 1;
                if n % 100 == 0 {
                    println!(
                        "Received {} events (latest: {})",
                        n,
                        String::from_utf8_lossy(&event.body)
                    );
                }
            })
        },
        None,
    )
    .await?;
stream_subscriber.rb
counter = 0
cancel = KubeMQ::CancellationToken.new

sub = KubeMQ::PubSub::EventsSubscription.new(channel: "order-stream")
client.subscribe_to_events(
  sub,
  cancellation_token: cancel,
  on_error: ->(e) { warn "Error: #{e.message}" }
) do |event|
  counter += 1
  puts "Received #{counter} events (latest: #{event.body})" if (counter % 100).zero?
end
stream_subscriber.exs
{:ok, counter} = Agent.start_link(fn -> 0 end)

{:ok, _sub} =
  KubeMQ.Client.subscribe_to_events(client, "order-stream",
    on_event: fn event ->
      n = Agent.get_and_update(counter, fn c -> {c + 1, c + 1} end)
      if rem(n, 100) == 0 do
        IO.puts("Received #{n} events (latest: #{event.body})")
      end
    end,
    on_error: fn err -> IO.puts(:stderr, "Error: #{inspect(err)}") end
  )

Handle Backpressure

When the subscriber cannot keep up, the stream provides flow control. Monitor for send errors and implement retry or throttling.

Stream publishing does not change the at-most-once delivery guarantee. Events dropped due to slow consumers are still lost. For guaranteed delivery, use Events Store.

Best Practices

PracticeRecommendation
Buffer sizeSet channel buffer to match expected burst size (e.g., 100–1000)
Error handlingAlways check send results — log failures and consider retry
Stream lifetimeKeep streams open for the duration of high-throughput phases; close when done
Subscriber scalingPair with consumer groups for parallel processing

Next Steps

Was this page helpful?

On this page