Struct mz_persist_client::read::Listen
source · pub struct Listen<K: Codec, V: Codec, T, D> {
handle: ReadHandle<K, V, T, D>,
watch: StateWatch<K, V, T, D>,
as_of: Antichain<T>,
since: Antichain<T>,
frontier: Antichain<T>,
}
Expand description
An ongoing subscription of updates to a shard.
Fields§
§handle: ReadHandle<K, V, T, D>
§watch: StateWatch<K, V, T, D>
§as_of: Antichain<T>
§since: Antichain<T>
§frontier: Antichain<T>
Implementations§
source§impl<K, V, T, D> Listen<K, V, T, D>
impl<K, V, T, D> Listen<K, V, T, D>
async fn new(handle: ReadHandle<K, V, T, D>, as_of: Antichain<T>) -> Self
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>
sourceasync fn fetch_batch_part(
&mut self,
part: LeasedBatchPart<T>,
) -> FetchedPart<K, V, T, D> ⓘ
async fn fetch_batch_part( &mut self, part: LeasedBatchPart<T>, ) -> FetchedPart<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 for ReadHandle
to expire the
ReadHandle
held by the 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> Freeze for Listen<K, V, T, D>where
T: Freeze,
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<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
.source§impl<'a, S, T> Semigroup<&'a S> for Twhere
T: Semigroup<S>,
impl<'a, S, T> Semigroup<&'a S> for Twhere
T: Semigroup<S>,
source§fn plus_equals(&mut self, rhs: &&'a S)
fn plus_equals(&mut self, rhs: &&'a S)
std::ops::AddAssign
, for types that do not implement AddAssign
.