Start New Only
Subscribe to only newly published KubeMQ Events Store messages, ignoring history, using StartNewOnly in Rust.
Overview
Start-from-new turns a durable Events Store channel into a live-only feed — reach for it when a consumer only cares what happens from this moment forward and would rather skip a large backlog than pay to replay it. Dashboards, live notification fan-outs, and freshly-deployed services that don't need to catch up on history are the classic cases: any of the replay-from-start variants would mean churning through every historical event just to reach the live tail.
It works by passing EventsStoreSubscription::StartNewOnly to subscribe_to_events_store — the broker stamps the subscription's registration time as a watermark and delivers only events published after it, ignoring everything already stored. Gotchas: there's a race between registering and the publisher sending — a publish that lands before the broker fully registers you is silently skipped, so give the subscription a moment to settle before publishing; this position can never see anything published earlier, so use a start-from-first or start-from-sequence variant when you need guaranteed replay; and reconnecting doesn't resume where you left off — a fresh StartNewOnly subscription starts from "now" again, with no cursor persisted across restarts.
Prerequisites
- KubeMQ server running on
localhost:50000 - Rust SDK installed (
cargo add kubemq)
Code
use kubemq::prelude::*;
use kubemq::EventsStoreSubscription;
use std::time::Duration;
#[tokio::main]
async fn main() -> kubemq::Result<()> {
let client = KubemqClient::builder()
.host("localhost")
.port(50000)
.build()
.await?;
let sub = client.subscribe_to_events_store(
"rust-events-store.start-new-only", "",
EventsStoreSubscription::StartNewOnly,
|event| Box::pin(async move {
println!("New event: seq={}, body={}", event.sequence, String::from_utf8_lossy(&event.body));
}),
None,
).await?;
tokio::time::sleep(Duration::from_secs(10)).await;
sub.unsubscribe().await;
client.close().await?;
Ok(())
}How It Works
StartNewOnlyskips all historical events and only delivers events published after subscribing.- This is the lightest start mode — no replay overhead.
- Review timeouts, channel names, and client IDs before running against shared environments.
- Run the program while the server from the prerequisites is available.
Related
Was this page helpful?