Replay from Sequence
Replay events starting 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 startPosition = StartPosition.StartAtSequence(5) 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 StartAtSequence(5) still delivers event 5 — 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 - Kotlin SDK installed (
implementation("io.kubemq.sdk:kubemq-sdk-kotlin:1.0.1"))
Code
package io.kubemq.sdk.examples.eventsstore
import io.kubemq.sdk.client.KubeMQClient
import io.kubemq.sdk.pubsub.StartPosition
import io.kubemq.sdk.pubsub.eventStoreMessage
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.take
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
private const val ADDRESS = "localhost:50000"
private const val CLIENT_ID = "kotlin-events-store-replay-from-sequence"
private const val CHANNEL = "kotlin-events-store.replay-from-sequence"
fun main() = runBlocking {
val client = KubeMQClient.pubSub {
address = ADDRESS
clientId = CLIENT_ID
}
client.use {
// Publish 10 events
repeat(10) { i ->
val result = client.publishEventStore(eventStoreMessage {
channel = CHANNEL
body = "Event #${i + 1}".toByteArray()
})
println("Published #${i + 1}, sent=${result.sent}")
}
delay(500)
// Subscribe from sequence 5 -- should receive events 5-10
println("\nSubscribing from sequence 5:")
val subJob = launch {
client.subscribeToEventsStore {
channel = CHANNEL
startPosition = StartPosition.StartAtSequence(5)
}.take(6).collect { msg ->
println(" seq=${msg.sequence}: ${String(msg.body)}")
}
}
subJob.join()
println("Done.")
}
}How It Works
StartPosition.StartAtSequence(5)begins replay from sequence number 5.- Events with sequence numbers 5 through 10 are delivered, plus any new events.
- Sequence numbers are monotonically increasing per channel.
- This is useful for resuming processing after a known checkpoint.
Related
Was this page helpful?