pub struct RedisGroupSeeker { /* private fields */ }Expand description
Moves a consumer group’s cursor.
The handle behind the core Seekable capability, minted with
Seekable::seeker or injected into a handler as a
Seek(seeker): Seek<RedisGroupSeeker> parameter.
§A seek is group-wide
Redis keeps one cursor per consumer group, so XGROUP SETID moves the read position for
every consumer of that group, not only for the subscription this seeker came from. That is
unlike a partitioned log, where a seek is scoped to one consumer instance. Treat a seek as an
operation on the group: replaying a range replays it for the whole worker pool, and skipping
forward skips for all of them.
§What a seek does not do
- It does not clear the pending entries list. Entries already delivered and not acknowledged
stay pending and remain reachable through the reclaim path
(
RedisStream::reclaim), whichever way the cursor moved. - It does not cancel delayed retries. Copies already scheduled in a
DelayedRetry::DurableZsetqueue are keyed by their due time, not by the cursor, so they are appended to the stream when they fall due regardless of where the group is reading. - It does not reset delivery counts. A replayed entry is delivered again, so its native
delivery count (the one the reclaim path reads and
RedisStream::max_deliveriescaps) grows with each replay; the framework retry-count header only moves on an actualnack.
§Timing
The cursor changes as soon as seek returns, but a subscription parked in a blocking
XREADGROUP observes it only on its next read - within one
RedisStream::block interval. Entries selected under the old
cursor are discarded rather than delivered, so a seek never yields a message from the position
it moved away from.
§Examples
use ruststream::{Broker, Seekable, Seeker};
use ruststream_fred::{RedisBroker, RedisGroupPosition, RedisStream};
let connected = RedisBroker::standalone("redis://localhost:6379").connect().await?;
let subscriber = connected
.subscribe(RedisStream::new("orders").group("workers"))
.await?;
// Minted before the stream opens; usable while it runs.
let seeker = subscriber.seeker();
seeker.seek(RedisGroupPosition::beginning()).await?;Trait Implementations§
Source§impl Clone for RedisGroupSeeker
impl Clone for RedisGroupSeeker
Source§fn clone(&self) -> RedisGroupSeeker
fn clone(&self) -> RedisGroupSeeker
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreSource§impl Debug for RedisGroupSeeker
impl Debug for RedisGroupSeeker
Source§impl Seeker for RedisGroupSeeker
impl Seeker for RedisGroupSeeker
Source§async fn seek(&self, to: RedisGroupPosition) -> Result<(), RedisError>
async fn seek(&self, to: RedisGroupPosition) -> Result<(), RedisError>
Moves the group cursor with XGROUP SETID.
§Errors
Returns RedisError::Stream when the group or the stream does not exist, or the
command fails.
Source§type Position = RedisGroupPosition
type Position = RedisGroupPosition
Kafka partition offsets, a Redis entry id, a byte
offset), constructed by the broker crate or captured from a delivered message via
Positioned::position.Source§type Error = RedisError
type Error = RedisError
Auto Trait Implementations§
impl !RefUnwindSafe for RedisGroupSeeker
impl !UnwindSafe for RedisGroupSeeker
impl Freeze for RedisGroupSeeker
impl Send for RedisGroupSeeker
impl Sync for RedisGroupSeeker
impl Unpin for RedisGroupSeeker
impl UnsafeUnpin for RedisGroupSeeker
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
fn into_either(self, into_left: bool) -> Either<Self, Self>
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more