KubeMQ

Getting Started with Events Store

Publish persistent events and subscribe with replay in 5 minutes.

This guide walks you through publishing persistent events and subscribing with replay. By the end, you will see how Events Store preserves messages for subscribers that connect after the event was published.

This is the quickstart — a single publisher and subscriber in 5 minutes. For the deep-dive covering multiple subscriber types (audit replay vs. real-time-only) and durable subscriptions, see Persistent Publish & Subscribe.

Prerequisites

  • KubeMQ server running on localhost:50000
  • One of the supported SDKs installed

Need to install KubeMQ? Run it with Docker in seconds:

docker run -d \  --name kubemq \  -p 50000:50000 \  -p 9090:9090 \  -p 8080:8080 \  -e KUBEMQ_TOKEN=YOUR_LICENSE_KEY \  europe-docker.pkg.dev/kubemq/images/kubemq:next

What You Will Build

An order tracking system where:

  • A publisher stores order events to a persistent channel
  • A subscriber connects later and replays the full order history using StartFromFirst

Events are stored on publish; a later subscriber replays the full history, then continues live.

Step-by-Step Guide

Install the SDK

go get github.com/kubemq-io/kubemq-go/v2
pip install kubemq
npm install kubemq-js
<dependency>
    <groupId>io.kubemq.sdk</groupId>
    <artifactId>kubemq-sdk-Java</artifactId>
    <version>2.1.1</version>
</dependency>
dotnet add package KubeMQ.SDK.CSharp
implementation("io.kubemq.sdk:kubemq-sdk-kotlin:2.1.0")
vcpkg install kubemq
cargo add kubemq
gem install kubemq
# mix.exs
def deps do
  [{:kubemq, "~> 1.0"}]
end

Publish Persistent Events

Publish several order events before starting the subscriber. With Events Store, messages are persisted and available for replay.

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()

    orders := []string{
        `{"action":"order.created","orderId":"ORD-1001","total":149.99}`,
        `{"action":"order.paid","orderId":"ORD-1001","method":"credit_card"}`,
        `{"action":"order.shipped","orderId":"ORD-1001","carrier":"fedex"}`,
    }

    for i, body := range orders {
        result, err := client.SendEventStore(ctx, kubemq.NewEvent().
            SetChannel("orders.events").
            SetBody([]byte(body)),
        )
        if err != nil {
            log.Printf("Failed to store event %d: %v", i+1, err)
            continue
        }
        log.Printf("Stored event %d: ID=%s", i+1, result.EventID)
    }
}
publisher.py
from kubemq import PubSubClient, EventStoreMessage

orders = [
    b'{"action":"order.created","orderId":"ORD-1001","total":149.99}',
    b'{"action":"order.paid","orderId":"ORD-1001","method":"credit_card"}',
    b'{"action":"order.shipped","orderId":"ORD-1001","carrier":"fedex"}',
]

with PubSubClient(address="localhost:50000") as client:
    for i, body in enumerate(orders, 1):
        result = client.publish_event_store(
            EventStoreMessage(channel="orders.events", body=body)
        )
        print(f"Stored event {i}: ID={result.id}")
publisher.ts
import { KubeMQClient, createEventStoreMessage } from 'kubemq-js';

const client = await KubeMQClient.create({ address: 'localhost:50000' });

const orders = [
  '{"action":"order.created","orderId":"ORD-1001","total":149.99}',
  '{"action":"order.paid","orderId":"ORD-1001","method":"credit_card"}',
  '{"action":"order.shipped","orderId":"ORD-1001","carrier":"fedex"}',
];

for (const [i, body] of orders.entries()) {
  const result = await client.sendEventStore(
    createEventStoreMessage({ channel: 'orders.events', body })
  );
  console.log(`Stored event ${i + 1}: ID=${result.id}`);
}
Publisher.java
PubSubClient client = PubSubClient.builder()
    .address("localhost:50000")
    .clientId("order-publisher")
    .build();

String[] orders = {
    "{\"action\":\"order.created\",\"orderId\":\"ORD-1001\",\"total\":149.99}",
    "{\"action\":\"order.paid\",\"orderId\":\"ORD-1001\",\"method\":\"credit_card\"}",
    "{\"action\":\"order.shipped\",\"orderId\":\"ORD-1001\",\"carrier\":\"fedex\"}"
};

for (int i = 0; i < orders.length; i++) {
    EventSendResult result = client.sendEventsStoreMessage(
        EventStoreMessage.builder()
            .channel("orders.events")
            .body(orders[i].getBytes())
            .build());
    System.out.printf("Stored event %d: ID=%s%n", i + 1, result.getId());
}
client.close();
Publisher.cs
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

string[] orders = {
    "{\"action\":\"order.created\",\"orderId\":\"ORD-1001\",\"total\":149.99}",
    "{\"action\":\"order.paid\",\"orderId\":\"ORD-1001\",\"method\":\"credit_card\"}",
    "{\"action\":\"order.shipped\",\"orderId\":\"ORD-1001\",\"carrier\":\"fedex\"}"
};

for (var i = 0; i < orders.Length; i++)
{
    var result = await client.SendEventStoreAsync(new EventStoreMessage
    {
        Channel = "orders.events",
        Body = Encoding.UTF8.GetBytes(orders[i]),
    });
    Console.WriteLine($"Stored event {i + 1}: ID={result.Id}");
}
Publisher.kt
val client = KubeMQClient.pubSub {
    address = "localhost:50000"
    clientId = "order-publisher"
}

val orders = listOf(
    """{"action":"order.created","orderId":"ORD-1001","total":149.99}""",
    """{"action":"order.paid","orderId":"ORD-1001","method":"credit_card"}""",
    """{"action":"order.shipped","orderId":"ORD-1001","carrier":"fedex"}""",
)

client.use {
    orders.forEachIndexed { i, body ->
        val result = client.sendEventStore(eventStoreMessage {
            channel = "orders.events"
            this.body = body.toByteArray()
        })
        println("Stored event ${i + 1}: ID=${result.id}")
    }
}
publisher.cc
kubemq::ClientOptions options;
options.set_address("localhost", 50000);
options.set_client_id("order-publisher");

auto client = kubemq::Client::Create(options).value();

std::vector<std::string> orders = {
    R"({"action":"order.created","orderId":"ORD-1001","total":149.99})",
    R"({"action":"order.paid","orderId":"ORD-1001","method":"credit_card"})",
    R"({"action":"order.shipped","orderId":"ORD-1001","carrier":"fedex"})"
};

for (size_t i = 0; i < orders.size(); ++i) {
    kubemq::EventStoreMessage msg;
    msg.set_channel("orders.events");
    msg.set_body(orders[i]);
    auto result = client->SendEventStore(msg);
    if (result.ok()) {
        std::cout << "Stored event " << i + 1 << ": ID=" << result->id() << std::endl;
    }
}
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 orders = [
        r#"{"action":"order.created","orderId":"ORD-1001","total":149.99}"#,
        r#"{"action":"order.paid","orderId":"ORD-1001","method":"credit_card"}"#,
        r#"{"action":"order.shipped","orderId":"ORD-1001","carrier":"fedex"}"#,
    ];

    for (i, body) in orders.iter().enumerate() {
        let event = EventStoreBuilder::new()
            .channel("orders.events")
            .body(body.as_bytes().to_vec())
            .build();
        let result = client.send_event_store(event).await?;
        println!("Stored event {}: id={}", i + 1, result.id);
    }

    client.close().await?;
    Ok(())
}
publisher.rb
require 'kubemq'

client = KubeMQ::PubSubClient.new(address: 'localhost:50000', client_id: 'order-publisher')

orders = [
  '{"action":"order.created","orderId":"ORD-1001","total":149.99}',
  '{"action":"order.paid","orderId":"ORD-1001","method":"credit_card"}',
  '{"action":"order.shipped","orderId":"ORD-1001","carrier":"fedex"}'
]

orders.each_with_index do |body, i|
  msg = KubeMQ::PubSub::EventStoreMessage.new(channel: 'orders.events', body: body)
  result = client.send_event_store(msg)
  puts "Stored event #{i + 1}: sent=#{result.sent}"
end

client.close
publisher.exs
{:ok, client} =
  KubeMQ.Client.start_link(address: "localhost:50000", client_id: "order-publisher")

orders = [
  ~s({"action":"order.created","orderId":"ORD-1001","total":149.99}),
  ~s({"action":"order.paid","orderId":"ORD-1001","method":"credit_card"}),
  ~s({"action":"order.shipped","orderId":"ORD-1001","carrier":"fedex"})
]

orders
|> Enum.with_index(1)
|> Enum.each(fn {body, i} ->
  event = KubeMQ.EventStore.new(channel: "orders.events", body: body)
  {:ok, result} = KubeMQ.Client.send_event_store(client, event)
  IO.puts("Stored event #{i}: sent=#{result.sent}")
end)

KubeMQ.Client.close(client)

Subscribe with Replay from Beginning

Start the subscriber after all events have been published. Using StartFromFirst, it replays the entire history.

subscriber.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()

    sub, err := client.SubscribeToEventsStore(ctx, "orders.events", "",
        kubemq.StartFromFirst(),
        kubemq.WithOnEvent(func(event *kubemq.Event) {
            fmt.Printf("Received: seq=%d body=%s\n",
                event.Sequence, string(event.Body))
        }),
        kubemq.WithOnError(func(err error) {
            log.Println("Error:", err)
        }),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer sub.Unsubscribe()

    log.Println("Replaying order history...")
    <-ctx.Done()
}
subscriber.py
import time
from kubemq import (
    PubSubClient, EventsStoreSubscription,
    EventStoreStartPosition, CancellationToken,
)

def on_event(event):
    print(f"Received: seq={event.sequence} body={event.body.decode('utf-8')}")

with PubSubClient(address="localhost:50000") as client:
    client.subscribe_to_events_store(
        subscription=EventsStoreSubscription(
            channel="orders.events",
            start_position=EventStoreStartPosition.StartFromFirst,
            on_receive_event_callback=on_event,
            on_error_callback=lambda e: print(f"Error: {e}"),
        ),
        cancel=CancellationToken(),
    )
    print("Replaying order history...")
    time.sleep(120)
subscriber.ts
import { KubeMQClient, EventStoreStartPosition } from 'kubemq-js';

const client = await KubeMQClient.create({ address: 'localhost:50000' });

client.subscribeToEventsStore({
  channel: 'orders.events',
  startPosition: EventStoreStartPosition.StartFromFirst,
  onEvent: (msg) =>
    console.log(
      `Received: seq=${msg.sequence} body=${new TextDecoder().decode(msg.body)}`
    ),
  onError: (err) => console.error('Error:', err.message),
});

console.log('Replaying order history...');
Subscriber.java
PubSubClient client = PubSubClient.builder()
    .address("localhost:50000")
    .clientId("order-subscriber")
    .build();

client.subscribeToEventsStore(EventsStoreSubscription.builder()
    .channel("orders.events")
    .startPosition(EventStoreStartPosition.StartFromFirst)
    .onReceiveEventCallback(event ->
        System.out.printf("Received: seq=%d body=%s%n",
            event.getSequence(), new String(event.getBody())))
    .onErrorCallback(err ->
        System.err.println("Error: " + err.getMessage()))
    .build());

System.out.println("Replaying order history...");
Thread.sleep(120_000);
client.close();
Subscriber.cs
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

Console.WriteLine("Replaying order history...");
await foreach (var msg in client.SubscribeToEventsStoreAsync(
    new EventsStoreSubscription
    {
        Channel = "orders.events",
        StartPosition = EventStoreStartPosition.StartFromFirst,
    }))
{
    Console.WriteLine($"Received: seq={msg.Sequence} "
        + $"body={Encoding.UTF8.GetString(msg.Body.Span)}");
}
Subscriber.kt
val client = KubeMQClient.pubSub {
    address = "localhost:50000"
    clientId = "order-subscriber"
}

client.use {
    val flow = client.subscribeToEventsStore {
        channel = "orders.events"
        startPosition = StartPosition.StartFromFirst
    }
    println("Replaying order history...")
    flow.collect { msg ->
        println("Received: seq=${msg.sequence} body=${String(msg.body)}")
    }
}
subscriber.cc
kubemq::ClientOptions options;
options.set_address("localhost", 50000);
options.set_client_id("order-subscriber");

auto client = kubemq::Client::Create(options).value();

std::cout << "Replaying order history..." << std::endl;
client->SubscribeToEventsStore(
    "orders.events", "",
    kubemq::StartPosition::StartFromFirst,
    [](const kubemq::EventStoreReceived& msg) {
        std::cout << "Received: seq=" << msg.sequence()
                  << " body=" << msg.body() << std::endl;
    },
    [](const std::string& err) {
        std::cerr << "Error: " << err << std::endl;
    });
subscriber.rs
use kubemq::prelude::*;
use kubemq::EventsStoreSubscription;
use std::time::Duration;

#[tokio::main]
async fn main() -> kubemq::Result<()> {
    let client = KubemqClient::builder()
        .host("localhost")
        .port(50000)
        .build()
        .await?;

    println!("Replaying order history...");
    let sub = client
        .subscribe_to_events_store(
            "orders.events",
            "",
            EventsStoreSubscription::StartFromFirst,
            |event| {
                Box::pin(async move {
                    println!(
                        "Received: seq={} body={}",
                        event.sequence,
                        String::from_utf8_lossy(&event.body)
                    );
                })
            },
            None,
        )
        .await?;

    tokio::time::sleep(Duration::from_secs(120)).await;
    sub.unsubscribe().await;
    client.close().await?;
    Ok(())
}
subscriber.rb
require 'kubemq'

client = KubeMQ::PubSubClient.new(address: 'localhost:50000', client_id: 'order-subscriber')
cancel = KubeMQ::CancellationToken.new

sub = KubeMQ::PubSub::EventsStoreSubscription.new(
  channel: 'orders.events',
  start_position: KubeMQ::PubSub::EventStoreStartPosition::START_FROM_FIRST
)

client.subscribe_to_events_store(sub, cancellation_token: cancel, on_error: lambda { |e|
  puts "Error: #{e.message}"
}) do |event|
  puts "Received: seq=#{event.sequence} body=#{event.body}"
end

puts 'Replaying order history...'
sleep 120

cancel.cancel
client.close
subscriber.exs
{:ok, client} =
  KubeMQ.Client.start_link(address: "localhost:50000", client_id: "order-subscriber")

IO.puts("Replaying order history...")

{:ok, _sub} =
  KubeMQ.Client.subscribe_to_events_store(client, "orders.events",
    start_at: :start_from_first,
    on_event: fn event ->
      IO.puts("Received: seq=#{event.sequence} body=#{event.body}")
    end
  )

Process.sleep(120_000)
KubeMQ.Client.close(client)

Verify the Output

The subscriber receives all 3 previously published events, then continues waiting for new ones:

Replaying order history...
Received: seq=1 body={"action":"order.created","orderId":"ORD-1001","total":149.99}
Received: seq=2 body={"action":"order.paid","orderId":"ORD-1001","method":"credit_card"}
Received: seq=3 body={"action":"order.shipped","orderId":"ORD-1001","carrier":"fedex"}

Any new events published after the subscriber connects are delivered in real time.

Understanding What Happened

The store assigns each event a sequence number; a StartFromFirst subscriber replays the full history, then receives live events.

  1. The publisher stored 3 events in the orders.events channel
  2. Each event received a sequence number (1, 2, 3) and a server timestamp
  3. The subscriber connected later and requested StartFromFirst
  4. KubeMQ replayed the entire history, then continues delivering new events
  5. The subscription is durable — if the subscriber disconnects and reconnects with the same client ID and group, it resumes from the last delivered sequence

With plain Events, the subscriber would have received nothing because Events are not persisted. Events Store is essential when subscribers must not miss messages.

What's Next

Was this page helpful?

On this page