KubeMQ
LearnEvents StoreTutorials

Event Sourcing Pattern

Implement event sourcing with KubeMQ Events Store to rebuild state from stored events.

Event sourcing stores every state change as an immutable event rather than overwriting the current state. KubeMQ Events Store is a natural fit because it provides persistent, sequenced, and replayable event streams.

What You Will Build

An order management service that:

  • Stores every order state change (created, paid, shipped, delivered) as an event
  • Reconstructs the current order state by replaying the full event history
  • Supports checkpoint-based recovery to avoid full replays

Every order state change is appended to the event store; replaying the full sequence reconstructs the current order state.

Prerequisites

Step-by-Step

Define the Event Schema

Each event includes:

  • type — the event kind (e.g., order.created, order.paid)
  • orderId — the aggregate identifier
  • data — event-specific payload
  • timestamp — when the state change occurred
{
  "type": "order.created",
  "orderId": "ORD-1001",
  "data": { "customer": "C-500", "items": [{"sku": "WIDGET-A", "qty": 2, "price": 49.99}] },
  "timestamp": "2026-03-26T10:00:00Z"
}

Publish Order Events

Store a series of state changes for an order. Each event is immutable once stored.

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

    channel := "order.ORD-1001"
    events := []string{
        `{"type":"order.created","orderId":"ORD-1001","data":{"customer":"C-500","total":149.99}}`,
        `{"type":"order.paid","orderId":"ORD-1001","data":{"method":"credit_card","txId":"TX-789"}}`,
        `{"type":"order.shipped","orderId":"ORD-1001","data":{"carrier":"fedex","tracking":"FX-123"}}`,
    }

    for _, body := range events {
        result, err := client.SendEventStore(ctx, kubemq.NewEvent().
            SetChannel(channel).
            SetBody([]byte(body)),
        )
        if err != nil {
            log.Fatal(err)
        }
        log.Printf("Stored: %s (seq: %s)", body[:40], result.EventID)
    }
}
order_service.py
import json
from kubemq import PubSubClient, EventStoreMessage

events = [
    {"type": "order.created", "orderId": "ORD-1001",
     "data": {"customer": "C-500", "total": 149.99}},
    {"type": "order.paid", "orderId": "ORD-1001",
     "data": {"method": "credit_card", "txId": "TX-789"}},
    {"type": "order.shipped", "orderId": "ORD-1001",
     "data": {"carrier": "fedex", "tracking": "FX-123"}},
]

with PubSubClient(address="localhost:50000") as client:
    for event in events:
        result = client.publish_event_store(
            EventStoreMessage(
                channel="order.ORD-1001",
                body=json.dumps(event).encode("utf-8"),
            )
        )
        print(f"Stored: {event['type']} (ID: {result.id})")
order_service.ts
import { KubeMQClient, createEventStoreMessage } from 'kubemq-js';

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

const events = [
  { type: 'order.created', orderId: 'ORD-1001', data: { customer: 'C-500', total: 149.99 } },
  { type: 'order.paid', orderId: 'ORD-1001', data: { method: 'credit_card', txId: 'TX-789' } },
  { type: 'order.shipped', orderId: 'ORD-1001', data: { carrier: 'fedex', tracking: 'FX-123' } },
];

for (const event of events) {
  await client.sendEventStore(
    createEventStoreMessage({
      channel: 'order.ORD-1001',
      body: JSON.stringify(event),
    })
  );
  console.log(`Stored: ${event.type}`);
}
OrderService.java
PubSubClient client = PubSubClient.builder()
    .address("localhost:50000")
    .clientId("order-service")
    .build();

String[] events = {
    "{\"type\":\"order.created\",\"orderId\":\"ORD-1001\",\"data\":{\"customer\":\"C-500\",\"total\":149.99}}",
    "{\"type\":\"order.paid\",\"orderId\":\"ORD-1001\",\"data\":{\"method\":\"credit_card\"}}",
    "{\"type\":\"order.shipped\",\"orderId\":\"ORD-1001\",\"data\":{\"carrier\":\"fedex\"}}"
};

for (String body : events) {
    client.sendEventsStoreMessage(EventStoreMessage.builder()
        .channel("order.ORD-1001")
        .body(body.getBytes())
        .build());
    System.out.println("Stored event");
}
client.close();
OrderService.cs
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

string[] events = {
    "{\"type\":\"order.created\",\"orderId\":\"ORD-1001\",\"data\":{\"customer\":\"C-500\",\"total\":149.99}}",
    "{\"type\":\"order.paid\",\"orderId\":\"ORD-1001\",\"data\":{\"method\":\"credit_card\"}}",
    "{\"type\":\"order.shipped\",\"orderId\":\"ORD-1001\",\"data\":{\"carrier\":\"fedex\"}}"
};

foreach (var body in events)
{
    await client.SendEventStoreAsync(new EventStoreMessage
    {
        Channel = "order.ORD-1001",
        Body = Encoding.UTF8.GetBytes(body),
    });
    Console.WriteLine("Stored event");
}
OrderService.kt
val client = KubeMQClient.pubSub {
    address = "localhost:50000"
    clientId = "order-service"
}

val events = listOf(
    """{"type":"order.created","orderId":"ORD-1001","data":{"customer":"C-500","total":149.99}}""",
    """{"type":"order.paid","orderId":"ORD-1001","data":{"method":"credit_card"}}""",
    """{"type":"order.shipped","orderId":"ORD-1001","data":{"carrier":"fedex"}}""",
)

client.use {
    for (body in events) {
        client.sendEventStore(eventStoreMessage {
            channel = "order.ORD-1001"
            this.body = body.toByteArray()
        })
        println("Stored event")
    }
}
order_service.cc
kubemq::ClientOptions options;
options.set_address("localhost", 50000);
options.set_client_id("order-service");
auto client = kubemq::Client::Create(options).value();

std::vector<std::string> events = {
    R"({"type":"order.created","orderId":"ORD-1001","data":{"customer":"C-500","total":149.99}})",
    R"({"type":"order.paid","orderId":"ORD-1001","data":{"method":"credit_card"}})",
    R"({"type":"order.shipped","orderId":"ORD-1001","data":{"carrier":"fedex"}})"
};

for (const auto& body : events) {
    kubemq::EventStoreMessage msg;
    msg.set_channel("order.ORD-1001");
    msg.set_body(body);
    client->SendEventStore(msg);
    std::cout << "Stored event" << std::endl;
}
order_service.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 channel = "order.ORD-1001";
    let events = [
        r#"{"type":"order.created","orderId":"ORD-1001","data":{"customer":"C-500","total":149.99}}"#,
        r#"{"type":"order.paid","orderId":"ORD-1001","data":{"method":"credit_card","txId":"TX-789"}}"#,
        r#"{"type":"order.shipped","orderId":"ORD-1001","data":{"carrier":"fedex","tracking":"FX-123"}}"#,
    ];

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

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

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

channel = 'order.ORD-1001'
events = [
  '{"type":"order.created","orderId":"ORD-1001","data":{"customer":"C-500","total":149.99}}',
  '{"type":"order.paid","orderId":"ORD-1001","data":{"method":"credit_card","txId":"TX-789"}}',
  '{"type":"order.shipped","orderId":"ORD-1001","data":{"carrier":"fedex","tracking":"FX-123"}}'
]

events.each do |body|
  result = client.send_event_store(
    KubeMQ::PubSub::EventStoreMessage.new(channel: channel, body: body)
  )
  puts "Stored event: sent=#{result.sent}"
end

client.close
order_service.exs
channel = "order.ORD-1001"
{:ok, client} = KubeMQ.Client.start_link(address: "localhost:50000", client_id: "order-service")

events = [
  ~s({"type":"order.created","orderId":"ORD-1001","data":{"customer":"C-500","total":149.99}}),
  ~s({"type":"order.paid","orderId":"ORD-1001","data":{"method":"credit_card","txId":"TX-789"}}),
  ~s({"type":"order.shipped","orderId":"ORD-1001","data":{"carrier":"fedex","tracking":"FX-123"}})
]

for body <- events do
  {:ok, result} =
    KubeMQ.Client.send_event_store(client, KubeMQ.EventStore.new(channel: channel, body: body))

  IO.puts("Stored event: sent=#{result.sent}")
end

KubeMQ.Client.close(client)

Rebuild State from Event History

Subscribe with StartFromFirst to replay all events and compute the current order state.

state_rebuilder.go
package main

import (
    "context"
    "encoding/json"
    "fmt"
    "log"
    "sync"
    "time"

    "github.com/kubemq-io/kubemq-go/v2"
)

type OrderState struct {
    OrderID  string
    Status   string
    Customer string
    Total    float64
    Carrier  string
    Tracking string
}

func main() {
    ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
    defer cancel()

    client, err := kubemq.NewClient(ctx,
        kubemq.WithAddress("localhost", 50000),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer client.Close()

    var mu sync.Mutex
    order := OrderState{OrderID: "ORD-1001"}

    sub, err := client.SubscribeToEventsStore(ctx, "order.ORD-1001", "",
        kubemq.StartFromFirst(),
        kubemq.WithOnEvent(func(event *kubemq.Event) {
            mu.Lock()
            defer mu.Unlock()

            var evt map[string]interface{}
            json.Unmarshal(event.Body, &evt)
            eventType := evt["type"].(string)
            data := evt["data"].(map[string]interface{})

            switch eventType {
            case "order.created":
                order.Status = "created"
                order.Customer = data["customer"].(string)
                order.Total = data["total"].(float64)
            case "order.paid":
                order.Status = "paid"
            case "order.shipped":
                order.Status = "shipped"
                order.Carrier = data["carrier"].(string)
            }

            fmt.Printf("seq=%d %s -> status=%s\n", event.Sequence, eventType, order.Status)
        }),
        kubemq.WithOnError(func(err error) { log.Println("Error:", err) }),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer sub.Unsubscribe()

    <-ctx.Done()
    fmt.Printf("\nOrder State: %+v\n", order)
}
state_rebuilder.py
import json
import time
from kubemq import (
    PubSubClient, EventsStoreSubscription,
    EventStoreStartPosition, CancellationToken,
)

order = {"orderId": "ORD-1001", "status": "unknown"}

def on_event(event):
    global order
    evt = json.loads(event.body.decode("utf-8"))
    event_type = evt["type"]
    data = evt.get("data", {})

    if event_type == "order.created":
        order["status"] = "created"
        order["customer"] = data.get("customer")
        order["total"] = data.get("total")
    elif event_type == "order.paid":
        order["status"] = "paid"
    elif event_type == "order.shipped":
        order["status"] = "shipped"
        order["carrier"] = data.get("carrier")

    print(f"seq={event.sequence} {event_type} -> status={order['status']}")

with PubSubClient(address="localhost:50000") as client:
    client.subscribe_to_events_store(
        subscription=EventsStoreSubscription(
            channel="order.ORD-1001",
            start_position=EventStoreStartPosition.StartFromFirst,
            on_receive_event_callback=on_event,
            on_error_callback=lambda e: print(f"Error: {e}"),
        ),
        cancel=CancellationToken(),
    )
    time.sleep(5)
    print(f"\nOrder State: {order}")
state_rebuilder.ts
import { KubeMQClient, EventStoreStartPosition } from 'kubemq-js';

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

const order: Record<string, unknown> = { orderId: 'ORD-1001', status: 'unknown' };

client.subscribeToEventsStore({
  channel: 'order.ORD-1001',
  startFrom: EventStoreStartPosition.StartFromFirst,
  onEvent: (msg) => {
    const evt = JSON.parse(new TextDecoder().decode(msg.body));
    const { type, data } = evt;

    if (type === 'order.created') {
      Object.assign(order, { status: 'created', customer: data.customer, total: data.total });
    } else if (type === 'order.paid') {
      order.status = 'paid';
    } else if (type === 'order.shipped') {
      Object.assign(order, { status: 'shipped', carrier: data.carrier });
    }
    console.log(`seq=${msg.sequence} ${type} -> status=${order.status}`);
  },
  onError: (err) => console.error('Error:', err.message),
});

setTimeout(() => console.log('\nOrder State:', order), 5000);
StateRebuilder.java
PubSubClient client = PubSubClient.builder()
    .address("localhost:50000")
    .clientId("state-rebuilder")
    .build();

Map<String, Object> order = new ConcurrentHashMap<>();
order.put("orderId", "ORD-1001");

client.subscribeToEventsStore(EventsStoreSubscription.builder()
    .channel("order.ORD-1001")
    .startPosition(EventStoreStartPosition.StartFromFirst)
    .onReceiveEventCallback(event -> {
        var evt = new ObjectMapper().readTree(event.getBody());
        String type = evt.get("type").asText();
        if ("order.created".equals(type)) {
            order.put("status", "created");
            order.put("total", evt.get("data").get("total").asDouble());
        } else if ("order.paid".equals(type)) {
            order.put("status", "paid");
        } else if ("order.shipped".equals(type)) {
            order.put("status", "shipped");
        }
        System.out.printf("seq=%d %s -> status=%s%n",
            event.getSequence(), type, order.get("status"));
    })
    .onErrorCallback(err -> System.err.println(err.getMessage()))
    .build());

Thread.sleep(5000);
System.out.println("\nOrder State: " + order);
client.close();
StateRebuilder.cs
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

var order = new Dictionary<string, object> { ["orderId"] = "ORD-1001" };

await foreach (var msg in client.SubscribeToEventsStoreAsync(
    new EventsStoreSubscription
    {
        Channel = "order.ORD-1001",
        StartPosition = EventStoreStartPosition.StartFromFirst,
    }))
{
    var evt = JsonSerializer.Deserialize<JsonElement>(msg.Body.Span);
    var type = evt.GetProperty("type").GetString()!;

    order["status"] = type switch
    {
        "order.created" => "created",
        "order.paid" => "paid",
        "order.shipped" => "shipped",
        _ => order.GetValueOrDefault("status", "unknown")!
    };
    Console.WriteLine($"seq={msg.Sequence} {type} -> status={order["status"]}");
}
StateRebuilder.kt
val client = KubeMQClient.pubSub {
    address = "localhost:50000"
    clientId = "state-rebuilder"
}

data class OrderState(var status: String = "unknown", var total: Double = 0.0)
val order = OrderState()

client.use {
    client.subscribeToEventsStore {
        channel = "order.ORD-1001"
        startPosition = StartPosition.StartFromFirst
    }.collect { msg ->
        val evt = JSONObject(String(msg.body))
        val type = evt.getString("type")
        when (type) {
            "order.created" -> { order.status = "created"; order.total = evt.getJSONObject("data").getDouble("total") }
            "order.paid" -> order.status = "paid"
            "order.shipped" -> order.status = "shipped"
        }
        println("seq=${msg.sequence} $type -> status=${order.status}")
    }
}
state_rebuilder.cc
kubemq::ClientOptions options;
options.set_address("localhost", 50000);
auto client = kubemq::Client::Create(options).value();

std::string status = "unknown";

client->SubscribeToEventsStore("order.ORD-1001", "",
    kubemq::StartPosition::StartFromFirst,
    [&status](const kubemq::EventStoreReceived& msg) {
        auto body = msg.body();
        if (body.find("order.created") != std::string::npos) status = "created";
        else if (body.find("order.paid") != std::string::npos) status = "paid";
        else if (body.find("order.shipped") != std::string::npos) status = "shipped";
        std::cout << "seq=" << msg.sequence() << " -> status=" << status << std::endl;
    },
    [](const std::string& err) { std::cerr << err << std::endl; });
state_rebuilder.rs
use kubemq::prelude::*;
use kubemq::EventsStoreSubscription;
use serde_json::Value;
use std::sync::{Arc, Mutex};
use std::time::Duration;

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

    let status = Arc::new(Mutex::new(String::from("unknown")));
    let status_cb = status.clone();

    // Replay every stored event from the beginning to rebuild state.
    let sub = client
        .subscribe_to_events_store(
            "order.ORD-1001",
            "",
            EventsStoreSubscription::StartFromFirst,
            move |event| {
                let status = status_cb.clone();
                Box::pin(async move {
                    let evt: Value = serde_json::from_slice(&event.body).unwrap_or(Value::Null);
                    let event_type = evt["type"].as_str().unwrap_or("");
                    let mut s = status.lock().unwrap();
                    match event_type {
                        "order.created" => *s = "created".into(),
                        "order.paid" => *s = "paid".into(),
                        "order.shipped" => *s = "shipped".into(),
                        _ => {}
                    }
                    println!("seq={} {} -> status={}", event.sequence, event_type, *s);
                })
            },
            None,
        )
        .await?;

    tokio::time::sleep(Duration::from_secs(5)).await;
    println!("Order State: status={}", *status.lock().unwrap());

    sub.unsubscribe().await;
    client.close().await?;
    Ok(())
}
state_rebuilder.rb
require 'kubemq'
require 'json'

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

order = { 'orderId' => 'ORD-1001', 'status' => 'unknown' }

# Replay all stored events from the first to rebuild current state.
sub = KubeMQ::PubSub::EventsStoreSubscription.new(
  channel: 'order.ORD-1001',
  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|
  evt = JSON.parse(event.body)
  case evt['type']
  when 'order.created' then order['status'] = 'created'
  when 'order.paid'    then order['status'] = 'paid'
  when 'order.shipped' then order['status'] = 'shipped'
  end
  puts "seq=#{event.sequence} #{evt['type']} -> status=#{order['status']}"
end

sleep 5
puts "Order State: #{order}"
cancel.cancel
client.close
state_rebuilder.exs
{:ok, client} = KubeMQ.Client.start_link(address: "localhost:50000", client_id: "state-rebuilder")

{:ok, agent} = Agent.start_link(fn -> "unknown" end)

# Replay all stored events from the first to rebuild current state.
{:ok, sub} =
  KubeMQ.Client.subscribe_to_events_store(client, "order.ORD-1001",
    start_at: :start_from_first,
    on_event: fn event ->
      evt = Jason.decode!(event.body)

      status =
        case evt["type"] do
          "order.created" -> "created"
          "order.paid" -> "paid"
          "order.shipped" -> "shipped"
          _ -> Agent.get(agent, & &1)
        end

      Agent.update(agent, fn _ -> status end)
      IO.puts("seq=#{event.sequence} #{evt["type"]} -> status=#{status}")
    end
  )

Process.sleep(5_000)
IO.puts("Order State: status=#{Agent.get(agent, & &1)}")

KubeMQ.Subscription.cancel(sub)
KubeMQ.Client.close(client)

Expected output:

seq=1 order.created -> status=created
seq=2 order.paid -> status=paid
seq=3 order.shipped -> status=shipped

Order State: {orderId=ORD-1001, status=shipped, customer=C-500, total=149.99}

Checkpoint-Based Recovery

For aggregates with long histories, save the last processed sequence number as a checkpoint. On restart, subscribe from that sequence instead of replaying the entire history.

A checkpoint records the last processed sequence so restarts replay only the unprocessed tail instead of the full history.

The subscriber uses StartAtSequence(lastCheckpoint + 1) to pick up only unprocessed events.

Event Sourcing Best Practices

Channel Per Aggregate

Use one channel per aggregate instance (e.g., order.ORD-1001, order.ORD-1002). This provides independent sequencing and clean replay per entity.

Immutable Events

Never modify or delete events. The event stream is the source of truth. To correct an error, append a compensating event (e.g., an order.refunded event).

Snapshots for Performance

For aggregates with long histories, periodically save a snapshot (the computed state at a sequence number). On startup, load the snapshot and replay only events after the snapshot sequence.

Event Schema Versioning

Include a version field in your event schema. When the schema changes, handle both old and new formats in your event handler.

For event sourcing, configure retention to unlimited (Store.MaxRetention=0) or set it longer than your maximum replay window. See Configure Retention.

Next Steps

Was this page helpful?

On this page