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
| Aspect | Single Send | Stream Send |
|---|---|---|
| Connection | New request per event | Persistent bidirectional stream |
| Throughput | Moderate | High (batch-friendly) |
| Acknowledgment | Per-call response | Async ack per event on stream |
| Use case | Occasional publishes | Bulk ingestion, high-frequency data |
Prerequisites
- KubeMQ server running on
localhost:50000 - SDK installed (Getting Started)
Step-by-Step
Open a Stream and Publish Events
Open a persistent stream connection and send order events at high throughput.
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")
}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()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();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();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();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()
}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();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(())
}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# 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 allowedBest Practices
| Practice | Recommendation |
|---|---|
| Batch size | Send events as fast as the stream allows; backpressure is handled automatically |
| Error handling | Log failed acks and retry with exponential backoff |
| Stream lifecycle | Reuse a single stream for the lifetime of your publisher process |
| Channel separation | Use 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
- Learn about event sourcing patterns
- Configure retention policies for high-volume streams
- Monitor storage utilization thresholds
- See the Events Store Reference for stream protocol details
Was this page helpful?