KubeMQ
Client SDKsElixirHow-to guidesEvents Store

Stream Send

Efficiently batch-send persistent events to the KubeMQ Events Store using the Elixir SDK stream API.

Overview

Stream send covers publishing a batch of persistent events back-to-back, without pausing your producer between messages beyond what you choose to impose. A single send_event_store/2 call is fine for one-off writes, but if you're bulk-loading history, replicating a firehose of records, or backfilling an Events Store channel, sending events in rapid succession turns network latency into your throughput ceiling instead of an app-level bottleneck.

Each KubeMQ.Client.send_event_store/2 call still confirms storage synchronously, returning {:ok, %EventStoreResult{sent: true}} with the broker-assigned sequence number, or {:error, err} if persistence failed. Iterating the batch and pattern-matching on each result lets you confirm every event landed before moving on. Gotchas: each call blocks on its own confirmation, so this loop is latency-bound — very high fan-out workloads need concurrent sends rather than one tight loop; a single {:error, err} doesn't stop the batch, so you must check every result if you need all-or-nothing delivery; and don't rely on a fixed Process.sleep between sends in production — it's only there to make the demo's ordering legible.

Prerequisites

  • KubeMQ server running on localhost:50000
  • Elixir SDK installed ({:kubemq, "~> 1.0"} in mix.exs)

Code

main.exs
channel = "elixir-events-store.stream-send"
{:ok, client} = KubeMQ.Client.start_link(address: "localhost:50000", client_id: "elixir-es-stream")

IO.puts("Sending 5 events to Event Store...")

for i <- 1..5 do
  event = KubeMQ.EventStore.new(channel: channel, body: "Streamed store event #{i}")

  case KubeMQ.Client.send_event_store(client, event) do
    {:ok, result} ->
      IO.puts("Event #{i} confirmed: sent=#{result.sent}")

    {:error, err} ->
      IO.puts("Event #{i} failed: #{err.message}")
  end

  Process.sleep(100)
end

KubeMQ.Client.close(client)
IO.puts("All events persisted. Done.")

How It Works

  • Each send_event_store/2 call returns {:ok, %EventStoreResult{sent: true}} on success
  • Events are persisted sequentially with a small delay for demonstration
  • Each event gets a unique sequence number assigned by the server

Was this page helpful?

On this page