Cross-Service State Sync
Synchronize state across microservices using KubeMQ Events Store as a shared event log.
This scenario synchronizes state across independent microservices using Events Store as a shared event log. When the Order Service creates an order, the Inventory Service, Shipping Service, and Notification Service each independently consume the event stream and update their local state.
Architecture
One durable orders.lifecycle log fans out to independent services, each tracking its own position via a consumer group; the ungrouped Analytics subscriber sees every event.
Design Decisions
- Single event channel —
orders.lifecycleserves as the shared event log - Consumer groups per service — each service has its own group for independent processing and position tracking
- Durable subscriptions — services that restart catch up from their last position automatically
- Analytics fan-out — an ungrouped subscriber receives every event for reporting
Implementation
Publish Order Lifecycle Events
The Order Service publishes state changes as they occur.
func publishOrderEvent(ctx context.Context, client *kubemq.Client,
orderId, eventType string, data map[string]interface{}) error {
payload := map[string]interface{}{
"orderId": orderId,
"eventType": eventType,
"data": data,
"timestamp": time.Now().UTC().Format(time.RFC3339),
}
body, _ := json.Marshal(payload)
result, err := client.SendEventStore(ctx, kubemq.NewEvent().
SetChannel("orders.lifecycle").
SetMetadata(eventType).
SetBody(body).
SetTags(map[string]string{"orderId": orderId}),
)
if err != nil {
return err
}
log.Printf("Published %s for %s (seq: %s)", eventType, orderId, result.EventID)
return nil
}def publish_order_event(client, order_id, event_type, data):
payload = {
"orderId": order_id,
"eventType": event_type,
"data": data,
"timestamp": datetime.utcnow().isoformat() + "Z",
}
result = client.publish_event_store(
EventStoreMessage(
channel="orders.lifecycle",
metadata=event_type,
body=json.dumps(payload).encode("utf-8"),
tags={"orderId": order_id},
)
)
print(f"Published {event_type} for {order_id} (ID: {result.id})")async function publishOrderEvent(
client: KubeMQClient, orderId: string, eventType: string, data: object
) {
const payload = { orderId, eventType, data, timestamp: new Date().toISOString() };
const result = await client.sendEventStore(
createEventStoreMessage({
channel: 'orders.lifecycle',
metadata: eventType,
body: JSON.stringify(payload),
tags: { orderId },
})
);
console.log(`Published ${eventType} for ${orderId} (ID: ${result.id})`);
}public void publishOrderEvent(PubSubClient client, String orderId,
String eventType, Map<String, Object> data) {
var payload = Map.of(
"orderId", orderId,
"eventType", eventType,
"data", data,
"timestamp", Instant.now().toString());
var body = objectMapper.writeValueAsBytes(payload);
var result = client.sendEventsStoreMessage(EventStoreMessage.builder()
.channel("orders.lifecycle")
.metadata(eventType)
.body(body)
.tags(Map.of("orderId", orderId))
.build());
System.out.printf("Published %s for %s%n", eventType, orderId);
}public async Task PublishOrderEventAsync(KubeMQClient client,
string orderId, string eventType, object data)
{
var payload = new { orderId, eventType, data, timestamp = DateTime.UtcNow };
await client.SendEventStoreAsync(new EventStoreMessage
{
Channel = "orders.lifecycle",
Metadata = eventType,
Body = JsonSerializer.SerializeToUtf8Bytes(payload),
Tags = { ["orderId"] = orderId },
});
Console.WriteLine($"Published {eventType} for {orderId}");
}suspend fun publishOrderEvent(client: KubeMQClient, orderId: String,
eventType: String, data: Map<String, Any>) {
val payload = mapOf("orderId" to orderId, "eventType" to eventType, "data" to data)
client.sendEventStore(eventStoreMessage {
channel = "orders.lifecycle"
metadata = eventType
body = Json.encodeToString(payload).toByteArray()
tags = mapOf("orderId" to orderId)
})
println("Published $eventType for $orderId")
}void publish_order_event(kubemq::Client& client,
const std::string& order_id, const std::string& event_type) {
kubemq::EventStoreMessage msg;
msg.set_channel("orders.lifecycle");
msg.set_metadata(event_type);
msg.set_body("{\"orderId\":\"" + order_id + "\",\"eventType\":\"" + event_type + "\"}");
msg.set_tag("orderId", order_id);
client.SendEventStore(msg);
std::cout << "Published " << event_type << " for " << order_id << std::endl;
}async fn publish_order_event(
client: &KubemqClient,
order_id: &str,
event_type: &str,
data: &str,
) -> kubemq::Result<()> {
let body = format!(
r#"{{"orderId":"{}","eventType":"{}","data":{}}}"#,
order_id, event_type, data
);
let event = EventStoreBuilder::new()
.channel("orders.lifecycle")
.metadata(event_type)
.body(body.into_bytes())
.add_tag("orderId", order_id)
.build();
let result = client.send_event_store(event).await?;
println!(
"Published {} for {} (id: {})",
event_type, order_id, result.id
);
Ok(())
}def publish_order_event(client, order_id, event_type, data)
payload = {
orderId: order_id,
eventType: event_type,
data: data,
timestamp: Time.now.utc.iso8601
}
msg = KubeMQ::PubSub::EventStoreMessage.new(
channel: 'orders.lifecycle',
metadata: event_type,
body: payload.to_json,
tags: { 'orderId' => order_id }
)
result = client.send_event_store(msg)
puts "Published #{event_type} for #{order_id} (sent: #{result.sent})"
enddef publish_order_event(client, order_id, event_type, data) do
payload =
Jason.encode!(%{
orderId: order_id,
eventType: event_type,
data: data,
timestamp: DateTime.utc_now() |> DateTime.to_iso8601()
})
event =
KubeMQ.EventStore.new(
channel: "orders.lifecycle",
metadata: event_type,
body: payload,
tags: %{"orderId" => order_id}
)
{:ok, result} = KubeMQ.Client.send_event_store(client, event)
IO.puts("Published #{event_type} for #{order_id} (sent: #{result.sent})")
endInventory Service Consumer
The Inventory Service subscribes with its own group to reserve or release stock based on order events.
sub, err := client.SubscribeToEventsStore(ctx,
"orders.lifecycle",
"inventory-service",
kubemq.StartFromFirst(),
kubemq.WithOnEvent(func(event *kubemq.Event) {
switch event.Metadata {
case "order.created":
log.Printf("[Inventory] Reserving stock for order %s",
event.Tags["orderId"])
case "order.cancelled":
log.Printf("[Inventory] Releasing stock for order %s",
event.Tags["orderId"])
}
}),
kubemq.WithOnError(func(err error) {
log.Println("[Inventory] Error:", err)
}),
)def on_order_event(event):
order_id = event.tags.get("orderId", "unknown")
if event.metadata == "order.created":
print(f"[Inventory] Reserving stock for {order_id}")
elif event.metadata == "order.cancelled":
print(f"[Inventory] Releasing stock for {order_id}")
client.subscribe_to_events_store(
subscription=EventsStoreSubscription(
channel="orders.lifecycle",
group="inventory-service",
start_position=EventStoreStartPosition.StartFromFirst,
on_receive_event_callback=on_order_event,
on_error_callback=lambda e: print(f"[Inventory] Error: {e}"),
),
cancel=CancellationToken(),
)client.subscribeToEventsStore({
channel: 'orders.lifecycle',
group: 'inventory-service',
startPosition: EventStoreStartPosition.StartFromFirst,
onEvent: (msg) => {
const orderId = msg.tags?.orderId ?? 'unknown';
switch (msg.metadata) {
case 'order.created':
console.log(`[Inventory] Reserving stock for ${orderId}`);
break;
case 'order.cancelled':
console.log(`[Inventory] Releasing stock for ${orderId}`);
break;
}
},
onError: (err) => console.error('[Inventory] Error:', err.message),
});client.subscribeToEventsStore(EventsStoreSubscription.builder()
.channel("orders.lifecycle")
.group("inventory-service")
.startPosition(EventStoreStartPosition.StartFromFirst)
.onReceiveEventCallback(event -> {
String orderId = event.getTags().getOrDefault("orderId", "unknown");
switch (event.getMetadata()) {
case "order.created" ->
System.out.printf("[Inventory] Reserving stock for %s%n", orderId);
case "order.cancelled" ->
System.out.printf("[Inventory] Releasing stock for %s%n", orderId);
}
})
.onErrorCallback(err ->
System.err.println("[Inventory] Error: " + err.getMessage()))
.build());await foreach (var msg in client.SubscribeToEventsStoreAsync(
new EventsStoreSubscription
{
Channel = "orders.lifecycle",
Group = "inventory-service",
StartPosition = EventStoreStartPosition.StartFromFirst,
}))
{
var orderId = msg.Tags.GetValueOrDefault("orderId", "unknown");
switch (msg.Metadata)
{
case "order.created":
Console.WriteLine($"[Inventory] Reserving stock for {orderId}");
break;
case "order.cancelled":
Console.WriteLine($"[Inventory] Releasing stock for {orderId}");
break;
}
}client.subscribeToEventsStore {
channel = "orders.lifecycle"
group = "inventory-service"
startPosition = StartPosition.StartFromFirst
}.collect { msg ->
val orderId = msg.tags["orderId"] ?: "unknown"
when (msg.metadata) {
"order.created" -> println("[Inventory] Reserving stock for $orderId")
"order.cancelled" -> println("[Inventory] Releasing stock for $orderId")
}
}client->SubscribeToEventsStore("orders.lifecycle", "inventory-service",
kubemq::StartPosition::StartFromFirst,
[](const kubemq::EventStoreReceived& msg) {
if (msg.metadata() == "order.created") {
std::cout << "[Inventory] Reserving stock" << std::endl;
} else if (msg.metadata() == "order.cancelled") {
std::cout << "[Inventory] Releasing stock" << std::endl;
}
},
[](const std::string& err) {
std::cerr << "[Inventory] Error: " << err << std::endl;
});let sub = client
.subscribe_to_events_store(
"orders.lifecycle",
"inventory-service",
EventsStoreSubscription::StartFromFirst,
|event| {
Box::pin(async move {
let order_id = event
.tags
.get("orderId")
.map(String::as_str)
.unwrap_or("unknown");
match event.metadata.as_str() {
"order.created" => {
println!("[Inventory] Reserving stock for {}", order_id)
}
"order.cancelled" => {
println!("[Inventory] Releasing stock for {}", order_id)
}
_ => {}
}
})
},
Some(|err| {
Box::pin(async move { eprintln!("[Inventory] Error: {}", err) })
}),
)
.await?;cancel = KubeMQ::CancellationToken.new
sub = KubeMQ::PubSub::EventsStoreSubscription.new(
channel: 'orders.lifecycle',
group: 'inventory-service',
start_position: KubeMQ::PubSub::EventStoreStartPosition::START_FROM_FIRST
)
client.subscribe_to_events_store(sub, cancellation_token: cancel, on_error: lambda { |e|
warn "[Inventory] Error: #{e.message}"
}) do |event|
order_id = event.tags.fetch('orderId', 'unknown')
case event.metadata
when 'order.created'
puts "[Inventory] Reserving stock for #{order_id}"
when 'order.cancelled'
puts "[Inventory] Releasing stock for #{order_id}"
end
end{:ok, sub} =
KubeMQ.Client.subscribe_to_events_store(client, "orders.lifecycle",
start_at: :start_from_first,
group: "inventory-service",
on_event: fn event ->
order_id = Map.get(event.tags, "orderId", "unknown")
case event.metadata do
"order.created" -> IO.puts("[Inventory] Reserving stock for #{order_id}")
"order.cancelled" -> IO.puts("[Inventory] Releasing stock for #{order_id}")
_ -> :ok
end
end,
on_error: fn err -> IO.warn("[Inventory] Error: #{err}") end
)Deploy Multiple Consumers
Each service subscribes independently with its own group. All services can be started, stopped, and scaled independently.
# Each service runs independently and tracks its own position
./inventory-service & # group: inventory-service
./shipping-service & # group: shipping-service
./notification-service & # group: notification-service
./analytics-service & # no group (receives all events)If a service restarts, it automatically catches up from its last processed position.
Production Considerations
Related
- Consumer Groups for scaling within a service
- Resume After Disconnect for durable position tracking
- Audit Trail for compliance logging across services
Was this page helpful?