KubeMQ
LearnEvents StoreScenarios

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 channelorders.lifecycle serves 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.

order_service.go
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
}
order_service.py
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})")
order_service.ts
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})`);
}
OrderService.java
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);
}
OrderService.cs
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}");
}
OrderService.kt
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")
}
order_service.cc
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;
}
order_service.rs
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(())
}
order_service.rb
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})"
end
order_service.exs
def 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})")
end

Inventory Service Consumer

The Inventory Service subscribes with its own group to reserve or release stock based on order events.

inventory_service.go
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)
    }),
)
inventory_service.py
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(),
)
inventory_service.ts
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),
});
InventoryService.java
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());
InventoryService.cs
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;
    }
}
InventoryService.kt
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")
    }
}
inventory_service.cc
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;
    });
inventory_service.rs
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?;
inventory_service.rb
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
inventory_service.exs
{: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

Was this page helpful?

On this page