Replay from Sequence
Resume events store from a specific sequence number.
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. Setting start_position: KubeMQ::PubSub::EventStoreStartPosition::START_AT_SEQUENCE with start_position_value: 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_position_value: 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 - Ruby SDK installed (
gem install kubemq)
Code
require 'kubemq'
address = ENV.fetch('KUBEMQ_ADDRESS', 'localhost:50000')
channel = 'ruby-events-store.replay-from-sequence'
begin
client = KubeMQ::PubSubClient.new(address: address, client_id: 'es-seq-example')
puts "Connected to #{address}"
5.times do |i|
msg = KubeMQ::PubSub::EventStoreMessage.new(channel: channel, metadata: "event #{i}", body: "data-#{i}")
client.send_event_store(msg)
end
puts 'Pre-stored 5 events'
sleep 1
cancel = KubeMQ::CancellationToken.new
resume_sequence = 3
sub = KubeMQ::PubSub::EventsStoreSubscription.new(
channel: channel,
start_position: KubeMQ::PubSub::EventStoreStartPosition::START_AT_SEQUENCE,
start_position_value: resume_sequence
)
client.subscribe_to_events_store(sub, cancellation_token: cancel, on_error: lambda { |e|
puts "Error: #{e.message}"
}) do |event|
puts "Received (from seq #{resume_sequence}): #{event.metadata}"
end
puts "Subscribed with START_AT_SEQUENCE=#{resume_sequence}"
sleep 3
rescue KubeMQ::Error => e
puts "KubeMQ error: #{e.message}"
ensure
cancel&.cancel
client&.close
puts 'Done'
endHow It Works
START_AT_SEQUENCEreplays events starting from the given sequence number.- Track the last processed sequence to implement exactly-once processing on restart.
- On reconnect, the subscription automatically resumes from the last received sequence.
- Review timeouts, channel names, and client IDs before running against shared environments.
Related
Was this page helpful?