# Stream Publishing (/learn/events/tutorials/stream-publishing)



## What You Will Build [#what-you-will-build]

<Callout type="info">
  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](/learn/events-store/tutorials/stream-publishing).
</Callout>

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.

<Mermaid
  chart="sequenceDiagram
  participant C as Publisher
  participant K as KubeMQ

  C->>K: open events stream
  loop while open
    C->>K: send event
    K-->>C: send result
  end
  C->>K: close stream"
/>

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

## Prerequisites [#prerequisites]

* KubeMQ server running on `localhost:50000`
* SDK installed ([Getting Started](/learn/events/getting-started))

## Stream vs Single Send [#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 [#steps]

<Steps>
  <Step>
    ### Open an Event Stream [#open-an-event-stream]

    Create a persistent streaming connection and send events through it.

    <Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
      <Tab value="Go">
        ```go title="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")
        }
        ```
      </Tab>

      <Tab value="Python">
        ```python title="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()
        ```
      </Tab>

      <Tab value="Node.js">
        ```javascript title="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();
        ```
      </Tab>

      <Tab value="Java">
        ```java title="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();
        ```
      </Tab>

      <Tab value="C#">
        ```csharp title="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");
        ```
      </Tab>

      <Tab value="Kotlin">
        ```kotlin title="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()
        ```
      </Tab>

      <Tab value="C++">
        ```cpp title="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;
        ```
      </Tab>

      <Tab value="Rust">
        ```rust title="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(())
        }
        ```
      </Tab>

      <Tab value="Ruby">
        ```ruby title="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"
        ```
      </Tab>

      <Tab value="Elixir">
        ```elixir title="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")
        ```
      </Tab>
    </Tabs>
  </Step>

  <Step>
    ### Create a Subscriber [#create-a-subscriber]

    Subscribe to the stream channel to verify delivery.

    <Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
      <Tab value="Go">
        ```go title="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)
            }),
        )
        ```
      </Tab>

      <Tab value="Python">
        ```python title="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(),
        )
        ```
      </Tab>

      <Tab value="Node.js">
        ```javascript title="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),
        });
        ```
      </Tab>

      <Tab value="Java">
        ```java title="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());
        ```
      </Tab>

      <Tab value="C#">
        ```csharp title="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)})");
            }
        }
        ```
      </Tab>

      <Tab value="Kotlin">
        ```kotlin title="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}") }
        )
        ```
      </Tab>

      <Tab value="C++">
        ```cpp title="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;
            }
        );
        ```
      </Tab>

      <Tab value="Rust">
        ```rust title="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?;
        ```
      </Tab>

      <Tab value="Ruby">
        ```ruby title="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
        ```
      </Tab>

      <Tab value="Elixir">
        ```elixir title="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
          )
        ```
      </Tab>
    </Tabs>
  </Step>

  <Step>
    ### Handle Backpressure [#handle-backpressure]

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

    <Callout type="warn">
      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](/learn/events-store).
    </Callout>
  </Step>
</Steps>

## Best Practices [#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](/learn/events/tutorials/consumer-groups) for parallel processing |

## Next Steps [#next-steps]

<Cards>
  <Card title="Consumer Groups" href="/learn/events/tutorials/consumer-groups" description="Distribute streamed events across workers." />

  <Card title="Handle Slow Consumers" href="/learn/events/how-to/handle-slow-consumers" description="Mitigate message drops under load." />
</Cards>
