pub struct GlobalStreamManager {
pub env: MetaSrvEnv,
pub metadata_manager: MetadataManager,
pub barrier_scheduler: BarrierScheduler,
pub source_manager: SourceManagerRef,
creating_job_info: Arc<CreatingStreamingJobInfo>,
pub scale_controller: ScaleControllerRef,
}
Expand description
GlobalStreamManager
manages all the streams in the system.
Fields§
§env: MetaSrvEnv
§metadata_manager: MetadataManager
§barrier_scheduler: BarrierScheduler
Broadcasts and collect barriers
source_manager: SourceManagerRef
Maintains streaming sources from external system like kafka
creating_job_info: Arc<CreatingStreamingJobInfo>
Creating streaming job info.
scale_controller: ScaleControllerRef
Implementations§
source§impl GlobalStreamManager
impl GlobalStreamManager
pub async fn reschedule_lock_read_guard(&self) -> RwLockReadGuard<'_, ()>
pub async fn reschedule_lock_write_guard(&self) -> RwLockWriteGuard<'_, ()>
sourcepub async fn reschedule_actors(
&self,
database_id: DatabaseId,
reschedules: HashMap<FragmentId, WorkerReschedule>,
options: RescheduleOptions,
table_parallelism: Option<HashMap<TableId, TableParallelism>>,
) -> MetaResult<()>
pub async fn reschedule_actors( &self, database_id: DatabaseId, reschedules: HashMap<FragmentId, WorkerReschedule>, options: RescheduleOptions, table_parallelism: Option<HashMap<TableId, TableParallelism>>, ) -> MetaResult<()>
The entrypoint of rescheduling actors.
Used by:
- The directly exposed low-level API
risingwave_meta_service::scale_service::ScaleService
(risectl meta reschedule
) - High-level parallelism control API
- manual
ALTER [TABLE | INDEX | MATERIALIZED VIEW | SINK] SET PARALLELISM
- automatic parallelism control for
TableParallelism::Adaptive
when worker nodes changed
- manual
sourceasync fn trigger_parallelism_control(&self) -> MetaResult<bool>
async fn trigger_parallelism_control(&self) -> MetaResult<bool>
When new worker nodes joined, or the parallelism of existing worker nodes changed, examines if there are any jobs can be scaled, and scales them if found.
This method will iterate over all CREATED
jobs, and can be repeatedly called.
Returns
Ok(false)
if no jobs can be scaled;Ok(true)
if some jobs are scaled, and it is possible that there are more jobs can be scaled.
sourceasync fn run(&self, shutdown_rx: Receiver<()>)
async fn run(&self, shutdown_rx: Receiver<()>)
Handles notification of worker node activation and deletion, and triggers parallelism control.
pub fn start_auto_parallelism_monitor( self: Arc<Self>, ) -> (JoinHandle<()>, Sender<()>)
source§impl GlobalStreamManager
impl GlobalStreamManager
pub fn new( env: MetaSrvEnv, metadata_manager: MetadataManager, barrier_scheduler: BarrierScheduler, source_manager: SourceManagerRef, scale_controller: ScaleControllerRef, ) -> MetaResult<Self>
sourcepub async fn create_streaming_job(
self: &Arc<Self>,
table_fragments: TableFragments,
ctx: CreateStreamingJobContext,
) -> MetaResult<NotificationVersion>
pub async fn create_streaming_job( self: &Arc<Self>, table_fragments: TableFragments, ctx: CreateStreamingJobContext, ) -> MetaResult<NotificationVersion>
Create streaming job, it works as follows:
- Broadcast the actor info based on the scheduling result in the context, build the hanging channels in upstream worker nodes.
- (optional) Get the split information of the
StreamSource
via source manager and patch actors. - Notify related worker nodes to update and build the actors.
- Store related meta data.
async fn create_streaming_job_impl( &self, table_fragments: TableFragments, __arg2: CreateStreamingJobContext, ) -> MetaResult<NotificationVersion>
pub async fn replace_table( &self, table_fragments: TableFragments, __arg2: ReplaceTableContext, ) -> MetaResult<()>
sourcepub async fn drop_streaming_jobs(
&self,
database_id: DatabaseId,
removed_actors: Vec<ActorId>,
streaming_job_ids: Vec<ObjectId>,
state_table_ids: Vec<TableId>,
fragment_ids: HashSet<FragmentId>,
)
pub async fn drop_streaming_jobs( &self, database_id: DatabaseId, removed_actors: Vec<ActorId>, streaming_job_ids: Vec<ObjectId>, state_table_ids: Vec<TableId>, fragment_ids: HashSet<FragmentId>, )
Drop streaming jobs by barrier manager, and clean up all related resources. The error will
be ignored because the recovery process will take over it in cleaning part. Check
Command::DropStreamingJobs
for details.
sourcepub async fn cancel_streaming_jobs(
&self,
table_ids: Vec<TableId>,
) -> Vec<TableId>
pub async fn cancel_streaming_jobs( &self, table_ids: Vec<TableId>, ) -> Vec<TableId>
Cancel streaming jobs and return the canceled table ids.
- Send cancel message to stream jobs (via
cancel_jobs
). - Send cancel message to recovered stream jobs (via
barrier_scheduler
).
Cleanup of their state will be cleaned up after the CancelStreamJob
command succeeds,
by the barrier manager for both of them.
pub(crate) async fn alter_table_parallelism( &self, table_id: u32, parallelism: TableParallelism, deferred: bool, ) -> MetaResult<()>
pub async fn create_subscription( self: &Arc<Self>, subscription: &Subscription, ) -> MetaResult<()>
pub async fn drop_subscription( self: &Arc<Self>, database_id: DatabaseId, subscription_id: u32, table_id: u32, )
Auto Trait Implementations§
impl Freeze for GlobalStreamManager
impl !RefUnwindSafe for GlobalStreamManager
impl Send for GlobalStreamManager
impl Sync for GlobalStreamManager
impl Unpin for GlobalStreamManager
impl !UnwindSafe for GlobalStreamManager
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> 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.