Struct mz_persist_client::read::Listen
source · pub struct Listen<K, V, T, D>where
T: Timestamp + Lattice + Codec64,
K: Debug + Codec,
V: Debug + Codec,
D: Semigroup + Codec64 + Send + Sync,{ /* private fields */ }
Expand description
An ongoing subscription of updates to a shard.
Implementations§
source§impl<K, V, T, D> Listen<K, V, T, D>
impl<K, V, T, D> Listen<K, V, T, D>
sourcepub fn frontier(&self) -> &Antichain<T>
pub fn frontier(&self) -> &Antichain<T>
An exclusive upper bound on the progress of this Listen.
sourcepub async fn next(
&mut self,
retry: Option<RetryParameters>
) -> (Vec<LeasedBatchPart<T>>, Antichain<T>)
pub async fn next( &mut self, retry: Option<RetryParameters> ) -> (Vec<LeasedBatchPart<T>>, Antichain<T>)
Attempt to pull out the next values of this subscription.
The returned LeasedBatchPart
is appropriate to use with
crate::fetch::fetch_leased_part
.
The returned Antichain
represents the subscription progress as it will
be after the returned parts are fetched.
source§impl<K, V, T, D> Listen<K, V, T, D>
impl<K, V, T, D> Listen<K, V, T, D>
sourcepub async fn fetch_next(
&mut self
) -> Vec<ListenEvent<T, ((Result<K, String>, Result<V, String>), T, D)>>
pub async fn fetch_next( &mut self ) -> Vec<ListenEvent<T, ((Result<K, String>, Result<V, String>), T, D)>>
Attempt to pull out the next values of this subscription.
The updates received in ListenEvent::Updates should be assumed to be in arbitrary order
and not necessarily consolidated. However, the timestamp of each individual update will be
greater than or equal to the last received ListenEvent::Progress frontier (or this
Listen’s initial as_of
frontier if no progress event has been emitted yet) and less
than the next ListenEvent::Progress frontier.
If you have a use for consolidated listen output, given that snapshots can’t be consolidated, come talk to us!
sourcepub fn into_stream(
self
) -> impl Stream<Item = ListenEvent<T, ((Result<K, String>, Result<V, String>), T, D)>>
pub fn into_stream( self ) -> impl Stream<Item = ListenEvent<T, ((Result<K, String>, Result<V, String>), T, D)>>
Convert listener into futures::Stream
source§impl<K, V, T, D> Listen<K, V, T, D>
impl<K, V, T, D> Listen<K, V, T, D>
sourcepub async fn expire(self)
pub async fn expire(self)
Politely expires this listen, releasing its lease.
There is a best-effort impl in Drop to expire a listen that wasn’t explicitly expired with this method. When possible, explicit expiry is still preferred because the Drop one is best effort and is dependant on a tokio Handle being available in the TLC at the time of drop (which is a bit subtle). Also, explicit expiry allows for control over when it happens.
Trait Implementations§
Auto Trait Implementations§
impl<K, V, T, D> !RefUnwindSafe for Listen<K, V, T, D>
impl<K, V, T, D> Send for Listen<K, V, T, D>
impl<K, V, T, D> Sync for Listen<K, V, T, D>
impl<K, V, T, D> Unpin for Listen<K, V, T, D>where
T: Unpin,
impl<K, V, T, D> !UnwindSafe for Listen<K, V, T, D>
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
source§impl<R, O, T> CopyOnto<ConsecutiveOffsetPairs<R, O>> for T
impl<R, O, T> CopyOnto<ConsecutiveOffsetPairs<R, O>> for T
source§fn copy_onto(
self,
target: &mut ConsecutiveOffsetPairs<R, O>
) -> <ConsecutiveOffsetPairs<R, O> as Region>::Index
fn copy_onto( self, target: &mut ConsecutiveOffsetPairs<R, O> ) -> <ConsecutiveOffsetPairs<R, O> as Region>::Index
source§impl<T> FutureExt for T
impl<T> FutureExt for T
source§fn with_context(self, otel_cx: Context) -> WithContext<Self>
fn with_context(self, otel_cx: Context) -> WithContext<Self>
source§fn with_current_context(self) -> WithContext<Self>
fn with_current_context(self) -> WithContext<Self>
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> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
T
in a tonic::Request
source§impl<T> Pointable for T
impl<T> Pointable for T
source§impl<P, R> ProtoType<R> for Pwhere
R: RustType<P>,
impl<P, R> ProtoType<R> for Pwhere
R: RustType<P>,
source§fn into_rust(self) -> Result<R, TryFromProtoError>
fn into_rust(self) -> Result<R, TryFromProtoError>
RustType::from_proto
.source§fn from_rust(rust: &R) -> P
fn from_rust(rust: &R) -> P
RustType::into_proto
.