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:nextWhat 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/v2pip install kubemqnpm 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.CSharpimplementation("io.kubemq.sdk:kubemq-sdk-kotlin:2.1.0")vcpkg install kubemqcargo add kubemqgem install kubemq# mix.exs
def deps do
[{:kubemq, "~> 1.0"}]
endPublish Persistent Events
Publish several order events before starting the subscriber. With Events Store, messages are persisted and available for replay.
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)
}
}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}")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}`);
}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();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}");
}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}")
}
}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;
}
}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(())
}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{: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.
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()
}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)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...');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();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)}");
}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)}")
}
}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;
});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(())
}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{: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.
- The publisher stored 3 events in the
orders.eventschannel - Each event received a sequence number (1, 2, 3) and a server timestamp
- The subscriber connected later and requested
StartFromFirst - KubeMQ replayed the entire history, then continues delivering new events
- 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?