1mod prelude;
16
17use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
18use std::fmt::Debug;
19use std::future::pending;
20use std::hash::Hash;
21use std::pin::Pin;
22use std::sync::Arc;
23use std::task::Poll;
24use std::vec;
25
26use await_tree::InstrumentAwait;
27use enum_as_inner::EnumAsInner;
28use futures::future::try_join_all;
29use futures::stream::{BoxStream, FusedStream, FuturesUnordered, StreamFuture};
30use futures::{FutureExt, Stream, StreamExt, TryStreamExt};
31use itertools::Itertools;
32use prometheus::core::{AtomicU64, GenericCounter};
33use risingwave_common::array::StreamChunk;
34use risingwave_common::bitmap::Bitmap;
35use risingwave_common::catalog::{Schema, TableId};
36use risingwave_common::config::StreamingConfig;
37use risingwave_common::metrics::LabelGuardedMetric;
38use risingwave_common::row::OwnedRow;
39use risingwave_common::types::{DataType, Datum, DefaultOrd, ScalarImpl};
40use risingwave_common::util::epoch::{Epoch, EpochPair};
41use risingwave_common::util::tracing::TracingContext;
42use risingwave_common::util::value_encoding::{DatumFromProtoExt, DatumToProtoExt};
43use risingwave_common_estimate_size::EstimateSize;
44use risingwave_connector::source::SplitImpl;
45use risingwave_expr::expr::NonStrictExpression;
46use risingwave_pb::common::ThrottleType;
47use risingwave_pb::data::PbEpoch;
48use risingwave_pb::expr::PbInputRef;
49use risingwave_pb::stream_plan::add_mutation::PbNewUpstreamSink;
50use risingwave_pb::stream_plan::barrier::BarrierKind;
51use risingwave_pb::stream_plan::barrier_mutation::Mutation as PbMutation;
52use risingwave_pb::stream_plan::stream_node::PbStreamKind;
53use risingwave_pb::stream_plan::throttle_mutation::ThrottleConfig;
54use risingwave_pb::stream_plan::update_mutation::{DispatcherUpdate, MergeUpdate};
55use risingwave_pb::stream_plan::{
56 IcebergPkIndexCompactionContext, PbBarrier, PbBarrierMutation, PbDispatcher,
57 PbSinkSchemaChange, PbStreamMessageBatch, PbWatermark, SubscriptionUpstreamInfo,
58};
59use smallvec::SmallVec;
60use tokio::sync::mpsc;
61use tokio::time::{Duration, Instant};
62
63use crate::error::StreamResult;
64use crate::executor::exchange::input::{
65 BoxedActorInput, BoxedInput, assert_equal_dispatcher_barrier, new_input,
66};
67use crate::executor::monitor::ActorInputMetrics;
68use crate::executor::prelude::StreamingMetrics;
69use crate::executor::watermark::BufferedWatermarks;
70use crate::task::{ActorId, FragmentId, LocalBarrierManager};
71
72mod actor;
73mod barrier_align;
74pub mod exchange;
75pub mod monitor;
76
77pub mod aggregate;
78pub mod asof_join;
79mod backfill;
80mod barrier_recv;
81mod batch_query;
82mod chain;
83mod changelog;
84mod dedup;
85mod dispatch;
86pub mod dml;
87mod dynamic_filter;
88pub mod eowc;
89pub mod error;
90mod expand;
91mod filter;
92mod gap_fill;
93pub mod hash_join;
94mod hop_window;
95pub(crate) mod iceberg_with_pk_index;
96mod join;
97pub mod locality_provider;
98mod lookup;
99mod lookup_union;
100pub mod match_recognize;
101mod merge;
102mod mview;
103mod nested_loop_temporal_join;
104mod no_op;
105mod now;
106mod over_window;
107pub mod project;
108mod receiver;
109pub mod row_id_gen;
110mod sink;
111pub mod source;
112mod stream_reader;
113pub mod subtask;
114mod temporal_join;
115mod top_n;
116mod troublemaker;
117mod union;
118mod upstream_sink_union;
119mod values;
120mod watermark;
121mod watermark_filter;
122mod wrapper;
123
124mod approx_percentile;
125
126mod row_merge;
127
128#[cfg(test)]
129mod integration_tests;
130mod sync_kv_log_store;
131#[cfg(any(test, feature = "test"))]
132pub mod test_utils;
133mod utils;
134mod vector;
135
136pub use actor::{Actor, ActorContext, ActorContextRef};
137use anyhow::{Context, anyhow};
138pub use approx_percentile::global::GlobalApproxPercentileExecutor;
139pub use approx_percentile::local::LocalApproxPercentileExecutor;
140pub use backfill::arrangement_backfill::*;
141pub use backfill::cdc::{
142 CdcBackfillExecutor, ExternalStorageTable, ParallelizedCdcBackfillExecutor,
143};
144pub use backfill::no_shuffle_backfill::*;
145pub use backfill::snapshot_backfill::*;
146pub use barrier_recv::BarrierRecvExecutor;
147pub use batch_query::BatchQueryExecutor;
148pub use chain::ChainExecutor;
149pub use changelog::{ChangeLogExecutor, ChangeLogMode};
150pub use dedup::AppendOnlyDedupExecutor;
151pub use dispatch::{DispatchExecutor, SyncLogStoreDispatchExecutor};
152pub use dynamic_filter::DynamicFilterExecutor;
153pub use error::{StreamExecutorError, StreamExecutorResult};
154pub use expand::ExpandExecutor;
155pub use filter::{FilterExecutor, UpsertFilterExecutor};
156pub use gap_fill::{GapFillExecutor, GapFillExecutorArgs};
157pub use hash_join::*;
158pub use hop_window::HopWindowExecutor;
159pub use iceberg_with_pk_index::{
160 CompactionResolverExecutor, IcebergWriterImpl, PositionDeleteHandlerImpl,
161 PositionDeleteMergerExecutor, WriterExecutor,
162};
163pub use join::asof_join::{AsOfCpuEncoding, AsOfMemoryEncoding};
164pub use join::row::{CachedJoinRow, CpuEncoding, JoinEncoding, MemoryEncoding};
165pub use join::{AsOfDesc, AsOfJoinType, JoinType};
166pub use lookup::*;
167pub use lookup_union::LookupUnionExecutor;
168pub use merge::MergeExecutor;
169pub(crate) use merge::{MergeExecutorInput, MergeExecutorUpstream};
170pub use mview::{MaterializeExecutor, RefreshableMaterializeArgs};
171pub use nested_loop_temporal_join::NestedLoopTemporalJoinExecutor;
172pub use no_op::NoOpExecutor;
173pub use now::*;
174pub use over_window::*;
175pub use receiver::ReceiverExecutor;
176use risingwave_common::id::SourceId;
177pub use row_merge::RowMergeExecutor;
178pub use sink::SinkExecutor;
179pub use sync_kv_log_store::SyncedKvLogStoreExecutor;
180pub use sync_kv_log_store::metrics::SyncedKvLogStoreMetrics;
181pub use temporal_join::TemporalJoinExecutor;
182pub use top_n::{
183 AppendOnlyGroupTopNExecutor, AppendOnlyTopNExecutor, GroupTopNExecutor, TopNExecutor,
184};
185pub use troublemaker::TroublemakerExecutor;
186pub use union::UnionExecutor;
187pub use upstream_sink_union::{UpstreamFragmentInfo, UpstreamSinkUnionExecutor};
188pub use utils::DummyExecutor;
189pub use values::ValuesExecutor;
190pub use vector::*;
191pub use watermark_filter::{UpsertWatermarkFilterExecutor, WatermarkFilterExecutor};
192pub use wrapper::WrapperExecutor;
193
194use self::barrier_align::AlignedMessageStream;
195
196pub type MessageStreamItemInner<M> = StreamExecutorResult<MessageInner<M>>;
197pub type MessageStreamItem = MessageStreamItemInner<BarrierMutationType>;
198pub type DispatcherMessageStreamItem = StreamExecutorResult<DispatcherMessage>;
199pub type BoxedMessageStream = BoxStream<'static, MessageStreamItem>;
200
201pub use risingwave_common::util::epoch::task_local::{curr_epoch, epoch, prev_epoch};
202use risingwave_connector::sink::catalog::SinkId;
203use risingwave_connector::source::cdc::{
204 CdcTableSnapshotSplitAssignmentWithGeneration,
205 build_actor_cdc_table_snapshot_splits_with_generation,
206};
207use risingwave_pb::id::{ExecutorId, SubscriberId};
208use risingwave_pb::stream_plan::stream_message_batch::{BarrierBatch, StreamMessageBatch};
209
210pub trait MessageStreamInner<M> = Stream<Item = MessageStreamItemInner<M>> + Send;
211pub trait MessageStream = Stream<Item = MessageStreamItem> + Send;
212pub trait DispatcherMessageStream = Stream<Item = DispatcherMessageStreamItem> + Send;
213
214#[derive(Debug, Default, Clone)]
216pub struct ExecutorInfo {
217 pub schema: Schema,
219
220 pub stream_key: StreamKey,
222
223 pub stream_kind: PbStreamKind,
225
226 pub identity: String,
228
229 pub id: ExecutorId,
231}
232
233impl ExecutorInfo {
234 pub fn for_test(schema: Schema, stream_key: StreamKey, identity: String, id: u64) -> Self {
235 Self {
236 schema,
237 stream_key,
238 stream_kind: PbStreamKind::Retract, identity,
240 id: id.into(),
241 }
242 }
243}
244
245pub trait Execute: Send + 'static {
247 fn execute(self: Box<Self>) -> BoxedMessageStream;
248
249 fn boxed(self) -> Box<dyn Execute>
250 where
251 Self: Sized + Send + 'static,
252 {
253 Box::new(self)
254 }
255}
256
257pub struct Executor {
260 info: ExecutorInfo,
261 execute: Box<dyn Execute>,
262}
263
264impl Executor {
265 pub fn new(info: ExecutorInfo, execute: Box<dyn Execute>) -> Self {
266 Self { info, execute }
267 }
268
269 pub fn info(&self) -> &ExecutorInfo {
270 &self.info
271 }
272
273 pub fn schema(&self) -> &Schema {
274 &self.info.schema
275 }
276
277 pub fn stream_key(&self) -> StreamKeyRef<'_> {
278 &self.info.stream_key
279 }
280
281 pub fn stream_kind(&self) -> PbStreamKind {
282 self.info.stream_kind
283 }
284
285 pub fn identity(&self) -> &str {
286 &self.info.identity
287 }
288
289 pub fn execute(self) -> BoxedMessageStream {
290 self.execute.execute()
291 }
292}
293
294impl std::fmt::Debug for Executor {
295 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
296 f.write_str(self.identity())
297 }
298}
299
300impl From<(ExecutorInfo, Box<dyn Execute>)> for Executor {
301 fn from((info, execute): (ExecutorInfo, Box<dyn Execute>)) -> Self {
302 Self::new(info, execute)
303 }
304}
305
306impl<E> From<(ExecutorInfo, E)> for Executor
307where
308 E: Execute,
309{
310 fn from((info, execute): (ExecutorInfo, E)) -> Self {
311 Self::new(info, execute.boxed())
312 }
313}
314
315pub const INVALID_EPOCH: u64 = 0;
316
317type UpstreamFragmentId = FragmentId;
318type SplitAssignments = HashMap<ActorId, Vec<SplitImpl>>;
319
320#[derive(Debug, Clone, PartialEq)]
321#[cfg_attr(any(test, feature = "test"), derive(Default))]
322pub struct UpdateMutation {
323 pub dispatchers: HashMap<ActorId, Vec<DispatcherUpdate>>,
324 pub merges: HashMap<(ActorId, UpstreamFragmentId), MergeUpdate>,
325 pub vnode_bitmaps: HashMap<ActorId, Arc<Bitmap>>,
326 pub dropped_actors: HashSet<ActorId>,
327 pub actor_splits: SplitAssignments,
328 pub actor_new_dispatchers: HashMap<ActorId, Vec<PbDispatcher>>,
329 pub actor_cdc_table_snapshot_splits: CdcTableSnapshotSplitAssignmentWithGeneration,
330 pub sink_schema_change: HashMap<SinkId, PbSinkSchemaChange>,
331 pub subscriptions_to_drop: Vec<SubscriptionUpstreamInfo>,
332 pub iceberg_pk_index_compaction: Option<IcebergPkIndexCompactionContext>,
333}
334
335#[derive(Debug, Clone, PartialEq)]
336#[cfg_attr(any(test, feature = "test"), derive(Default))]
337pub struct AddMutation {
338 pub adds: HashMap<ActorId, Vec<PbDispatcher>>,
339 pub added_actors: HashSet<ActorId>,
340 pub dropped_actors: HashSet<ActorId>,
341 pub splits: SplitAssignments,
343 pub pause: bool,
344 pub subscriptions_to_add: Vec<(TableId, SubscriberId)>,
346 pub backfill_nodes_to_pause: HashSet<FragmentId>,
348 pub actor_cdc_table_snapshot_splits: CdcTableSnapshotSplitAssignmentWithGeneration,
349 pub new_upstream_sinks: HashMap<FragmentId, PbNewUpstreamSink>,
350 pub sink_log_store_flush: HashSet<SinkId>,
351}
352
353#[derive(Debug, Clone, PartialEq)]
354#[cfg_attr(any(test, feature = "test"), derive(Default))]
355pub struct StopMutation {
356 pub dropped_actors: HashSet<ActorId>,
357 pub dropped_sink_fragments: HashSet<FragmentId>,
358}
359
360#[derive(Debug, Clone, PartialEq)]
362pub enum Mutation {
363 Stop(StopMutation),
364 Update(UpdateMutation),
365 Add(AddMutation),
366 SourceChangeSplit(SplitAssignments),
367 Pause,
368 Resume,
369 Throttle(HashMap<FragmentId, ThrottleConfig>),
370 ConnectorPropsChange(HashMap<u32, HashMap<String, String>>),
371 DropSubscriptions {
372 subscriptions_to_drop: Vec<SubscriptionUpstreamInfo>,
374 },
375 StartFragmentBackfill {
376 fragment_ids: HashSet<FragmentId>,
377 },
378 RefreshStart {
379 table_id: TableId,
380 associated_source_id: SourceId,
381 },
382 ListFinish {
383 associated_source_id: SourceId,
384 },
385 LoadFinish {
386 associated_source_id: SourceId,
387 },
388 ResetSource {
389 source_id: SourceId,
390 },
391 InjectSourceOffsets {
392 source_id: SourceId,
393 split_offsets: HashMap<String, String>,
395 },
396}
397
398#[derive(Debug, Clone)]
403pub struct BarrierInner<M> {
404 pub epoch: EpochPair,
405 pub mutation: M,
406 pub kind: BarrierKind,
407
408 pub tracing_context: TracingContext,
410}
411
412pub type BarrierMutationType = Option<Arc<Mutation>>;
413pub type Barrier = BarrierInner<BarrierMutationType>;
414pub type DispatcherBarrier = BarrierInner<()>;
415
416impl<M: Default> BarrierInner<M> {
417 pub fn new_test_barrier(epoch: u64) -> Self {
419 Self {
420 epoch: EpochPair::new_test_epoch(epoch),
421 kind: BarrierKind::Checkpoint,
422 tracing_context: TracingContext::none(),
423 mutation: Default::default(),
424 }
425 }
426
427 pub fn with_prev_epoch_for_test(epoch: u64, prev_epoch: u64) -> Self {
428 Self {
429 epoch: EpochPair::new(epoch, prev_epoch),
430 kind: BarrierKind::Checkpoint,
431 tracing_context: TracingContext::none(),
432 mutation: Default::default(),
433 }
434 }
435}
436
437impl Barrier {
438 pub fn into_dispatcher(self) -> DispatcherBarrier {
439 DispatcherBarrier {
440 epoch: self.epoch,
441 mutation: (),
442 kind: self.kind,
443 tracing_context: self.tracing_context,
444 }
445 }
446
447 #[must_use]
448 pub fn with_mutation(self, mutation: Mutation) -> Self {
449 Self {
450 mutation: Some(Arc::new(mutation)),
451 ..self
452 }
453 }
454
455 #[cfg(any(test, feature = "test"))]
456 #[must_use]
457 pub fn with_iceberg_pk_index_compaction(self, update: IcebergPkIndexCompactionContext) -> Self {
458 self.with_mutation(Mutation::Update(UpdateMutation {
459 iceberg_pk_index_compaction: Some(update),
460 ..Default::default()
461 }))
462 }
463
464 pub fn iceberg_pk_index_compaction(&self) -> Option<&IcebergPkIndexCompactionContext> {
465 match self.mutation.as_deref() {
466 Some(Mutation::Update(update)) => update.iceberg_pk_index_compaction.as_ref(),
467 _ => None,
468 }
469 }
470
471 #[must_use]
472 pub fn with_stop(self) -> Self {
473 self.with_mutation(Mutation::Stop(StopMutation {
474 dropped_actors: Default::default(),
475 dropped_sink_fragments: Default::default(),
476 }))
477 }
478
479 pub fn is_with_stop_mutation(&self) -> bool {
481 matches!(self.mutation.as_deref(), Some(Mutation::Stop(_)))
482 }
483
484 pub fn is_stop(&self, actor_id: ActorId) -> bool {
486 self.all_stop_actors()
487 .is_some_and(|actors| actors.contains(&actor_id))
488 }
489
490 pub fn is_checkpoint(&self) -> bool {
491 self.kind == BarrierKind::Checkpoint
492 }
493
494 pub fn initial_split_assignment(&self, actor_id: ActorId) -> Option<&[SplitImpl]> {
504 match self.mutation.as_deref()? {
505 Mutation::Update(UpdateMutation { actor_splits, .. })
506 | Mutation::Add(AddMutation {
507 splits: actor_splits,
508 ..
509 }) => actor_splits.get(&actor_id),
510
511 _ => {
512 if cfg!(debug_assertions) {
513 panic!(
514 "the initial mutation of the barrier should not be {:?}",
515 self.mutation
516 );
517 }
518 None
519 }
520 }
521 .map(|s| s.as_slice())
522 }
523
524 pub fn all_stop_actors(&self) -> Option<&HashSet<ActorId>> {
526 self.mutation.as_deref()?.all_stop_actors()
527 }
528
529 pub fn is_newly_added(&self, actor_id: ActorId) -> bool {
535 match self.mutation.as_deref() {
536 Some(Mutation::Add(AddMutation { added_actors, .. })) => {
537 added_actors.contains(&actor_id)
538 }
539 _ => false,
540 }
541 }
542
543 pub fn should_start_fragment_backfill(&self, fragment_id: FragmentId) -> bool {
544 if let Some(Mutation::StartFragmentBackfill { fragment_ids }) = self.mutation.as_deref() {
545 fragment_ids.contains(&fragment_id)
546 } else {
547 false
548 }
549 }
550
551 pub fn has_more_downstream_fragments(&self, upstream_actor_id: ActorId) -> bool {
569 let Some(mutation) = self.mutation.as_deref() else {
570 return false;
571 };
572 match mutation {
573 Mutation::Add(AddMutation { adds, .. }) => adds.get(&upstream_actor_id).is_some(),
575 Mutation::Update(_)
576 | Mutation::Stop(_)
577 | Mutation::Pause
578 | Mutation::Resume
579 | Mutation::SourceChangeSplit(_)
580 | Mutation::Throttle { .. }
581 | Mutation::DropSubscriptions { .. }
582 | Mutation::ConnectorPropsChange(_)
583 | Mutation::StartFragmentBackfill { .. }
584 | Mutation::RefreshStart { .. }
585 | Mutation::ListFinish { .. }
586 | Mutation::LoadFinish { .. }
587 | Mutation::ResetSource { .. }
588 | Mutation::InjectSourceOffsets { .. } => false,
589 }
590 }
591
592 pub fn is_pause_on_startup(&self) -> bool {
594 match self.mutation.as_deref() {
595 Some(Mutation::Add(AddMutation { pause, .. })) => *pause,
596 _ => false,
597 }
598 }
599
600 pub fn is_backfill_pause_on_startup(&self, backfill_fragment_id: FragmentId) -> bool {
601 match self.mutation.as_deref() {
602 Some(Mutation::Add(AddMutation {
603 backfill_nodes_to_pause,
604 ..
605 })) => backfill_nodes_to_pause.contains(&backfill_fragment_id),
606 Some(Mutation::Update(_)) => false,
607 _ => {
608 tracing::warn!(
609 "expected an AddMutation or UpdateMutation on Startup, instead got {:?}",
610 self
611 );
612 false
613 }
614 }
615 }
616
617 pub fn is_resume(&self) -> bool {
619 matches!(self.mutation.as_deref(), Some(Mutation::Resume))
620 }
621
622 pub fn as_update_merge(
625 &self,
626 actor_id: ActorId,
627 upstream_fragment_id: UpstreamFragmentId,
628 ) -> Option<&MergeUpdate> {
629 self.mutation
630 .as_deref()
631 .and_then(|mutation| match mutation {
632 Mutation::Update(UpdateMutation { merges, .. }) => {
633 merges.get(&(actor_id, upstream_fragment_id))
634 }
635 _ => None,
636 })
637 }
638
639 pub fn as_new_upstream_sink(&self, fragment_id: FragmentId) -> Option<&PbNewUpstreamSink> {
642 self.mutation
643 .as_deref()
644 .and_then(|mutation| match mutation {
645 Mutation::Add(AddMutation {
646 new_upstream_sinks, ..
647 }) => new_upstream_sinks.get(&fragment_id),
648 _ => None,
649 })
650 }
651
652 pub fn as_dropped_upstream_sinks(&self) -> Option<&HashSet<FragmentId>> {
654 self.mutation
655 .as_deref()
656 .and_then(|mutation| match mutation {
657 Mutation::Stop(StopMutation {
658 dropped_sink_fragments,
659 ..
660 }) => Some(dropped_sink_fragments),
661 _ => None,
662 })
663 }
664
665 pub fn as_update_vnode_bitmap(&self, actor_id: ActorId) -> Option<Arc<Bitmap>> {
671 self.mutation
672 .as_deref()
673 .and_then(|mutation| match mutation {
674 Mutation::Update(UpdateMutation { vnode_bitmaps, .. }) => {
675 vnode_bitmaps.get(&actor_id).cloned()
676 }
677 _ => None,
678 })
679 }
680
681 pub fn assume_no_update_vnode_bitmap(&self, actor_id: ActorId) -> StreamExecutorResult<()> {
682 if self.as_update_vnode_bitmap(actor_id).is_some() {
683 return Err(anyhow!("updating vnode bitmap in place is not supported").into());
684 }
685 Ok(())
686 }
687
688 pub fn as_sink_schema_change(&self, sink_id: SinkId) -> Option<PbSinkSchemaChange> {
689 self.mutation
690 .as_deref()
691 .and_then(|mutation| match mutation {
692 Mutation::Update(UpdateMutation {
693 sink_schema_change, ..
694 }) => sink_schema_change.get(&sink_id).cloned(),
695 _ => None,
696 })
697 }
698
699 pub fn should_flush_sink_log_store(&self, sink_id: SinkId) -> bool {
700 self.mutation
701 .as_deref()
702 .is_some_and(|mutation| match mutation {
703 Mutation::Add(AddMutation {
704 sink_log_store_flush,
705 ..
706 }) => sink_log_store_flush.contains(&sink_id),
707 _ => false,
708 })
709 }
710
711 pub fn as_subscriptions_to_drop(&self) -> Option<&[SubscriptionUpstreamInfo]> {
712 match self.mutation.as_deref() {
713 Some(Mutation::DropSubscriptions {
714 subscriptions_to_drop,
715 })
716 | Some(Mutation::Update(UpdateMutation {
717 subscriptions_to_drop,
718 ..
719 })) => Some(subscriptions_to_drop.as_slice()),
720 _ => None,
721 }
722 }
723
724 pub fn get_curr_epoch(&self) -> Epoch {
725 Epoch(self.epoch.curr)
726 }
727
728 pub fn tracing_context(&self) -> &TracingContext {
730 &self.tracing_context
731 }
732
733 pub fn added_subscriber_on_mv_table(
734 &self,
735 mv_table_id: TableId,
736 ) -> impl Iterator<Item = SubscriberId> + '_ {
737 if let Some(Mutation::Add(add)) = self.mutation.as_deref() {
738 Some(add)
739 } else {
740 None
741 }
742 .into_iter()
743 .flat_map(move |add| {
744 add.subscriptions_to_add.iter().filter_map(
745 move |(upstream_mv_table_id, subscriber_id)| {
746 if *upstream_mv_table_id == mv_table_id {
747 Some(*subscriber_id)
748 } else {
749 None
750 }
751 },
752 )
753 })
754 }
755}
756
757impl<M: PartialEq> PartialEq for BarrierInner<M> {
758 fn eq(&self, other: &Self) -> bool {
759 self.epoch == other.epoch && self.mutation == other.mutation
760 }
761}
762
763impl Mutation {
764 pub fn backfill_throttle_config(&self, fragment_id: FragmentId) -> Option<&ThrottleConfig> {
766 let Mutation::Throttle(fragment_throttles) = self else {
767 return None;
768 };
769
770 fragment_throttles
771 .get(&fragment_id)
772 .filter(|config| config.throttle_type() == ThrottleType::Backfill)
773 }
774
775 pub fn all_stop_actors(&self) -> Option<&HashSet<ActorId>> {
777 match self {
778 Mutation::Stop(StopMutation { dropped_actors, .. })
779 | Mutation::Update(UpdateMutation { dropped_actors, .. })
780 | Mutation::Add(AddMutation { dropped_actors, .. }) => Some(dropped_actors),
781 _ => None,
782 }
783 }
784
785 pub fn is_stop(&self, actor_id: ActorId) -> bool {
787 self.all_stop_actors()
788 .is_some_and(|actors| actors.contains(&actor_id))
789 }
790
791 #[cfg(test)]
795 pub fn is_stop_mutation(&self) -> bool {
796 matches!(self, Mutation::Stop(_))
797 }
798
799 #[cfg(test)]
800 fn to_protobuf(&self) -> PbMutation {
801 use risingwave_pb::source::{
802 ConnectorSplit, ConnectorSplits, PbCdcTableSnapshotSplitsWithGeneration,
803 };
804 use risingwave_pb::stream_plan::connector_props_change_mutation::ConnectorPropsInfo;
805 use risingwave_pb::stream_plan::{
806 PbAddMutation, PbConnectorPropsChangeMutation, PbDispatchers,
807 PbDropSubscriptionsMutation, PbPauseMutation, PbResumeMutation,
808 PbSourceChangeSplitMutation, PbStartFragmentBackfillMutation, PbStopMutation,
809 PbThrottleMutation, PbUpdateMutation,
810 };
811 let actor_splits_to_protobuf = |actor_splits: &SplitAssignments| {
812 actor_splits
813 .iter()
814 .map(|(&actor_id, splits)| {
815 (
816 actor_id,
817 ConnectorSplits {
818 splits: splits.clone().iter().map(ConnectorSplit::from).collect(),
819 },
820 )
821 })
822 .collect::<HashMap<_, _>>()
823 };
824
825 match self {
826 Mutation::Stop(StopMutation {
827 dropped_actors,
828 dropped_sink_fragments,
829 }) => PbMutation::Stop(PbStopMutation {
830 actors: dropped_actors.iter().copied().collect(),
831 dropped_sink_fragments: dropped_sink_fragments.iter().copied().collect(),
832 }),
833 Mutation::Update(UpdateMutation {
834 dispatchers,
835 merges,
836 vnode_bitmaps,
837 dropped_actors,
838 actor_splits,
839 actor_new_dispatchers,
840 actor_cdc_table_snapshot_splits,
841 sink_schema_change,
842 subscriptions_to_drop,
843 iceberg_pk_index_compaction,
844 }) => PbMutation::Update(PbUpdateMutation {
845 dispatcher_update: dispatchers.values().flatten().cloned().collect(),
846 merge_update: merges.values().cloned().collect(),
847 actor_vnode_bitmap_update: vnode_bitmaps
848 .iter()
849 .map(|(&actor_id, bitmap)| (actor_id, bitmap.to_protobuf()))
850 .collect(),
851 dropped_actors: dropped_actors.iter().copied().collect(),
852 actor_splits: actor_splits_to_protobuf(actor_splits),
853 actor_new_dispatchers: actor_new_dispatchers
854 .iter()
855 .map(|(&actor_id, dispatchers)| {
856 (
857 actor_id,
858 PbDispatchers {
859 dispatchers: dispatchers.clone(),
860 },
861 )
862 })
863 .collect(),
864 actor_cdc_table_snapshot_splits: Some(PbCdcTableSnapshotSplitsWithGeneration {
865 splits:actor_cdc_table_snapshot_splits.splits.iter().map(|(actor_id,(splits, generation))| {
866 (*actor_id, risingwave_pb::source::PbCdcTableSnapshotSplits {
867 splits: splits.iter().map(risingwave_connector::source::cdc::build_cdc_table_snapshot_split).collect(),
868 generation: *generation,
869 })
870 }).collect()
871 }),
872 sink_schema_change: sink_schema_change
873 .iter()
874 .map(|(sink_id, change)| ((*sink_id).as_raw_id(), change.clone()))
875 .collect(),
876 subscriptions_to_drop: subscriptions_to_drop.clone(),
877 iceberg_pk_index_compaction: iceberg_pk_index_compaction.clone(),
878 }),
879 Mutation::Add(AddMutation {
880 adds,
881 added_actors,
882 dropped_actors,
883 splits,
884 pause,
885 subscriptions_to_add,
886 backfill_nodes_to_pause,
887 actor_cdc_table_snapshot_splits,
888 new_upstream_sinks,
889 sink_log_store_flush,
890 }) => PbMutation::Add(PbAddMutation {
891 actor_dispatchers: adds
892 .iter()
893 .map(|(&actor_id, dispatchers)| {
894 (
895 actor_id,
896 PbDispatchers {
897 dispatchers: dispatchers.clone(),
898 },
899 )
900 })
901 .collect(),
902 added_actors: added_actors.iter().copied().collect(),
903 actor_splits: actor_splits_to_protobuf(splits),
904 pause: *pause,
905 subscriptions_to_add: subscriptions_to_add
906 .iter()
907 .map(|(table_id, subscriber_id)| SubscriptionUpstreamInfo {
908 subscriber_id: *subscriber_id,
909 upstream_mv_table_id: *table_id,
910 })
911 .collect(),
912 backfill_nodes_to_pause: backfill_nodes_to_pause.iter().copied().collect(),
913 actor_cdc_table_snapshot_splits:
914 Some(PbCdcTableSnapshotSplitsWithGeneration {
915 splits:actor_cdc_table_snapshot_splits.splits.iter().map(|(actor_id,(splits, generation))| {
916 (*actor_id, risingwave_pb::source::PbCdcTableSnapshotSplits {
917 splits: splits.iter().map(risingwave_connector::source::cdc::build_cdc_table_snapshot_split).collect(),
918 generation: *generation,
919 })
920 }).collect()
921 }),
922 new_upstream_sinks: new_upstream_sinks
923 .iter()
924 .map(|(k, v)| (*k, v.clone()))
925 .collect(),
926 dropped_actors: dropped_actors.iter().copied().collect(),
927 sink_log_store_flush: sink_log_store_flush.iter().copied().collect(),
928 }),
929 Mutation::SourceChangeSplit(changes) => {
930 PbMutation::Splits(PbSourceChangeSplitMutation {
931 actor_splits: changes
932 .iter()
933 .map(|(&actor_id, splits)| {
934 (
935 actor_id,
936 ConnectorSplits {
937 splits: splits
938 .clone()
939 .iter()
940 .map(ConnectorSplit::from)
941 .collect(),
942 },
943 )
944 })
945 .collect(),
946 })
947 }
948 Mutation::Pause => PbMutation::Pause(PbPauseMutation {}),
949 Mutation::Resume => PbMutation::Resume(PbResumeMutation {}),
950 Mutation::Throttle (changes) => PbMutation::Throttle(PbThrottleMutation {
951 fragment_throttle: changes.clone(),
952 }),
953 Mutation::DropSubscriptions {
954 subscriptions_to_drop,
955 } => PbMutation::DropSubscriptions(PbDropSubscriptionsMutation {
956 info: subscriptions_to_drop.clone(),
957 }),
958 Mutation::ConnectorPropsChange(map) => {
959 PbMutation::ConnectorPropsChange(PbConnectorPropsChangeMutation {
960 connector_props_infos: map
961 .iter()
962 .map(|(actor_id, options)| {
963 (
964 *actor_id,
965 ConnectorPropsInfo {
966 connector_props_info: options
967 .iter()
968 .map(|(k, v)| (k.clone(), v.clone()))
969 .collect(),
970 },
971 )
972 })
973 .collect(),
974 })
975 }
976 Mutation::StartFragmentBackfill { fragment_ids } => {
977 PbMutation::StartFragmentBackfill(PbStartFragmentBackfillMutation {
978 fragment_ids: fragment_ids.iter().copied().collect(),
979 })
980 }
981 Mutation::RefreshStart {
982 table_id,
983 associated_source_id,
984 } => PbMutation::RefreshStart(risingwave_pb::stream_plan::RefreshStartMutation {
985 table_id: *table_id,
986 associated_source_id: *associated_source_id,
987 }),
988 Mutation::ListFinish {
989 associated_source_id,
990 } => PbMutation::ListFinish(risingwave_pb::stream_plan::ListFinishMutation {
991 associated_source_id: *associated_source_id,
992 }),
993 Mutation::LoadFinish {
994 associated_source_id,
995 } => PbMutation::LoadFinish(risingwave_pb::stream_plan::LoadFinishMutation {
996 associated_source_id: *associated_source_id,
997 }),
998 Mutation::ResetSource { source_id } => {
999 PbMutation::ResetSource(risingwave_pb::stream_plan::ResetSourceMutation {
1000 source_id: source_id.as_raw_id(),
1001 })
1002 }
1003 Mutation::InjectSourceOffsets {
1004 source_id,
1005 split_offsets,
1006 } => PbMutation::InjectSourceOffsets(
1007 risingwave_pb::stream_plan::InjectSourceOffsetsMutation {
1008 source_id: source_id.as_raw_id(),
1009 split_offsets: split_offsets.clone(),
1010 },
1011 ),
1012 }
1013 }
1014
1015 fn from_protobuf(prost: &PbMutation) -> StreamExecutorResult<Self> {
1016 let mutation = match prost {
1017 PbMutation::Stop(stop) => Mutation::Stop(StopMutation {
1018 dropped_actors: stop.actors.iter().copied().collect(),
1019 dropped_sink_fragments: stop.dropped_sink_fragments.iter().copied().collect(),
1020 }),
1021
1022 PbMutation::Update(update) => Mutation::Update(UpdateMutation {
1023 dispatchers: update
1024 .dispatcher_update
1025 .iter()
1026 .map(|u| (u.actor_id, u.clone()))
1027 .into_group_map(),
1028 merges: update
1029 .merge_update
1030 .iter()
1031 .map(|u| ((u.actor_id, u.upstream_fragment_id), u.clone()))
1032 .collect(),
1033 vnode_bitmaps: update
1034 .actor_vnode_bitmap_update
1035 .iter()
1036 .map(|(&actor_id, bitmap)| (actor_id, Arc::new(bitmap.into())))
1037 .collect(),
1038 dropped_actors: update.dropped_actors.iter().copied().collect(),
1039 actor_splits: update
1040 .actor_splits
1041 .iter()
1042 .map(|(&actor_id, splits)| {
1043 (
1044 actor_id,
1045 splits
1046 .splits
1047 .iter()
1048 .map(|split| split.try_into().unwrap())
1049 .collect(),
1050 )
1051 })
1052 .collect(),
1053 actor_new_dispatchers: update
1054 .actor_new_dispatchers
1055 .iter()
1056 .map(|(&actor_id, dispatchers)| (actor_id, dispatchers.dispatchers.clone()))
1057 .collect(),
1058 actor_cdc_table_snapshot_splits:
1059 build_actor_cdc_table_snapshot_splits_with_generation(
1060 update
1061 .actor_cdc_table_snapshot_splits
1062 .clone()
1063 .unwrap_or_default(),
1064 ),
1065 sink_schema_change: update
1066 .sink_schema_change
1067 .iter()
1068 .map(|(sink_id, change)| (SinkId::from(*sink_id), change.clone()))
1069 .collect(),
1070 subscriptions_to_drop: update.subscriptions_to_drop.clone(),
1071 iceberg_pk_index_compaction: update.iceberg_pk_index_compaction.clone(),
1072 }),
1073
1074 PbMutation::Add(add) => Mutation::Add(AddMutation {
1075 adds: add
1076 .actor_dispatchers
1077 .iter()
1078 .map(|(&actor_id, dispatchers)| (actor_id, dispatchers.dispatchers.clone()))
1079 .collect(),
1080 added_actors: add.added_actors.iter().copied().collect(),
1081 dropped_actors: add.dropped_actors.iter().copied().collect(),
1082 splits: add
1085 .actor_splits
1086 .iter()
1087 .map(|(&actor_id, splits)| {
1088 (
1089 actor_id,
1090 splits
1091 .splits
1092 .iter()
1093 .map(|split| split.try_into().unwrap())
1094 .collect(),
1095 )
1096 })
1097 .collect(),
1098 pause: add.pause,
1099 subscriptions_to_add: add
1100 .subscriptions_to_add
1101 .iter()
1102 .map(
1103 |SubscriptionUpstreamInfo {
1104 subscriber_id,
1105 upstream_mv_table_id,
1106 }| { (*upstream_mv_table_id, *subscriber_id) },
1107 )
1108 .collect(),
1109 backfill_nodes_to_pause: add.backfill_nodes_to_pause.iter().copied().collect(),
1110 actor_cdc_table_snapshot_splits:
1111 build_actor_cdc_table_snapshot_splits_with_generation(
1112 add.actor_cdc_table_snapshot_splits
1113 .clone()
1114 .unwrap_or_default(),
1115 ),
1116 new_upstream_sinks: add
1117 .new_upstream_sinks
1118 .iter()
1119 .map(|(k, v)| (*k, v.clone()))
1120 .collect(),
1121 sink_log_store_flush: add.sink_log_store_flush.iter().copied().collect(),
1122 }),
1123
1124 PbMutation::Splits(s) => {
1125 let mut change_splits: Vec<(ActorId, Vec<SplitImpl>)> =
1126 Vec::with_capacity(s.actor_splits.len());
1127 for (&actor_id, splits) in &s.actor_splits {
1128 if !splits.splits.is_empty() {
1129 change_splits.push((
1130 actor_id,
1131 splits
1132 .splits
1133 .iter()
1134 .map(SplitImpl::try_from)
1135 .try_collect()?,
1136 ));
1137 }
1138 }
1139 Mutation::SourceChangeSplit(change_splits.into_iter().collect())
1140 }
1141 PbMutation::Pause(_) => Mutation::Pause,
1142 PbMutation::Resume(_) => Mutation::Resume,
1143 PbMutation::Throttle(changes) => Mutation::Throttle(changes.fragment_throttle.clone()),
1144 PbMutation::DropSubscriptions(drop) => Mutation::DropSubscriptions {
1145 subscriptions_to_drop: drop.info.clone(),
1146 },
1147 PbMutation::ConnectorPropsChange(alter_connector_props) => {
1148 Mutation::ConnectorPropsChange(
1149 alter_connector_props
1150 .connector_props_infos
1151 .iter()
1152 .map(|(connector_id, options)| {
1153 (
1154 *connector_id,
1155 options
1156 .connector_props_info
1157 .iter()
1158 .map(|(k, v)| (k.clone(), v.clone()))
1159 .collect(),
1160 )
1161 })
1162 .collect(),
1163 )
1164 }
1165 PbMutation::StartFragmentBackfill(start_fragment_backfill) => {
1166 Mutation::StartFragmentBackfill {
1167 fragment_ids: start_fragment_backfill
1168 .fragment_ids
1169 .iter()
1170 .copied()
1171 .collect(),
1172 }
1173 }
1174 PbMutation::RefreshStart(refresh_start) => Mutation::RefreshStart {
1175 table_id: refresh_start.table_id,
1176 associated_source_id: refresh_start.associated_source_id,
1177 },
1178 PbMutation::ListFinish(list_finish) => Mutation::ListFinish {
1179 associated_source_id: list_finish.associated_source_id,
1180 },
1181 PbMutation::LoadFinish(load_finish) => Mutation::LoadFinish {
1182 associated_source_id: load_finish.associated_source_id,
1183 },
1184 PbMutation::ResetSource(reset_source) => Mutation::ResetSource {
1185 source_id: SourceId::from(reset_source.source_id),
1186 },
1187 PbMutation::InjectSourceOffsets(inject) => Mutation::InjectSourceOffsets {
1188 source_id: SourceId::from(inject.source_id),
1189 split_offsets: inject.split_offsets.clone(),
1190 },
1191 };
1192 Ok(mutation)
1193 }
1194}
1195
1196impl<M> BarrierInner<M> {
1197 fn to_protobuf_inner(&self, barrier_fn: impl FnOnce(&M) -> Option<PbMutation>) -> PbBarrier {
1198 let Self {
1199 epoch,
1200 mutation,
1201 kind,
1202 tracing_context,
1203 } = self;
1204
1205 PbBarrier {
1206 epoch: Some(PbEpoch {
1207 curr: epoch.curr,
1208 prev: epoch.prev,
1209 }),
1210 mutation: barrier_fn(mutation).map(|mutation| PbBarrierMutation {
1211 mutation: Some(mutation),
1212 }),
1213 tracing_context: tracing_context.to_protobuf(),
1214 kind: *kind as _,
1215 }
1216 }
1217
1218 fn from_protobuf_inner(
1219 prost: &PbBarrier,
1220 mutation_from_pb: impl FnOnce(Option<&PbMutation>) -> StreamExecutorResult<M>,
1221 ) -> StreamExecutorResult<Self> {
1222 let epoch = prost.get_epoch()?;
1223
1224 Ok(Self {
1225 kind: prost.kind(),
1226 epoch: EpochPair::new(epoch.curr, epoch.prev),
1227 mutation: mutation_from_pb(
1228 (prost.mutation.as_ref()).and_then(|mutation| mutation.mutation.as_ref()),
1229 )?,
1230 tracing_context: TracingContext::from_protobuf(&prost.tracing_context),
1231 })
1232 }
1233
1234 pub fn map_mutation<M2>(self, f: impl FnOnce(M) -> M2) -> BarrierInner<M2> {
1235 BarrierInner {
1236 epoch: self.epoch,
1237 mutation: f(self.mutation),
1238 kind: self.kind,
1239 tracing_context: self.tracing_context,
1240 }
1241 }
1242}
1243
1244impl DispatcherBarrier {
1245 pub fn to_protobuf(&self) -> PbBarrier {
1246 self.to_protobuf_inner(|_| None)
1247 }
1248}
1249
1250impl Barrier {
1251 #[cfg(test)]
1252 pub fn to_protobuf(&self) -> PbBarrier {
1253 self.to_protobuf_inner(|mutation| mutation.as_ref().map(|mutation| mutation.to_protobuf()))
1254 }
1255
1256 pub fn from_protobuf(prost: &PbBarrier) -> StreamExecutorResult<Self> {
1257 Self::from_protobuf_inner(prost, |mutation| {
1258 mutation
1259 .map(|m| Mutation::from_protobuf(m).map(Arc::new))
1260 .transpose()
1261 })
1262 }
1263}
1264
1265#[derive(Debug, PartialEq, Eq, Clone, EstimateSize)]
1266pub struct Watermark {
1267 pub col_idx: usize,
1268 #[estimate_size(ignore)]
1269 pub data_type: DataType,
1270 pub val: ScalarImpl,
1271}
1272
1273impl PartialOrd for Watermark {
1274 fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
1275 Some(self.cmp(other))
1276 }
1277}
1278
1279impl Ord for Watermark {
1280 fn cmp(&self, other: &Self) -> std::cmp::Ordering {
1281 self.val.default_cmp(&other.val)
1282 }
1283}
1284
1285impl Watermark {
1286 pub fn new(col_idx: usize, data_type: DataType, val: ScalarImpl) -> Self {
1287 Self {
1288 col_idx,
1289 data_type,
1290 val,
1291 }
1292 }
1293
1294 pub async fn transform_with_expr(
1295 self,
1296 expr: &NonStrictExpression,
1297 new_col_idx: usize,
1298 ) -> Option<Self> {
1299 let Self { col_idx, val, .. } = self;
1300 let row = {
1301 let mut row = vec![None; col_idx + 1];
1302 row[col_idx] = Some(val);
1303 OwnedRow::new(row)
1304 };
1305 let val = expr.eval_row_infallible(&row).await?;
1306 Some(Self::new(new_col_idx, expr.inner().return_type(), val))
1307 }
1308
1309 pub fn transform_with_indices(self, output_indices: &[usize]) -> Option<Self> {
1312 output_indices
1313 .iter()
1314 .position(|p| *p == self.col_idx)
1315 .map(|new_col_idx| self.with_idx(new_col_idx))
1316 }
1317
1318 pub fn to_protobuf(&self) -> PbWatermark {
1319 PbWatermark {
1320 column: Some(PbInputRef {
1321 index: self.col_idx as _,
1322 r#type: Some(self.data_type.to_protobuf()),
1323 }),
1324 val: Some(&self.val).to_protobuf().into(),
1325 }
1326 }
1327
1328 pub fn from_protobuf(prost: &PbWatermark) -> StreamExecutorResult<Self> {
1329 let col_ref = prost.get_column()?;
1330 let data_type = DataType::from(col_ref.get_type()?);
1331 let val = Datum::from_protobuf(prost.get_val()?, &data_type)?
1332 .expect("watermark value cannot be null");
1333 Ok(Self::new(col_ref.get_index() as _, data_type, val))
1334 }
1335
1336 pub fn with_idx(self, idx: usize) -> Self {
1337 Self::new(idx, self.data_type, self.val)
1338 }
1339}
1340
1341#[cfg_attr(any(test, feature = "test"), derive(PartialEq))]
1342#[derive(Debug, EnumAsInner, Clone)]
1343pub enum MessageInner<M> {
1344 Chunk(StreamChunk),
1345 Barrier(BarrierInner<M>),
1346 Watermark(Watermark),
1347}
1348
1349impl<M> MessageInner<M> {
1350 pub fn map_mutation<M2>(self, f: impl FnOnce(M) -> M2) -> MessageInner<M2> {
1351 match self {
1352 MessageInner::Chunk(chunk) => MessageInner::Chunk(chunk),
1353 MessageInner::Barrier(barrier) => MessageInner::Barrier(barrier.map_mutation(f)),
1354 MessageInner::Watermark(watermark) => MessageInner::Watermark(watermark),
1355 }
1356 }
1357}
1358
1359pub type Message = MessageInner<BarrierMutationType>;
1360pub type DispatcherMessage = MessageInner<()>;
1361
1362#[derive(Debug, EnumAsInner, Clone)]
1365pub enum MessageBatchInner<M> {
1366 Chunk(StreamChunk),
1367 BarrierBatch(Vec<BarrierInner<M>>),
1368 Watermark(Watermark),
1369}
1370pub type MessageBatch = MessageBatchInner<BarrierMutationType>;
1371pub type DispatcherBarriers = Vec<DispatcherBarrier>;
1372pub type DispatcherMessageBatch = MessageBatchInner<()>;
1373
1374impl<M> From<MessageInner<M>> for MessageBatchInner<M> {
1375 fn from(m: MessageInner<M>) -> Self {
1376 match m {
1377 MessageInner::Chunk(c) => Self::Chunk(c),
1378 MessageInner::Barrier(b) => Self::BarrierBatch(vec![b]),
1379 MessageInner::Watermark(w) => Self::Watermark(w),
1380 }
1381 }
1382}
1383
1384impl From<StreamChunk> for Message {
1385 fn from(chunk: StreamChunk) -> Self {
1386 Message::Chunk(chunk)
1387 }
1388}
1389
1390impl<'a> TryFrom<&'a Message> for &'a Barrier {
1391 type Error = ();
1392
1393 fn try_from(m: &'a Message) -> std::result::Result<Self, Self::Error> {
1394 match m {
1395 Message::Chunk(_) => Err(()),
1396 Message::Barrier(b) => Ok(b),
1397 Message::Watermark(_) => Err(()),
1398 }
1399 }
1400}
1401
1402impl Message {
1403 #[cfg(test)]
1408 pub fn is_stop(&self) -> bool {
1409 matches!(
1410 self,
1411 Message::Barrier(Barrier {
1412 mutation,
1413 ..
1414 }) if mutation.as_ref().unwrap().is_stop_mutation()
1415 )
1416 }
1417}
1418
1419impl DispatcherMessageBatch {
1420 pub fn to_protobuf(&self) -> PbStreamMessageBatch {
1421 let prost = match self {
1422 Self::Chunk(stream_chunk) => {
1423 let prost_stream_chunk = stream_chunk.to_protobuf();
1424 StreamMessageBatch::StreamChunk(prost_stream_chunk)
1425 }
1426 Self::BarrierBatch(barrier_batch) => StreamMessageBatch::BarrierBatch(BarrierBatch {
1427 barriers: barrier_batch.iter().map(|b| b.to_protobuf()).collect(),
1428 }),
1429 Self::Watermark(watermark) => StreamMessageBatch::Watermark(watermark.to_protobuf()),
1430 };
1431 PbStreamMessageBatch {
1432 stream_message_batch: Some(prost),
1433 }
1434 }
1435
1436 pub fn from_protobuf(prost: &PbStreamMessageBatch) -> StreamExecutorResult<Self> {
1437 let res = match prost.get_stream_message_batch()? {
1438 StreamMessageBatch::StreamChunk(chunk) => {
1439 Self::Chunk(StreamChunk::from_protobuf(chunk)?)
1440 }
1441 StreamMessageBatch::BarrierBatch(barrier_batch) => {
1442 let barriers = barrier_batch
1443 .barriers
1444 .iter()
1445 .map(|barrier| {
1446 DispatcherBarrier::from_protobuf_inner(barrier, |mutation| {
1447 if mutation.is_some() {
1448 if cfg!(debug_assertions) {
1449 panic!("should not receive message of barrier with mutation");
1450 } else {
1451 warn!(?barrier, "receive message of barrier with mutation");
1452 }
1453 }
1454 Ok(())
1455 })
1456 })
1457 .try_collect()?;
1458 Self::BarrierBatch(barriers)
1459 }
1460 StreamMessageBatch::Watermark(watermark) => {
1461 Self::Watermark(Watermark::from_protobuf(watermark)?)
1462 }
1463 };
1464 Ok(res)
1465 }
1466
1467 pub fn get_encoded_len(msg: &impl ::prost::Message) -> usize {
1468 ::prost::Message::encoded_len(msg)
1469 }
1470}
1471
1472pub type StreamKey = Vec<usize>;
1473pub type StreamKeyRef<'a> = &'a [usize];
1474pub type StreamKeyDataTypes = SmallVec<[DataType; 1]>;
1475
1476pub async fn expect_first_barrier<M: Debug>(
1478 stream: &mut (impl MessageStreamInner<M> + Unpin),
1479) -> StreamExecutorResult<BarrierInner<M>> {
1480 let message = stream
1481 .next()
1482 .instrument_await("expect_first_barrier")
1483 .await
1484 .context("failed to extract the first message: stream closed unexpectedly")??;
1485 let barrier = message
1486 .into_barrier()
1487 .expect("the first message must be a barrier");
1488 assert!(matches!(
1490 barrier.kind,
1491 BarrierKind::Checkpoint | BarrierKind::Initial
1492 ));
1493 Ok(barrier)
1494}
1495
1496pub async fn expect_first_barrier_from_aligned_stream(
1498 stream: &mut (impl AlignedMessageStream + Unpin),
1499) -> StreamExecutorResult<Barrier> {
1500 let message = stream
1501 .next()
1502 .instrument_await("expect_first_barrier")
1503 .await
1504 .context("failed to extract the first message: stream closed unexpectedly")??;
1505 let barrier = message
1506 .into_barrier()
1507 .expect("the first message must be a barrier");
1508 Ok(barrier)
1509}
1510
1511pub trait StreamConsumer: Send + 'static {
1513 type BarrierStream: Stream<Item = StreamResult<Barrier>> + Send;
1514
1515 fn execute(self: Box<Self>) -> Self::BarrierStream;
1516}
1517
1518type BoxedMessageInput<InputId, M> = BoxedInput<InputId, MessageStreamItemInner<M>>;
1519
1520pub struct DynamicReceivers<InputId, M> {
1524 barrier: Option<BarrierInner<M>>,
1526 start_ts: Option<Instant>,
1528 blocked: Vec<BoxedMessageInput<InputId, M>>,
1530 active: FuturesUnordered<StreamFuture<BoxedMessageInput<InputId, M>>>,
1532 buffered_watermarks: BTreeMap<usize, BufferedWatermarks<InputId>>,
1534 barrier_align_duration: Option<LabelGuardedMetric<GenericCounter<AtomicU64>>>,
1536 merge_barrier_align_duration: Option<LabelGuardedMetric<GenericCounter<AtomicU64>>>,
1538}
1539
1540impl<InputId: Clone + Ord + Hash + std::fmt::Debug + Unpin, M: Clone + Unpin> Stream
1541 for DynamicReceivers<InputId, M>
1542{
1543 type Item = MessageStreamItemInner<M>;
1544
1545 fn poll_next(
1546 mut self: Pin<&mut Self>,
1547 cx: &mut std::task::Context<'_>,
1548 ) -> Poll<Option<Self::Item>> {
1549 if self.is_empty() {
1550 return Poll::Ready(None);
1551 }
1552
1553 loop {
1554 match futures::ready!(self.active.poll_next_unpin(cx)) {
1555 Some((Some(Err(e)), _)) => {
1557 return Poll::Ready(Some(Err(e)));
1558 }
1559 Some((Some(Ok(message)), remaining)) => {
1561 let input_id = remaining.id();
1562 match message {
1563 MessageInner::Chunk(chunk) => {
1564 self.active.push(remaining.into_future());
1566 return Poll::Ready(Some(Ok(MessageInner::Chunk(chunk))));
1567 }
1568 MessageInner::Watermark(watermark) => {
1569 self.active.push(remaining.into_future());
1571 if let Some(watermark) = self.handle_watermark(input_id, watermark) {
1572 return Poll::Ready(Some(Ok(MessageInner::Watermark(watermark))));
1573 }
1574 }
1575 MessageInner::Barrier(barrier) => {
1576 if self.blocked.is_empty() {
1578 self.start_ts = Some(Instant::now());
1579 }
1580 self.blocked.push(remaining);
1581 if let Some(current_barrier) = self.barrier.as_ref() {
1582 if current_barrier.epoch != barrier.epoch {
1583 return Poll::Ready(Some(Err(
1584 StreamExecutorError::align_barrier(
1585 current_barrier.clone().map_mutation(|_| None),
1586 barrier.map_mutation(|_| None),
1587 ),
1588 )));
1589 }
1590 } else {
1591 self.barrier = Some(barrier);
1592 }
1593 }
1594 }
1595 }
1596 Some((None, remaining)) => {
1604 return Poll::Ready(Some(Err(StreamExecutorError::channel_closed(format!(
1605 "upstream input {:?} unexpectedly closed",
1606 remaining.id()
1607 )))));
1608 }
1609 None => {
1611 assert!(!self.blocked.is_empty());
1612
1613 let start_ts = self
1614 .start_ts
1615 .take()
1616 .expect("should have received at least one barrier");
1617 if let Some(barrier_align_duration) = &self.barrier_align_duration {
1618 barrier_align_duration.inc_by(start_ts.elapsed().as_nanos() as u64);
1619 }
1620 if let Some(merge_barrier_align_duration) = &self.merge_barrier_align_duration {
1621 merge_barrier_align_duration.inc_by(start_ts.elapsed().as_nanos() as u64);
1622 }
1623
1624 break;
1625 }
1626 }
1627 }
1628
1629 assert!(self.active.is_terminated());
1630
1631 let barrier = self.barrier.take().unwrap();
1632
1633 let upstreams = std::mem::take(&mut self.blocked);
1634 self.extend_active(upstreams);
1635 assert!(!self.active.is_terminated());
1636
1637 Poll::Ready(Some(Ok(MessageInner::Barrier(barrier))))
1638 }
1639}
1640
1641impl<InputId: Clone + Ord + Hash + std::fmt::Debug, M> DynamicReceivers<InputId, M> {
1642 pub fn new(
1643 upstreams: Vec<BoxedMessageInput<InputId, M>>,
1644 barrier_align_duration: Option<LabelGuardedMetric<GenericCounter<AtomicU64>>>,
1645 merge_barrier_align_duration: Option<LabelGuardedMetric<GenericCounter<AtomicU64>>>,
1646 ) -> Self {
1647 let mut this = Self {
1648 barrier: None,
1649 start_ts: None,
1650 blocked: Vec::with_capacity(upstreams.len()),
1651 active: Default::default(),
1652 buffered_watermarks: Default::default(),
1653 merge_barrier_align_duration,
1654 barrier_align_duration,
1655 };
1656 this.extend_active(upstreams);
1657 this
1658 }
1659
1660 pub fn extend_active(
1663 &mut self,
1664 upstreams: impl IntoIterator<Item = BoxedMessageInput<InputId, M>>,
1665 ) {
1666 assert!(self.blocked.is_empty() && self.barrier.is_none());
1667
1668 self.active
1669 .extend(upstreams.into_iter().map(|s| s.into_future()));
1670 }
1671
1672 pub fn handle_watermark(
1674 &mut self,
1675 input_id: InputId,
1676 watermark: Watermark,
1677 ) -> Option<Watermark> {
1678 let col_idx = watermark.col_idx;
1679 let upstream_ids: Vec<_> = self.upstream_input_ids().collect();
1681 let watermarks = self
1682 .buffered_watermarks
1683 .entry(col_idx)
1684 .or_insert_with(|| BufferedWatermarks::with_ids(upstream_ids));
1685 watermarks.handle_watermark(input_id, watermark)
1686 }
1687
1688 pub fn add_upstreams_from(
1691 &mut self,
1692 new_inputs: impl IntoIterator<Item = BoxedMessageInput<InputId, M>>,
1693 ) {
1694 assert!(self.blocked.is_empty() && self.barrier.is_none());
1695
1696 let new_inputs: Vec<_> = new_inputs.into_iter().collect();
1697 let input_ids = new_inputs.iter().map(|input| input.id());
1698 self.buffered_watermarks.values_mut().for_each(|buffers| {
1699 buffers.add_buffers(input_ids.clone());
1701 });
1702 self.active
1703 .extend(new_inputs.into_iter().map(|s| s.into_future()));
1704 }
1705
1706 pub fn remove_upstreams(&mut self, upstream_input_ids: &HashSet<InputId>) {
1710 assert!(self.blocked.is_empty() && self.barrier.is_none());
1711
1712 let new_upstreams = std::mem::take(&mut self.active)
1713 .into_iter()
1714 .map(|s| s.into_inner().unwrap())
1715 .filter(|u| !upstream_input_ids.contains(&u.id()));
1716 self.extend_active(new_upstreams);
1717 self.buffered_watermarks.values_mut().for_each(|buffers| {
1718 buffers.remove_buffer(upstream_input_ids.clone());
1721 });
1722 }
1723
1724 pub fn merge_barrier_align_duration(
1725 &self,
1726 ) -> Option<LabelGuardedMetric<GenericCounter<AtomicU64>>> {
1727 self.merge_barrier_align_duration.clone()
1728 }
1729
1730 pub fn flush_buffered_watermarks(&mut self) {
1731 self.buffered_watermarks
1732 .values_mut()
1733 .for_each(|buffers| buffers.clear());
1734 }
1735
1736 pub fn upstream_input_ids(&self) -> impl Iterator<Item = InputId> + '_ {
1737 self.blocked
1738 .iter()
1739 .map(|s| s.id())
1740 .chain(self.active.iter().map(|s| s.get_ref().unwrap().id()))
1741 }
1742
1743 pub fn is_empty(&self) -> bool {
1744 self.blocked.is_empty() && self.active.is_empty()
1745 }
1746}
1747
1748pub(crate) struct DispatchBarrierBuffer {
1769 buffer: VecDeque<(Barrier, Option<Vec<BoxedActorInput>>)>,
1770 barrier_rx: mpsc::UnboundedReceiver<Barrier>,
1771 recv_state: BarrierReceiverState,
1772 curr_upstream_fragment_id: FragmentId,
1773 actor_id: ActorId,
1774 build_input_ctx: Arc<BuildInputContext>,
1776}
1777
1778struct BuildInputContext {
1779 pub actor_id: ActorId,
1780 pub local_barrier_manager: LocalBarrierManager,
1781 pub metrics: Arc<StreamingMetrics>,
1782 pub fragment_id: FragmentId,
1783 pub actor_config: Arc<StreamingConfig>,
1784}
1785
1786type BoxedNewInputsFuture =
1787 Pin<Box<dyn Future<Output = StreamExecutorResult<Vec<BoxedActorInput>>> + Send>>;
1788
1789enum BarrierReceiverState {
1790 ReceivingBarrier,
1791 CreatingNewInput(Barrier, BoxedNewInputsFuture),
1792}
1793
1794impl DispatchBarrierBuffer {
1795 pub fn new(
1796 barrier_rx: mpsc::UnboundedReceiver<Barrier>,
1797 actor_id: ActorId,
1798 curr_upstream_fragment_id: FragmentId,
1799 local_barrier_manager: LocalBarrierManager,
1800 metrics: Arc<StreamingMetrics>,
1801 fragment_id: FragmentId,
1802 actor_config: Arc<StreamingConfig>,
1803 ) -> Self {
1804 Self {
1805 buffer: VecDeque::new(),
1806 barrier_rx,
1807 recv_state: BarrierReceiverState::ReceivingBarrier,
1808 curr_upstream_fragment_id,
1809 actor_id,
1810 build_input_ctx: Arc::new(BuildInputContext {
1811 actor_id,
1812 local_barrier_manager,
1813 metrics,
1814 fragment_id,
1815 actor_config,
1816 }),
1817 }
1818 }
1819
1820 pub async fn await_next_message(
1821 &mut self,
1822 stream: &mut (impl Stream<Item = StreamExecutorResult<DispatcherMessage>> + Unpin),
1823 metrics: &ActorInputMetrics,
1824 upstream_is_empty: bool,
1825 ) -> StreamExecutorResult<DispatcherMessage> {
1826 if upstream_is_empty {
1827 while self.buffer.is_empty() {
1828 self.try_fetch_barrier_rx(false).await?;
1829 }
1830 let (barrier, _) = self.buffer.front().unwrap();
1831 return Ok(DispatcherMessage::Barrier(
1832 barrier.clone().into_dispatcher(),
1833 ));
1834 }
1835
1836 let mut start_time = Instant::now();
1837 let interval_duration = Duration::from_secs(15);
1838 let mut interval =
1839 tokio::time::interval_at(start_time + interval_duration, interval_duration);
1840
1841 loop {
1842 tokio::select! {
1843 biased;
1844 msg = stream.try_next() => {
1845 metrics
1846 .actor_input_buffer_blocking_duration_ns
1847 .inc_by(start_time.elapsed().as_nanos() as u64);
1848 return msg?.ok_or_else(
1849 || StreamExecutorError::channel_closed("upstream executor closed unexpectedly")
1850 );
1851 }
1852
1853 e = self.continuously_fetch_barrier_rx() => {
1854 return Err(e);
1855 }
1856
1857 _ = interval.tick() => {
1858 start_time = Instant::now();
1859 metrics.actor_input_buffer_blocking_duration_ns.inc_by(interval_duration.as_nanos() as u64);
1860 continue;
1861 }
1862 }
1863 }
1864 }
1865
1866 pub async fn pop_barrier_with_inputs(
1867 &mut self,
1868 barrier: DispatcherBarrier,
1869 ) -> StreamExecutorResult<(Barrier, Option<Vec<BoxedActorInput>>)> {
1870 while self.buffer.is_empty() {
1871 self.try_fetch_barrier_rx(false).await?;
1872 }
1873 let (recv_barrier, inputs) = self.buffer.pop_front().unwrap();
1874 assert_equal_dispatcher_barrier(&recv_barrier, &barrier);
1875
1876 Ok((recv_barrier, inputs))
1877 }
1878
1879 async fn continuously_fetch_barrier_rx(&mut self) -> StreamExecutorError {
1880 loop {
1881 if let Err(e) = self.try_fetch_barrier_rx(true).await {
1882 return e;
1883 }
1884 }
1885 }
1886
1887 async fn try_fetch_barrier_rx(&mut self, pending_on_end: bool) -> StreamExecutorResult<()> {
1888 match &mut self.recv_state {
1889 BarrierReceiverState::ReceivingBarrier => {
1890 let Some(barrier) = self.barrier_rx.recv().await else {
1891 if pending_on_end {
1892 return pending().await;
1893 } else {
1894 return Err(StreamExecutorError::channel_closed(
1895 "barrier channel closed unexpectedly",
1896 ));
1897 }
1898 };
1899 if let Some(fut) = self.pre_apply_barrier(&barrier) {
1900 self.recv_state = BarrierReceiverState::CreatingNewInput(barrier, fut);
1901 } else {
1902 self.buffer.push_back((barrier, None));
1903 }
1904 }
1905 BarrierReceiverState::CreatingNewInput(barrier, fut) => {
1906 let new_inputs = fut.await?;
1907 self.buffer.push_back((barrier.clone(), Some(new_inputs)));
1908 self.recv_state = BarrierReceiverState::ReceivingBarrier;
1909 }
1910 }
1911 Ok(())
1912 }
1913
1914 fn pre_apply_barrier(&mut self, barrier: &Barrier) -> Option<BoxedNewInputsFuture> {
1915 let update = barrier.as_update_merge(self.actor_id, self.curr_upstream_fragment_id)?;
1916 let upstream_fragment_id = update
1917 .new_upstream_fragment_id
1918 .unwrap_or(self.curr_upstream_fragment_id);
1919 self.curr_upstream_fragment_id = upstream_fragment_id;
1923
1924 if !update.added_upstream_actors.is_empty() {
1925 let ctx = self.build_input_ctx.clone();
1926 let added_upstream_actors = update.added_upstream_actors.clone();
1927 let barrier = barrier.clone();
1928 let fut = async move {
1929 try_join_all(added_upstream_actors.iter().map(|upstream_actor| async {
1930 let mut new_input = new_input(
1931 &ctx.local_barrier_manager,
1932 ctx.metrics.clone(),
1933 ctx.actor_id,
1934 ctx.fragment_id,
1935 upstream_actor,
1936 upstream_fragment_id,
1937 ctx.actor_config.clone(),
1938 )
1939 .await?;
1940
1941 let first_barrier = expect_first_barrier(&mut new_input).await?;
1944 assert_equal_dispatcher_barrier(&barrier, &first_barrier);
1945
1946 StreamExecutorResult::Ok(new_input)
1947 }))
1948 .await
1949 }
1950 .boxed();
1951
1952 Some(fut)
1953 } else {
1954 None
1955 }
1956 }
1957}