Stream Send
Publish KubeMQ Events at high throughput by streaming sends with the Go SDK.
Overview
Which to use
This page covers high-throughput Events (fire-and-forget) streaming via SendEventStream. For the persistent, replayable equivalent, see Events Store Stream Send.
Publishing events one at a time means each SendEvent call pays its own round-trip: write the request, wait on the connection, then move to the next event. That's fine for occasional notifications, but it caps throughput when you need to push hundreds or thousands of events per second — log forwarding, sensor telemetry, change-data-capture feeds — where per-call overhead dominates.
SendEventStream opens one bidirectional gRPC stream up front and returns an *EventStreamHandle. Each subsequent handle.Send(ev) writes a frame directly onto that already-open stream instead of negotiating a new call, so the sender loop isn't blocked waiting on a broker round-trip for every event.
Gotchas: because sends don't wait on a per-message round-trip, write failures surface asynchronously on the handle's errors channel — you must drain it, or failures go unnoticed. Events are still fire-and-forget pub/sub underneath: no subscriber means a streamed event is dropped just like a regular one. Always close the handle when you're done; a stream left open holds a gRPC connection on the broker.
Prerequisites
- KubeMQ server running on
localhost:50000 - Go SDK installed (
go get github.com/kubemq-io/kubemq-go/v2)
Code
// Example: events/stream-send
//
// Demonstrates high-throughput event publishing using SendEventStream.
// A bidirectional stream is opened for sending multiple events efficiently.
//
// Channel: go-events.stream-send
// Client ID: go-events-stream-send-client
//
// Run with a KubeMQ server on localhost:50000
// (see https://docs.kubemq.io/deploy).
package main
import (
"context"
"fmt"
"log"
"time"
"github.com/kubemq-io/kubemq-go/v2"
)
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
client, err := kubemq.NewClient(ctx,
kubemq.WithAddress("localhost", 50000),
kubemq.WithClientId("go-events-stream-send-client"),
)
if err != nil {
log.Fatal(err)
}
defer client.Close()
channel := "go-events.stream-send"
received := make(chan struct{})
// Subscribe to verify events arrive.
sub, err := client.SubscribeToEvents(ctx, channel, "",
kubemq.WithOnEvent(func(event *kubemq.Event) {
fmt.Printf("Stream received: channel=%s body=%s\n", event.Channel, event.Body)
select {
case received <- struct{}{}:
default:
}
}),
kubemq.WithOnError(func(err error) {
log.Println("Subscription error:", err)
}),
)
if err != nil {
log.Fatal(err)
}
defer sub.Unsubscribe()
// Open a stream for high-throughput publishing.
handle, err := client.SendEventStream(ctx)
if err != nil {
log.Fatalf("SendEventStream: %v", err)
}
defer handle.Close()
// Drain errors in background.
go func() {
for err := range handle.Errors {
log.Println("Stream send error:", err)
}
}()
// Send multiple events via the stream.
for i := range 5 {
ev := kubemq.NewEvent().
SetChannel(channel).
SetBody(fmt.Appendf(nil, "stream-msg-%d", i)).
SetMetadata("stream-demo")
if err := handle.Send(ev); err != nil {
log.Fatalf("handle.Send: %v", err)
}
fmt.Printf("Stream sent event %d\n", i+1)
}
// Wait for at least one event to be received.
select {
case <-received:
fmt.Println("Stream send demo complete")
case <-ctx.Done():
log.Fatal("Timed out waiting for stream event")
}
}
How It Works
client.SendEventStream(ctx)opens a bidirectional gRPC stream for events, returning a*EventStreamHandle.- A background goroutine drains
handle.Errorsso write errors are logged without blocking the sender loop. - The
for i := range 5loop callshandle.Send(ev)five times — each call writes a frame to the stream without waiting for broker acknowledgement, giving higher throughput than repeatedSendEventround-trips. - A subscriber is registered first to verify that events arrive;
receivedis a buffered channel so the callback never blocks the gRPC stream goroutine. defer handle.Close()flushes and closes the stream;defer sub.Unsubscribe()stops the subscriber.
Related
Was this page helpful?