pub(super) struct LocalBarrierWorker {
pub(super) state: ManagedBarrierState,
await_epoch_completed_futures: FuturesOrdered<impl Future<Output = (PartialGraphId, Barrier, StreamResult<BarrierCompleteResult>)> + 'static>,
control_stream_handle: ControlStreamHandle,
pub(super) actor_manager: Arc<StreamActorManager>,
pub(super) current_shared_context: Arc<SharedContext>,
barrier_event_rx: UnboundedReceiver<LocalBarrierEvent>,
actor_failure_rx: UnboundedReceiver<(ActorId, StreamError)>,
}
Expand description
LocalBarrierWorker
manages barrier control flow, used by local stream manager.
Specifically, LocalBarrierWorker
serve barrier injection from meta server, send the
barriers to and collect them from all actors, and finally report the progress.
Fields§
§state: ManagedBarrierState
Current barrier collection state.
await_epoch_completed_futures: FuturesOrdered<impl Future<Output = (PartialGraphId, Barrier, StreamResult<BarrierCompleteResult>)> + 'static>
Futures will be finished in the order of epoch in ascending order.
control_stream_handle: ControlStreamHandle
§actor_manager: Arc<StreamActorManager>
§barrier_event_rx: UnboundedReceiver<LocalBarrierEvent>
§actor_failure_rx: UnboundedReceiver<(ActorId, StreamError)>
Implementations§
source§impl LocalBarrierWorker
impl LocalBarrierWorker
pub(crate) fn update_create_mview_progress( &mut self, epoch: EpochPair, actor: ActorId, state: BackfillState, )
source§impl LocalBarrierWorker
impl LocalBarrierWorker
pub(super) fn new( actor_manager: Arc<StreamActorManager>, initial_partial_graphs: Vec<InitialPartialGraph>, ) -> Self
fn to_debug_info(&self) -> LocalBarrierWorkerDebugInfo<'_>
async fn run(self, actor_op_rx: UnboundedReceiver<LocalActorOperation>)
fn handle_streaming_control_request( &mut self, request: StreamingControlStreamRequest, ) -> StreamResult<()>
fn handle_barrier_event( &mut self, event: LocalBarrierEvent, ) -> Result<(), (ActorId, StreamError)>
fn handle_actor_op(&mut self, actor_op: LocalActorOperation)
source§impl LocalBarrierWorker
impl LocalBarrierWorker
fn complete_barrier( &mut self, partial_graph_id: PartialGraphId, prev_epoch: u64, )
fn on_epoch_completed( &mut self, partial_graph_id: PartialGraphId, epoch: u64, result: BarrierCompleteResult, )
sourcefn send_barrier(
&mut self,
barrier: &Barrier,
request: InjectBarrierRequest,
) -> StreamResult<()>
fn send_barrier( &mut self, barrier: &Barrier, request: InjectBarrierRequest, ) -> StreamResult<()>
Broadcast a barrier to all senders. Save a receiver which will get notified when this barrier is finished, in managed mode.
Note that the error returned here is typically a StreamError::barrier_send
, which is not
the root cause of the failure. The caller should then call Self::try_find_root_failure
to find the root cause.
fn remove_partial_graphs( &mut self, partial_graph_ids: impl Iterator<Item = PartialGraphId>, )
pub(super) fn add_partial_graph(&mut self, partial_graph_id: PartialGraphId)
sourcepub(super) fn reset_state(
&mut self,
initial_partial_graphs: Vec<InitialPartialGraph>,
)
pub(super) fn reset_state( &mut self, initial_partial_graphs: Vec<InitialPartialGraph>, )
Reset all internal states.
sourcefn collect(&mut self, actor_id: ActorId, epoch: EpochPair)
fn collect(&mut self, actor_id: ActorId, epoch: EpochPair)
When a crate::executor::StreamConsumer
(typically crate::executor::DispatchExecutor
) get a barrier, it should report
and collect this barrier with its own actor_id
using this function.
sourceasync fn notify_actor_failure(
&mut self,
actor_id: ActorId,
err: StreamError,
err_context: &'static str,
)
async fn notify_actor_failure( &mut self, actor_id: ActorId, err: StreamError, err_context: &'static str, )
When a actor exit unexpectedly, the error is reported using this function. The control stream will be reset and the meta service will then trigger recovery.
sourceasync fn notify_other_failure(
&mut self,
err: StreamError,
message: impl Into<String>,
)
async fn notify_other_failure( &mut self, err: StreamError, message: impl Into<String>, )
When some other failure happens (like failed to send barrier), the error is reported using this function. The control stream will be reset and the meta service will then trigger recovery.
This is similar to Self::notify_actor_failure
, but since there’s not always an actor failure,
the given err
will be used if there’s no root failure found.
sourceasync fn try_find_root_failure(
&mut self,
first_err: StreamError,
) -> ScoredStreamError
async fn try_find_root_failure( &mut self, first_err: StreamError, ) -> ScoredStreamError
Collect actor errors for a while and find the one that might be the root cause.
Returns None
if there’s no actor error received.
source§impl LocalBarrierWorker
impl LocalBarrierWorker
sourcepub fn spawn(
env: StreamEnvironment,
streaming_metrics: Arc<StreamingMetrics>,
await_tree_reg: Option<Registry>,
watermark_epoch: AtomicU64Ref,
actor_op_rx: UnboundedReceiver<LocalActorOperation>,
) -> JoinHandle<()> ⓘ
pub fn spawn( env: StreamEnvironment, streaming_metrics: Arc<StreamingMetrics>, await_tree_reg: Option<Registry>, watermark_epoch: AtomicU64Ref, actor_op_rx: UnboundedReceiver<LocalActorOperation>, ) -> JoinHandle<()> ⓘ
Create a LocalBarrierWorker
with managed mode.
source§impl LocalBarrierWorker
impl LocalBarrierWorker
sourcepub(super) async fn reset(
&mut self,
initial_partial_graphs: Vec<InitialPartialGraph>,
)
pub(super) async fn reset( &mut self, initial_partial_graphs: Vec<InitialPartialGraph>, )
Force stop all actors on this worker, and then drop their resources.
source§impl LocalBarrierWorker
impl LocalBarrierWorker
sourcepub fn update_actor_info(
&self,
new_actor_infos: impl Iterator<Item = ActorInfo>,
) -> StreamResult<()>
pub fn update_actor_info( &self, new_actor_infos: impl Iterator<Item = ActorInfo>, ) -> StreamResult<()>
This function could only be called once during the lifecycle of LocalStreamManager
for
now.
Auto Trait Implementations§
impl !Freeze for LocalBarrierWorker
impl !RefUnwindSafe for LocalBarrierWorker
impl Send for LocalBarrierWorker
impl !Sync for LocalBarrierWorker
impl Unpin for LocalBarrierWorker
impl !UnwindSafe for LocalBarrierWorker
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<T> Conv for T
impl<T> Conv for T
§impl<Choices> CoproductSubsetter<CNil, HNil> for Choices
impl<Choices> CoproductSubsetter<CNil, HNil> for Choices
§impl<T> Downcast for Twhere
T: Any,
impl<T> Downcast for Twhere
T: Any,
§fn into_any(self: Box<T>) -> Box<dyn Any>
fn into_any(self: Box<T>) -> Box<dyn Any>
Box<dyn Trait>
(where Trait: Downcast
) to Box<dyn Any>
. Box<dyn Any>
can
then be further downcast
into Box<ConcreteType>
where ConcreteType
implements Trait
.§fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>
fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>
Rc<Trait>
(where Trait: Downcast
) to Rc<Any>
. Rc<Any>
can then be
further downcast
into Rc<ConcreteType>
where ConcreteType
implements Trait
.§fn as_any(&self) -> &(dyn Any + 'static)
fn as_any(&self) -> &(dyn Any + 'static)
&Trait
(where Trait: Downcast
) to &Any
. This is needed since Rust cannot
generate &Any
’s vtable from &Trait
’s.§fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)
fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)
&mut Trait
(where Trait: Downcast
) to &Any
. This is needed since Rust cannot
generate &mut Any
’s vtable from &mut Trait
’s.§impl<T> FmtForward for T
impl<T> FmtForward for T
§fn fmt_binary(self) -> FmtBinary<Self>where
Self: Binary,
fn fmt_binary(self) -> FmtBinary<Self>where
Self: Binary,
self
to use its Binary
implementation when Debug
-formatted.§fn fmt_display(self) -> FmtDisplay<Self>where
Self: Display,
fn fmt_display(self) -> FmtDisplay<Self>where
Self: Display,
self
to use its Display
implementation when
Debug
-formatted.§fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
Self: LowerExp,
fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
Self: LowerExp,
self
to use its LowerExp
implementation when
Debug
-formatted.§fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
Self: LowerHex,
fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
Self: LowerHex,
self
to use its LowerHex
implementation when
Debug
-formatted.§fn fmt_octal(self) -> FmtOctal<Self>where
Self: Octal,
fn fmt_octal(self) -> FmtOctal<Self>where
Self: Octal,
self
to use its Octal
implementation when Debug
-formatted.§fn fmt_pointer(self) -> FmtPointer<Self>where
Self: Pointer,
fn fmt_pointer(self) -> FmtPointer<Self>where
Self: Pointer,
self
to use its Pointer
implementation when
Debug
-formatted.§fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
Self: UpperExp,
fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
Self: UpperExp,
self
to use its UpperExp
implementation when
Debug
-formatted.§fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
Self: UpperHex,
fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
Self: UpperHex,
self
to use its UpperHex
implementation when
Debug
-formatted.§fn fmt_list(self) -> FmtList<Self>where
&'a Self: for<'a> IntoIterator,
fn fmt_list(self) -> FmtList<Self>where
&'a Self: for<'a> IntoIterator,
§impl<T> FutureExt for T
impl<T> FutureExt for T
§fn with_context(self, otel_cx: Context) -> WithContext<Self>
fn with_context(self, otel_cx: Context) -> WithContext<Self>
§fn with_current_context(self) -> WithContext<Self>
fn with_current_context(self) -> WithContext<Self>
§impl<T> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<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> 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 moresource§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
§impl<T> IntoResult<T> for T
impl<T> IntoResult<T> for T
type Err = Infallible
fn into_result(self) -> Result<T, <T as IntoResult<T>>::Err>
§impl<T, U, I> LiftInto<U, I> for Twhere
U: LiftFrom<T, I>,
impl<T, U, I> LiftInto<U, I> for Twhere
U: LiftFrom<T, I>,
source§impl<M> MetricVecRelabelExt for M
impl<M> MetricVecRelabelExt for M
source§fn relabel(
self,
metric_level: MetricLevel,
relabel_threshold: MetricLevel,
) -> RelabeledMetricVec<M>
fn relabel( self, metric_level: MetricLevel, relabel_threshold: MetricLevel, ) -> RelabeledMetricVec<M>
RelabeledMetricVec::with_metric_level
.source§fn relabel_n(
self,
metric_level: MetricLevel,
relabel_threshold: MetricLevel,
relabel_num: usize,
) -> RelabeledMetricVec<M>
fn relabel_n( self, metric_level: MetricLevel, relabel_threshold: MetricLevel, relabel_num: usize, ) -> RelabeledMetricVec<M>
RelabeledMetricVec::with_metric_level_relabel_n
.source§fn relabel_debug_1(
self,
relabel_threshold: MetricLevel,
) -> RelabeledMetricVec<M>
fn relabel_debug_1( self, relabel_threshold: MetricLevel, ) -> RelabeledMetricVec<M>
RelabeledMetricVec::with_metric_level_relabel_n
with metric_level
set to
MetricLevel::Debug
and relabel_num
set to 1.§impl<T> Pipe for Twhere
T: ?Sized,
impl<T> Pipe for Twhere
T: ?Sized,
§fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
Self: Sized,
fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
Self: Sized,
§fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
R: 'a,
fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
R: 'a,
self
and passes that borrow into the pipe function. Read more§fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
R: 'a,
fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
R: 'a,
self
and passes that borrow into the pipe function. Read more§fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
§fn pipe_borrow_mut<'a, B, R>(
&'a mut self,
func: impl FnOnce(&'a mut B) -> R,
) -> R
fn pipe_borrow_mut<'a, B, R>( &'a mut self, func: impl FnOnce(&'a mut B) -> R, ) -> R
§fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
self
, then passes self.as_ref()
into the pipe function.§fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
self
, then passes self.as_mut()
into the pipe
function.§fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
self
, then passes self.deref()
into the pipe function.§impl<T> Pointable for T
impl<T> Pointable for T
§impl<Source> Sculptor<HNil, HNil> for Source
impl<Source> Sculptor<HNil, HNil> for Source
§impl<T> Tap for T
impl<T> Tap for T
§fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
Borrow<B>
of a value. Read more§fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
BorrowMut<B>
of a value. Read more§fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
AsRef<R>
view of a value. Read more§fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
AsMut<R>
view of a value. Read more§fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
Deref::Target
of a value. Read more§fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
Deref::Target
of a value. Read more§fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
.tap()
only in debug builds, and is erased in release builds.§fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
.tap_mut()
only in debug builds, and is erased in release
builds.§fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
.tap_borrow()
only in debug builds, and is erased in release
builds.§fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
.tap_borrow_mut()
only in debug builds, and is erased in release
builds.§fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
.tap_ref()
only in debug builds, and is erased in release
builds.§fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
.tap_ref_mut()
only in debug builds, and is erased in release
builds.§fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
.tap_deref()
only in debug builds, and is erased in release
builds.