Replay from Time
Replay KubeMQ events store messages from a specific timestamp with StartAtTime using the Kotlin SDK.
Overview
Replaying from a timestamp lets a consumer recover a window of history without knowing exact sequence numbers — you reach for it after a deploy, an outage, or any gap where you know roughly when you went dark but not where you left off in the stream. It turns an Events Store channel into a rewindable log: resubscribe with a point in time and the broker replays every event stored at or after it, then hands off to live delivery.
The subscription's startPosition is set to StartPosition.StartAtTime(nanos), given nanoseconds since epoch derived from an Instant — the broker compares this against the storage timestamp it assigned to each event, not any timestamp embedded in the payload. Because it's wall-clock based, the window is approximate rather than exact: pass a time far enough back to be safe.
Gotchas: clock skew between your subscriber's clock and the server's matters — favor a generous buffer over a precise cutoff. Storage timestamps reflect when the broker persisted the event, not when the producer created it, so under load the two can drift. And unlike sequence-based replay, a time-based start position has no way to guarantee "no gaps, no duplicates" across a network hiccup — use StartPosition.StartAtSequence instead if you need exact resumption.
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
import java.time.Instant
private const val ADDRESS = "localhost:50000"
private const val CLIENT_ID = "kotlin-events-store-replay-from-time"
private const val CHANNEL = "kotlin-events-store.replay-from-time"
fun main() = runBlocking {
val client = KubeMQClient.pubSub {
address = ADDRESS
clientId = CLIENT_ID
}
client.use {
// Record the timestamp before publishing
val beforePublish = Instant.now()
// Publish 5 events
repeat(5) { i ->
val result = client.publishEventStore(eventStoreMessage {
channel = CHANNEL
body = "Time event #${i + 1}".toByteArray()
})
println("Published #${i + 1}, sent=${result.sent}")
}
delay(500)
// Subscribe from the recorded timestamp (nanos since epoch)
val timestampNanos = beforePublish.epochSecond * 1_000_000_000L + beforePublish.nano
println("\nSubscribing from time ${beforePublish}:")
val subJob = launch {
client.subscribeToEventsStore {
channel = CHANNEL
startPosition = StartPosition.StartAtTime(timestampNanos)
}.take(5).collect { msg ->
println(" seq=${msg.sequence}: ${String(msg.body)}")
}
}
subJob.join()
println("Done.")
}
}How It Works
StartPosition.StartAtTime(nanos)replays events from the specified timestamp (nanoseconds since epoch).- Use
java.time.Instantto capture timestamps and convert to nanoseconds. - All events published after the timestamp are delivered, plus any new events.
- This is useful for replaying events from a known point in time (e.g., after an outage).
Related
Was this page helpful?