KubeMQ
LearnEvents StoreTutorials

Persistent Publish & Subscribe

Publish persistent events and subscribe with guaranteed delivery and replay capability.

This tutorial demonstrates the persistent pub/sub pattern with KubeMQ Events Store. You will publish events that are stored on disk and subscribe with different start positions to control which events you receive.

This is the deep-dive — multiple subscriber types (full-history replay vs. new-events-only) and durable subscriptions. New to Events Store? Start with the 5-minute getting-started quickstart first.

What You Will Build

An order tracking system where:

  • An order service publishes order lifecycle events to a persistent channel
  • An audit dashboard replays the full history from the beginning
  • A real-time alerter receives only new events going forward

The order service persists each lifecycle event to the store; the audit dashboard replays the full history while the real-time alerter receives only new events.

Prerequisites

Step-by-Step

Create the Event Publisher

The publisher sends order lifecycle events with metadata describing the event type.

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

    events := []struct {
        Action  string
        OrderID string
        Detail  string
    }{
        {"order.created", "ORD-1001", "New order placed, total=$149.99"},
        {"order.paid", "ORD-1001", "Payment confirmed via credit card"},
        {"order.picked", "ORD-1001", "Items picked from warehouse"},
        {"order.shipped", "ORD-1001", "Shipped via FedEx, tracking=FX-9876"},
        {"order.delivered", "ORD-1001", "Delivered to customer"},
    }

    for _, e := range events {
        body := fmt.Sprintf(`{"action":"%s","orderId":"%s","detail":"%s"}`,
            e.Action, e.OrderID, e.Detail)

        result, err := client.SendEventStore(ctx, kubemq.NewEvent().
            SetChannel("orders.lifecycle").
            SetMetadata(e.Action).
            SetBody([]byte(body)),
        )
        if err != nil {
            log.Printf("Failed to store: %v", err)
            continue
        }
        log.Printf("Stored [%s]: %s (ID: %s)", e.Action, e.OrderID, result.EventID)
        time.Sleep(200 * time.Millisecond)
    }
}
order_publisher.py
import json
import time
from kubemq import PubSubClient, EventStoreMessage

events = [
    {"action": "order.created", "orderId": "ORD-1001", "detail": "New order, total=$149.99"},
    {"action": "order.paid", "orderId": "ORD-1001", "detail": "Payment confirmed"},
    {"action": "order.picked", "orderId": "ORD-1001", "detail": "Items picked"},
    {"action": "order.shipped", "orderId": "ORD-1001", "detail": "Shipped via FedEx"},
    {"action": "order.delivered", "orderId": "ORD-1001", "detail": "Delivered"},
]

with PubSubClient(address="localhost:50000") as client:
    for e in events:
        result = client.publish_event_store(
            EventStoreMessage(
                channel="orders.lifecycle",
                metadata=e["action"],
                body=json.dumps(e).encode("utf-8"),
            )
        )
        print(f"Stored [{e['action']}]: {e['orderId']} (ID: {result.id})")
        time.sleep(0.2)
order_publisher.ts
import { KubeMQClient, createEventStoreMessage } from 'kubemq-js';

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

const events = [
  { action: 'order.created', orderId: 'ORD-1001', detail: 'New order, total=$149.99' },
  { action: 'order.paid', orderId: 'ORD-1001', detail: 'Payment confirmed' },
  { action: 'order.picked', orderId: 'ORD-1001', detail: 'Items picked' },
  { action: 'order.shipped', orderId: 'ORD-1001', detail: 'Shipped via FedEx' },
  { action: 'order.delivered', orderId: 'ORD-1001', detail: 'Delivered' },
];

for (const e of events) {
  const result = await client.sendEventStore(
    createEventStoreMessage({
      channel: 'orders.lifecycle',
      metadata: e.action,
      body: JSON.stringify(e),
    })
  );
  console.log(`Stored [${e.action}]: ${e.orderId} (ID: ${result.id})`);
  await new Promise((r) => setTimeout(r, 200));
}
OrderPublisher.java
PubSubClient client = PubSubClient.builder()
    .address("localhost:50000")
    .clientId("order-publisher")
    .build();

String[][] events = {
    {"order.created", "ORD-1001", "New order, total=$149.99"},
    {"order.paid", "ORD-1001", "Payment confirmed"},
    {"order.picked", "ORD-1001", "Items picked"},
    {"order.shipped", "ORD-1001", "Shipped via FedEx"},
    {"order.delivered", "ORD-1001", "Delivered"},
};

for (String[] e : events) {
    String body = String.format(
        "{\"action\":\"%s\",\"orderId\":\"%s\",\"detail\":\"%s\"}", e[0], e[1], e[2]);
    EventSendResult result = client.sendEventsStoreMessage(
        EventStoreMessage.builder()
            .channel("orders.lifecycle")
            .metadata(e[0])
            .body(body.getBytes())
            .build());
    System.out.printf("Stored [%s]: %s (ID: %s)%n", e[0], e[1], result.getId());
    Thread.sleep(200);
}
client.close();
OrderPublisher.cs
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

var events = new[] {
    ("order.created", "ORD-1001", "New order, total=$149.99"),
    ("order.paid", "ORD-1001", "Payment confirmed"),
    ("order.picked", "ORD-1001", "Items picked"),
    ("order.shipped", "ORD-1001", "Shipped via FedEx"),
    ("order.delivered", "ORD-1001", "Delivered"),
};

foreach (var (action, orderId, detail) in events)
{
    var body = $"{{\"action\":\"{action}\",\"orderId\":\"{orderId}\",\"detail\":\"{detail}\"}}";
    var result = await client.SendEventStoreAsync(new EventStoreMessage
    {
        Channel = "orders.lifecycle",
        Metadata = action,
        Body = Encoding.UTF8.GetBytes(body),
    });
    Console.WriteLine($"Stored [{action}]: {orderId} (ID: {result.Id})");
    await Task.Delay(200);
}
OrderPublisher.kt
val client = KubeMQClient.pubSub {
    address = "localhost:50000"
    clientId = "order-publisher"
}

data class OrderEvent(val action: String, val orderId: String, val detail: String)

val events = listOf(
    OrderEvent("order.created", "ORD-1001", "New order, total=\$149.99"),
    OrderEvent("order.paid", "ORD-1001", "Payment confirmed"),
    OrderEvent("order.picked", "ORD-1001", "Items picked"),
    OrderEvent("order.shipped", "ORD-1001", "Shipped via FedEx"),
    OrderEvent("order.delivered", "ORD-1001", "Delivered"),
)

client.use {
    for (e in events) {
        val body = """{"action":"${e.action}","orderId":"${e.orderId}","detail":"${e.detail}"}"""
        val result = client.sendEventStore(eventStoreMessage {
            channel = "orders.lifecycle"
            metadata = e.action
            this.body = body.toByteArray()
        })
        println("Stored [${e.action}]: ${e.orderId} (ID: ${result.id})")
        delay(200)
    }
}
order_publisher.cc
kubemq::ClientOptions options;
options.set_address("localhost", 50000);
options.set_client_id("order-publisher");

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

struct OrderEvent { std::string action, orderId, detail; };
std::vector<OrderEvent> events = {
    {"order.created", "ORD-1001", "New order, total=$149.99"},
    {"order.paid", "ORD-1001", "Payment confirmed"},
    {"order.picked", "ORD-1001", "Items picked"},
    {"order.shipped", "ORD-1001", "Shipped via FedEx"},
    {"order.delivered", "ORD-1001", "Delivered"},
};

for (const auto& e : events) {
    kubemq::EventStoreMessage msg;
    msg.set_channel("orders.lifecycle");
    msg.set_metadata(e.action);
    msg.set_body("{\"action\":\"" + e.action + "\",\"orderId\":\"" + e.orderId + "\"}");
    auto result = client->SendEventStore(msg);
    if (result.ok()) {
        std::cout << "Stored [" << e.action << "]: " << e.orderId << std::endl;
    }
    std::this_thread::sleep_for(std::chrono::milliseconds(200));
}
order_publisher.rs
use kubemq::prelude::*;
use kubemq::EventStoreBuilder;
use std::time::Duration;

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

    let events = [
        ("order.created", "ORD-1001", "New order, total=$149.99"),
        ("order.paid", "ORD-1001", "Payment confirmed"),
        ("order.picked", "ORD-1001", "Items picked"),
        ("order.shipped", "ORD-1001", "Shipped via FedEx"),
        ("order.delivered", "ORD-1001", "Delivered"),
    ];

    for (action, order_id, detail) in events {
        let body = format!(
            r#"{{"action":"{action}","orderId":"{order_id}","detail":"{detail}"}}"#
        );
        let event = EventStoreBuilder::new()
            .channel("orders.lifecycle")
            .metadata(action)
            .body(body.into_bytes())
            .build();
        let result = client.send_event_store(event).await?;
        println!("Stored [{action}]: {order_id} (ID: {})", result.id);
        tokio::time::sleep(Duration::from_millis(200)).await;
    }

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

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

events = [
  ['order.created', 'ORD-1001', 'New order, total=$149.99'],
  ['order.paid', 'ORD-1001', 'Payment confirmed'],
  ['order.picked', 'ORD-1001', 'Items picked'],
  ['order.shipped', 'ORD-1001', 'Shipped via FedEx'],
  ['order.delivered', 'ORD-1001', 'Delivered']
]

events.each do |action, order_id, detail|
  body = { action: action, orderId: order_id, detail: detail }.to_json
  msg = KubeMQ::PubSub::EventStoreMessage.new(
    channel: 'orders.lifecycle',
    metadata: action,
    body: body
  )
  result = client.send_event_store(msg)
  puts "Stored [#{action}]: #{order_id} (sent: #{result.sent})"
  sleep 0.2
end

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

events = [
  {"order.created", "ORD-1001", "New order, total=$149.99"},
  {"order.paid", "ORD-1001", "Payment confirmed"},
  {"order.picked", "ORD-1001", "Items picked"},
  {"order.shipped", "ORD-1001", "Shipped via FedEx"},
  {"order.delivered", "ORD-1001", "Delivered"}
]

for {action, order_id, detail} <- events do
  body = Jason.encode!(%{action: action, orderId: order_id, detail: detail})

  event =
    KubeMQ.EventStore.new(channel: "orders.lifecycle", metadata: action, body: body)

  case KubeMQ.Client.send_event_store(client, event) do
    {:ok, result} -> IO.puts("Stored [#{action}]: #{order_id} (sent: #{result.sent})")
    {:error, err} -> IO.puts("Store failed: #{err.message}")
  end

  Process.sleep(200)
end

KubeMQ.Client.close(client)

Subscribe to Full History (StartFromFirst)

The audit dashboard connects after events are stored and replays the full history.

audit_dashboard.go
sub, err := client.SubscribeToEventsStore(ctx, "orders.lifecycle", "",
    kubemq.StartFromFirst(),
    kubemq.WithOnEvent(func(event *kubemq.Event) {
        fmt.Printf("[Audit] seq=%d action=%s body=%s\n",
            event.Sequence, event.Metadata, string(event.Body))
    }),
    kubemq.WithOnError(func(err error) {
        log.Println("[Audit] Error:", err)
    }),
)
audit_dashboard.py
client.subscribe_to_events_store(
    subscription=EventsStoreSubscription(
        channel="orders.lifecycle",
        start_position=EventStoreStartPosition.StartFromFirst,
        on_receive_event_callback=lambda e: print(
            f"[Audit] seq={e.sequence} action={e.metadata} "
            f"body={e.body.decode('utf-8')}"
        ),
        on_error_callback=lambda e: print(f"[Audit] Error: {e}"),
    ),
    cancel=CancellationToken(),
)
audit_dashboard.ts
client.subscribeToEventsStore({
  channel: 'orders.lifecycle',
  startPosition: EventStoreStartPosition.StartFromFirst,
  onEvent: (msg) =>
    console.log(
      `[Audit] seq=${msg.sequence} action=${msg.metadata} ` +
        `body=${new TextDecoder().decode(msg.body)}`
    ),
  onError: (err) => console.error('[Audit] Error:', err.message),
});
AuditDashboard.java
client.subscribeToEventsStore(EventsStoreSubscription.builder()
    .channel("orders.lifecycle")
    .startPosition(EventStoreStartPosition.StartFromFirst)
    .onReceiveEventCallback(event ->
        System.out.printf("[Audit] seq=%d action=%s body=%s%n",
            event.getSequence(), event.getMetadata(),
            new String(event.getBody())))
    .onErrorCallback(err ->
        System.err.println("[Audit] Error: " + err.getMessage()))
    .build());
AuditDashboard.cs
await foreach (var msg in client.SubscribeToEventsStoreAsync(
    new EventsStoreSubscription
    {
        Channel = "orders.lifecycle",
        StartPosition = EventStoreStartPosition.StartFromFirst,
    }))
{
    Console.WriteLine($"[Audit] seq={msg.Sequence} action={msg.Metadata} "
        + $"body={Encoding.UTF8.GetString(msg.Body.Span)}");
}
AuditDashboard.kt
client.subscribeToEventsStore {
    channel = "orders.lifecycle"
    startPosition = StartPosition.StartFromFirst
}.collect { msg ->
    println("[Audit] seq=${msg.sequence} action=${msg.metadata} body=${String(msg.body)}")
}
audit_dashboard.cc
client->SubscribeToEventsStore(
    "orders.lifecycle", "",
    kubemq::StartPosition::StartFromFirst,
    [](const kubemq::EventStoreReceived& msg) {
        std::cout << "[Audit] seq=" << msg.sequence()
                  << " action=" << msg.metadata()
                  << " body=" << msg.body() << std::endl;
    },
    [](const std::string& err) {
        std::cerr << "[Audit] Error: " << err << std::endl;
    });
audit_dashboard.rs
use kubemq::EventsStoreSubscription;

let sub = client
    .subscribe_to_events_store(
        "orders.lifecycle",
        "",
        EventsStoreSubscription::StartFromFirst,
        |event| {
            Box::pin(async move {
                println!(
                    "[Audit] seq={} action={} body={}",
                    event.sequence,
                    event.metadata,
                    String::from_utf8_lossy(&event.body)
                );
            })
        },
        None,
    )
    .await?;
audit_dashboard.rb
cancel = KubeMQ::CancellationToken.new

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

client.subscribe_to_events_store(sub, cancellation_token: cancel, on_error: lambda { |e|
  puts "[Audit] Error: #{e.message}"
}) do |event|
  puts "[Audit] seq=#{event.sequence} action=#{event.metadata} body=#{event.body}"
end
audit_dashboard.exs
{:ok, sub} =
  KubeMQ.Client.subscribe_to_events_store(client, "orders.lifecycle",
    start_at: :start_from_first,
    on_event: fn event ->
      IO.puts(
        "[Audit] seq=#{event.sequence} action=#{event.metadata} body=#{event.body}"
      )
    end,
    on_error: fn err -> IO.puts("[Audit] Error: #{err.message}") end
  )

Expected output — all 5 events replayed:

[Audit] seq=1 action=order.created body={"action":"order.created","orderId":"ORD-1001",...}
[Audit] seq=2 action=order.paid body={...}
[Audit] seq=3 action=order.picked body={...}
[Audit] seq=4 action=order.shipped body={...}
[Audit] seq=5 action=order.delivered body={...}

Subscribe to New Events Only (StartNewOnly)

The real-time alerter receives only events published after it subscribes.

realtime_alerter.go
sub, err := client.SubscribeToEventsStore(ctx, "orders.lifecycle", "",
    kubemq.StartNewOnly(),
    kubemq.WithOnEvent(func(event *kubemq.Event) {
        fmt.Printf("[Alert] New: %s\n", string(event.Body))
    }),
    kubemq.WithOnError(func(err error) {
        log.Println("[Alert] Error:", err)
    }),
)
realtime_alerter.py
client.subscribe_to_events_store(
    subscription=EventsStoreSubscription(
        channel="orders.lifecycle",
        start_position=EventStoreStartPosition.StartNewOnly,
        on_receive_event_callback=lambda e: print(
            f"[Alert] New: {e.body.decode('utf-8')}"
        ),
        on_error_callback=lambda e: print(f"[Alert] Error: {e}"),
    ),
    cancel=CancellationToken(),
)
realtime_alerter.ts
client.subscribeToEventsStore({
  channel: 'orders.lifecycle',
  startPosition: EventStoreStartPosition.StartNewOnly,
  onEvent: (msg) =>
    console.log(`[Alert] New: ${new TextDecoder().decode(msg.body)}`),
  onError: (err) => console.error('[Alert] Error:', err.message),
});
RealtimeAlerter.java
client.subscribeToEventsStore(EventsStoreSubscription.builder()
    .channel("orders.lifecycle")
    .startPosition(EventStoreStartPosition.StartNewOnly)
    .onReceiveEventCallback(event ->
        System.out.printf("[Alert] New: %s%n", new String(event.getBody())))
    .onErrorCallback(err ->
        System.err.println("[Alert] Error: " + err.getMessage()))
    .build());
RealtimeAlerter.cs
await foreach (var msg in client.SubscribeToEventsStoreAsync(
    new EventsStoreSubscription
    {
        Channel = "orders.lifecycle",
        StartPosition = EventStoreStartPosition.StartNewOnly,
    }))
{
    Console.WriteLine($"[Alert] New: {Encoding.UTF8.GetString(msg.Body.Span)}");
}
RealtimeAlerter.kt
client.subscribeToEventsStore {
    channel = "orders.lifecycle"
    startPosition = StartPosition.StartNewOnly
}.collect { msg ->
    println("[Alert] New: ${String(msg.body)}")
}
realtime_alerter.cc
client->SubscribeToEventsStore(
    "orders.lifecycle", "",
    kubemq::StartPosition::StartNewOnly,
    [](const kubemq::EventStoreReceived& msg) {
        std::cout << "[Alert] New: " << msg.body() << std::endl;
    },
    [](const std::string& err) {
        std::cerr << "[Alert] Error: " << err << std::endl;
    });
realtime_alerter.rs
use kubemq::EventsStoreSubscription;

let sub = client
    .subscribe_to_events_store(
        "orders.lifecycle",
        "",
        EventsStoreSubscription::StartNewOnly,
        |event| {
            Box::pin(async move {
                println!("[Alert] New: {}", String::from_utf8_lossy(&event.body));
            })
        },
        None,
    )
    .await?;
realtime_alerter.rb
cancel = KubeMQ::CancellationToken.new

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

client.subscribe_to_events_store(sub, cancellation_token: cancel, on_error: lambda { |e|
  puts "[Alert] Error: #{e.message}"
}) do |event|
  puts "[Alert] New: #{event.body}"
end
realtime_alerter.exs
{:ok, sub} =
  KubeMQ.Client.subscribe_to_events_store(client, "orders.lifecycle",
    start_at: :start_new_only,
    on_event: fn event -> IO.puts("[Alert] New: #{event.body}") end,
    on_error: fn err -> IO.puts("[Alert] Error: #{err.message}") end
  )

This subscriber receives nothing from the 5 previously stored events, but will receive any new events published after subscribing.

Key Concepts

Persistence Guarantee

Events Store writes events to the underlying store before acknowledging the publish. The publish call returns after the message is confirmed stored, providing at-least-once delivery semantics to subscribers.

Durable Subscriptions

Each Events Store subscription creates a durable name based on the channel and group:

DurableName = "{channel}-{group}"

If a subscriber disconnects and reconnects with the same durable name, the store resumes delivery from the last acknowledged position, regardless of the StartPosition specified.

Events vs Events Store

FeatureEventsEvents Store
PersistenceNo (memory only)Yes (disk-backed)
Late subscriber receives historyNoYes (via start positions)
Delivery guaranteeAt-most-onceAt-least-once
WildcardsYesNo

Sequence numbers are assigned per channel. Different channels have independent sequences starting from 1.

Next Steps

Was this page helpful?

On this page