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
- KubeMQ server running on
localhost:50000 - SDK installed (Getting Started)
Stream vs Single Send
| Aspect | Single Send | Stream Publishing |
|---|---|---|
| Connection | New RPC per event | Persistent bidirectional stream |
| Throughput | Moderate | High (batched I/O) |
| Overhead | Per-call gRPC overhead | Amortized over stream lifetime |
| Backpressure | None (fire-and-forget) | Built-in flow control |
| Use case | Low-to-moderate volume | High-frequency telemetry, log shipping |
Steps
Open an Event Stream
Create a persistent streaming connection and send events through it.
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")
}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()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();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();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");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()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;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(())
}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"{: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.
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)
}),
)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(),
)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),
});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());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)})");
}
}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}") }
)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;
}
);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?;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{: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
| Practice | Recommendation |
|---|---|
| Buffer size | Set channel buffer to match expected burst size (e.g., 100–1000) |
| Error handling | Always check send results — log failures and consider retry |
| Stream lifetime | Keep streams open for the duration of high-throughput phases; close when done |
| Subscriber scaling | Pair with consumer groups for parallel processing |
Next Steps
Was this page helpful?