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
- KubeMQ server running on
localhost:50000 - SDK installed (Getting Started)
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.
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)
}
}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})")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}`);
}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();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");
}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")
}
}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;
}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(())
}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.closechannel = "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.
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)
}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}")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);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();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"]}");
}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}")
}
}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; });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(())
}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{: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
- Configure retention policies for long-lived streams
- Scale processing with consumer groups
- Build an audit trail with Events Store
- See the Events Store Reference for configuration options
Was this page helpful?