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



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

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.

<Mermaid
  chart="graph LR
  PUB[&#x22;Publisher&#x22;]
  STREAM{{&#x22;Events Store stream<br/>orders.ingest&#x22;}}
  STORE[(&#x22;Persistent store&#x22;)]

  PUB == &#x22;send event 1..N&#x22; ==> STREAM
  STREAM == &#x22;persist&#x22; ==> STORE
  STORE -. &#x22;ack per event&#x22; .-> PUB

  class PUB client
  class STREAM stream
  class STORE store"
/>

*One persistent stream carries many events to the store; each is acknowledged back to the publisher as it is persisted.*

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

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

## Step-by-Step [#step-by-step]

<Steps>
  <Step>
    ### Open a Stream and Publish Events [#open-a-stream-and-publish-events]

    Open a persistent stream connection and send order events at high throughput.

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

            "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")
        }
        ```
      </Tab>

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

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

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

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

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

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

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

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

  <Step>
    ### Handle Acknowledgments [#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.

    ```text
    Ack: ID=abc123 Sent=true
    Ack: ID=def456 Sent=true
    Error storing: storage has reached to 96.5% utilization and is not allowed
    ```
  </Step>
</Steps>

## Best Practices [#best-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    |

<Callout type="info">
  Stream publishing is available only over gRPC (port 50000). The REST API does not support streaming.
</Callout>

## Next Steps [#next-steps]

* Learn about [event sourcing](/learn/events-store/tutorials/event-sourcing) patterns
* Configure [retention policies](/learn/events-store/how-to/configure-retention) for high-volume streams
* Monitor [storage utilization](/learn/events-store/how-to/monitor-storage) thresholds
* See the [Events Store Reference](/learn/events-store/reference) for stream protocol details
