KubeMQ
Client SDKsElixirHow-to guidesEvents Store

Replay from Sequence

Replay KubeMQ Events Store messages starting from a specific sequence number using the Elixir SDK.

Overview

Replaying from a sequence number lets a consumer resume an events-store subscription from an exact point in a channel's history, instead of re-reading everything or only catching new traffic. It's the checkpoint-recovery pattern: a worker persists the last sequence it processed, and after a crash or redeploy it reopens the subscription right there — no gap, no reprocessing everything that came before.

Sequence numbers are broker-assigned per channel, starting at 1 and increasing monotonically with every stored event; they never reset unless the channel is purged. Passing start_at: {:start_at_sequence, 3} tells the broker to begin delivery at that sequence inclusive, replaying stored events from that point, then transitioning the subscription to live delivery for anything published afterward.

Gotchas: the sequence value is inclusive, so {:start_at_sequence, 3} still delivers event 3 — off by one and you'll reprocess or silently drop a message; you must track and persist the "last processed" sequence yourself, KubeMQ doesn't checkpoint it for you; and requesting a sequence past the current head isn't an error — you'll just get nothing until new events catch up to it.

Prerequisites

  • KubeMQ server running on localhost:50000
  • Elixir SDK installed ({:kubemq, "~> 1.0"} in mix.exs)

Code

main.exs
channel = "elixir-events-store.replay-from-sequence"
{:ok, client} = KubeMQ.Client.start_link(address: "localhost:50000", client_id: "elixir-es-seq")

for i <- 1..5 do
  {:ok, _} = KubeMQ.Client.send_event_store(client,
    KubeMQ.EventStore.new(channel: channel, body: "Event #{i}"))
end

IO.puts("Stored 5 events. Subscribing from sequence 3...")
parent = self()

{:ok, sub} =
  KubeMQ.Client.subscribe_to_events_store(client, channel,
    start_at: {:start_at_sequence, 3},
    on_event: fn event ->
      IO.puts("Replayed seq #{event.sequence}: #{event.body}")
      send(parent, :replayed)
    end
  )

for _ <- 1..3 do
  receive do
    :replayed -> :ok
  after
    5_000 -> :timeout
  end
end

IO.puts("Replayed from sequence 3 onwards")
KubeMQ.Subscription.cancel(sub)
KubeMQ.Client.close(client)

How It Works

  • {:start_at_sequence, 3} tells the server to start delivering from sequence number 3
  • Events 1 and 2 are skipped, and events 3, 4, 5 are replayed
  • After replay, new events published to the channel are also delivered

Was this page helpful?

On this page