1use std::collections::{HashMap, HashSet};
16use std::marker::PhantomData;
17use std::mem::take;
18
19use anyhow::anyhow;
20use risingwave_common::bitmap::Bitmap;
21use risingwave_meta_model::WorkerId;
22use risingwave_meta_model::fragment::DistributionType;
23use risingwave_pb::common::{ActorInfo, HostAddress};
24use risingwave_pb::id::{PartialGraphId, SubscriberId};
25use risingwave_pb::stream_plan::update_mutation::{DispatcherUpdate, MergeUpdate};
26use risingwave_pb::stream_plan::{AddMutation, PbDispatcher, StreamNode, UpdateMutation};
27use tracing::warn;
28
29use crate::MetaResult;
30use crate::barrier::rpc::ControlStreamManager;
31use crate::controller::fragment::InflightFragmentInfo;
32use crate::controller::utils::compose_dispatchers;
33use crate::model::{
34 ActorId, ActorNewNoShuffle, ActorUpstreams, DownstreamFragmentRelation, Fragment,
35 FragmentActorDispatchers, FragmentDownstreamRelation, FragmentId, StreamActor,
36 StreamJobActorsToCreate,
37};
38
39type ComposedEdge = (
40 HashMap<ActorId, PbDispatcher>,
41 HashMap<ActorId, ActorUpstreams>,
42 Option<HashMap<ActorId, ActorId>>,
43);
44
45#[derive(Debug)]
49struct EdgeBuilderFragmentInfo {
50 distribution_type: DistributionType,
51 actors: HashMap<ActorId, Option<Bitmap>>,
52 actor_location: HashMap<ActorId, HostAddress>,
53 partial_graph_id: PartialGraphId,
54}
55
56#[derive(Debug)]
57enum FragmentStatus {
58 Existing(EdgeBuilderFragmentInfo),
59 New(EdgeBuilderFragmentInfo),
60 Changed {
61 before: EdgeBuilderFragmentInfo,
62 after: EdgeBuilderFragmentInfo,
63 },
64}
65
66impl EdgeBuilderFragmentInfo {
67 fn from_inflight(
69 info: &InflightFragmentInfo,
70 partial_graph_id: PartialGraphId,
71 control_stream_manager: &ControlStreamManager,
72 ) -> Self {
73 let (actors, actor_location) = info
74 .actors
75 .iter()
76 .map(|(&actor_id, actor)| {
77 (
78 (actor_id, actor.vnode_bitmap.clone()),
79 (actor_id, control_stream_manager.host_addr(actor.worker_id)),
80 )
81 })
82 .unzip();
83 Self {
84 distribution_type: info.distribution_type,
85 actors,
86 actor_location,
87 partial_graph_id,
88 }
89 }
90
91 fn from_fragment(
93 fragment: &Fragment,
94 stream_actors: &HashMap<FragmentId, Vec<StreamActor>>,
95 actor_worker: &HashMap<ActorId, WorkerId>,
96 partial_graph_id: PartialGraphId,
97 control_stream_manager: &ControlStreamManager,
98 ) -> Self {
99 let (actors, actor_location) = stream_actors
100 .get(&fragment.fragment_id)
101 .into_iter()
102 .flatten()
103 .map(|actor| {
104 (
105 (actor.actor_id, actor.vnode_bitmap.clone()),
106 (
107 actor.actor_id,
108 control_stream_manager.host_addr(actor_worker[&actor.actor_id]),
109 ),
110 )
111 })
112 .unzip();
113 Self {
114 distribution_type: fragment.distribution_type.into(),
115 actors,
116 actor_location,
117 partial_graph_id,
118 }
119 }
120}
121
122#[derive(Debug)]
123pub(crate) struct FragmentEdgeBuildResult {
124 upstreams: HashMap<FragmentId, HashMap<ActorId, ActorUpstreams>>,
125 dispatchers: FragmentActorDispatchers,
126 merge_updates: HashMap<FragmentId, Vec<MergeUpdate>>,
127 dispatcher_updates: Vec<DispatcherUpdate>,
128}
129
130impl FragmentEdgeBuildResult {
131 fn validate_terminal_consumption(&self) {
132 let remaining_upstreams = self.upstreams.values().map(HashMap::len).sum::<usize>();
133 let remaining_dispatchers = self.dispatchers.values().map(HashMap::len).sum::<usize>();
134 let unapplied_dispatcher_updates = self.dispatcher_updates.len();
135 let unapplied_merge_updates = self.merge_updates.values().map(Vec::len).sum::<usize>();
136 let has_unconsumed_fields = remaining_upstreams != 0
137 || remaining_dispatchers != 0
138 || unapplied_dispatcher_updates != 0
139 || unapplied_merge_updates != 0;
140
141 debug_assert!(
142 !has_unconsumed_fields,
143 "edge result has unconsumed fields: remaining_upstreams={remaining_upstreams}, remaining_dispatchers={remaining_dispatchers}, unapplied_dispatcher_updates={unapplied_dispatcher_updates}, unapplied_merge_updates={unapplied_merge_updates}"
144 );
145 #[cfg(not(debug_assertions))]
146 if has_unconsumed_fields {
147 warn!(
148 remaining_upstreams,
149 remaining_dispatchers,
150 unapplied_dispatcher_updates,
151 unapplied_merge_updates,
152 "edge result has unconsumed fields"
153 );
154 }
155 }
156
157 pub(crate) fn apply_to_add_mutation(mut self, mutation: &mut AddMutation) {
159 let dispatchers = take(&mut self.dispatchers);
160 self.validate_terminal_consumption();
161 for (actor_id, dispatchers) in dispatchers.into_values().flatten() {
162 mutation
163 .actor_dispatchers
164 .entry(actor_id)
165 .or_default()
166 .dispatchers
167 .extend(dispatchers);
168 }
169 }
170
171 pub(crate) fn apply_to_update_mutation(mut self, mutation: &mut UpdateMutation) {
173 let dispatchers = take(&mut self.dispatchers);
174 let dispatcher_updates = take(&mut self.dispatcher_updates);
175 let merge_updates = take(&mut self.merge_updates);
176 self.validate_terminal_consumption();
177 mutation.dispatcher_update.extend(dispatcher_updates);
178 mutation
179 .merge_update
180 .extend(merge_updates.into_values().flatten());
181 for (actor_id, dispatchers) in dispatchers.into_values().flatten() {
182 mutation
183 .actor_new_dispatchers
184 .entry(actor_id)
185 .or_default()
186 .dispatchers
187 .extend(dispatchers);
188 }
189 }
190
191 pub(crate) fn collect_actors_to_create(
192 &mut self,
193 actors: impl Iterator<
194 Item = (
195 FragmentId,
196 &StreamNode,
197 impl Iterator<Item = (&StreamActor, WorkerId)>,
198 impl IntoIterator<Item = SubscriberId>,
199 ),
200 >,
201 ) -> StreamJobActorsToCreate {
202 let mut actors_to_create = StreamJobActorsToCreate::default();
203 for (fragment_id, node, actors, subscriber_ids) in actors {
204 let subscriber_ids: HashSet<_> = subscriber_ids.into_iter().collect();
205 for (actor, worker_id) in actors {
206 let upstreams = self
207 .upstreams
208 .get_mut(&fragment_id)
209 .and_then(|upstreams| upstreams.remove(&actor.actor_id))
210 .unwrap_or_default();
211 let dispatchers = self
212 .dispatchers
213 .get_mut(&fragment_id)
214 .and_then(|upstreams| upstreams.remove(&actor.actor_id))
215 .unwrap_or_default();
216 actors_to_create
217 .entry(worker_id)
218 .or_default()
219 .entry(fragment_id)
220 .or_insert_with(|| (node.clone(), vec![], subscriber_ids.clone()))
221 .1
222 .push((actor.clone(), upstreams, dispatchers))
223 }
224 }
225 actors_to_create
226 }
227
228 pub(crate) fn take_common_upstream_actors(
234 &mut self,
235 upstream_fragment_id: FragmentId,
236 downstream_fragment_id: FragmentId,
237 ) -> MetaResult<Vec<ActorInfo>> {
238 let downstream_actor_upstreams = self
239 .upstreams
240 .get_mut(&downstream_fragment_id)
241 .ok_or_else(|| {
242 anyhow!("cannot find upstreams for downstream fragment {downstream_fragment_id}")
243 })?;
244 let mut common_upstream_actors = None;
245 for (&downstream_actor_id, actor_upstreams) in downstream_actor_upstreams.iter_mut() {
246 let upstream_actors = actor_upstreams
247 .remove(&upstream_fragment_id)
248 .ok_or_else(|| {
249 anyhow!(
250 "cannot find upstream fragment {upstream_fragment_id} for downstream actor {downstream_actor_id}"
251 )
252 })?;
253 if let Some(common) = &common_upstream_actors {
254 if common != &upstream_actors {
258 return Err(anyhow!(
259 "upstream actors from fragment {upstream_fragment_id} are not common across downstream fragment {downstream_fragment_id}"
260 )
261 .into());
262 }
263 } else {
264 common_upstream_actors = Some(upstream_actors);
265 }
266 }
267 let common_upstream_actors = common_upstream_actors
268 .ok_or_else(|| anyhow!("downstream fragment {downstream_fragment_id} has no actors"))?;
269
270 downstream_actor_upstreams.retain(|_, actor_upstreams| !actor_upstreams.is_empty());
271 if downstream_actor_upstreams.is_empty() {
272 self.upstreams.remove(&downstream_fragment_id);
273 }
274
275 let mut common_upstream_actors: Vec<_> = common_upstream_actors.into_values().collect();
276 common_upstream_actors.sort_unstable_by_key(|actor| actor.actor_id);
277 Ok(common_upstream_actors)
278 }
279}
280
281pub(crate) struct RegisteringFragments;
282pub(crate) struct AddingRelations;
283
284pub(crate) struct FragmentEdgeBuilder<State> {
285 fragments: HashMap<FragmentId, FragmentStatus>,
286 result: FragmentEdgeBuildResult,
287 actor_new_no_shuffle: ActorNewNoShuffle,
288 _state: PhantomData<State>,
289}
290
291impl FragmentEdgeBuilder<RegisteringFragments> {
292 pub(crate) fn new() -> Self {
293 Self {
294 fragments: Default::default(),
295 result: FragmentEdgeBuildResult {
296 upstreams: Default::default(),
297 dispatchers: Default::default(),
298 merge_updates: Default::default(),
299 dispatcher_updates: Default::default(),
300 },
301 actor_new_no_shuffle: Default::default(),
302 _state: PhantomData,
303 }
304 }
305
306 fn add_existing_fragment_infos(
307 &mut self,
308 fragment_infos: impl IntoIterator<Item = (FragmentId, EdgeBuilderFragmentInfo)>,
309 ) {
310 for (fragment_id, info) in fragment_infos {
311 self.fragments
312 .try_insert(fragment_id, FragmentStatus::Existing(info))
313 .expect("non-duplicate");
314 }
315 }
316
317 pub(crate) fn add_existing_fragments<'a>(
318 mut self,
319 fragments: impl IntoIterator<Item = &'a InflightFragmentInfo>,
320 partial_graph_id: PartialGraphId,
321 control_stream_manager: &ControlStreamManager,
322 ) -> Self {
323 self.add_existing_fragment_infos(fragments.into_iter().map(|fragment| {
324 (
325 fragment.fragment_id,
326 EdgeBuilderFragmentInfo::from_inflight(
327 fragment,
328 partial_graph_id,
329 control_stream_manager,
330 ),
331 )
332 }));
333 self
334 }
335
336 fn add_new_fragment_infos(
337 &mut self,
338 fragment_infos: impl IntoIterator<Item = (FragmentId, EdgeBuilderFragmentInfo)>,
339 ) {
340 for (fragment_id, info) in fragment_infos {
341 self.fragments
342 .try_insert(fragment_id, FragmentStatus::New(info))
343 .expect("new fragment must not already be registered");
344 }
345 }
346
347 pub(crate) fn add_new_fragments<'a>(
348 mut self,
349 fragments: impl IntoIterator<Item = &'a InflightFragmentInfo>,
350 partial_graph_id: PartialGraphId,
351 control_stream_manager: &ControlStreamManager,
352 ) -> Self {
353 self.add_new_fragment_infos(fragments.into_iter().map(|fragment| {
354 (
355 fragment.fragment_id,
356 EdgeBuilderFragmentInfo::from_inflight(
357 fragment,
358 partial_graph_id,
359 control_stream_manager,
360 ),
361 )
362 }));
363 self
364 }
365
366 pub(crate) fn add_new_logical_fragments<'a>(
367 mut self,
368 fragments: impl IntoIterator<Item = (PartialGraphId, &'a Fragment)>,
369 stream_actors: &HashMap<FragmentId, Vec<StreamActor>>,
370 actor_worker: &HashMap<ActorId, WorkerId>,
371 control_stream_manager: &ControlStreamManager,
372 ) -> Self {
373 self.add_new_fragment_infos(fragments.into_iter().map(|(partial_graph_id, fragment)| {
374 (
375 fragment.fragment_id,
376 EdgeBuilderFragmentInfo::from_fragment(
377 fragment,
378 stream_actors,
379 actor_worker,
380 partial_graph_id,
381 control_stream_manager,
382 ),
383 )
384 }));
385 self
386 }
387
388 fn replace_existing_fragment_actor_infos(
389 &mut self,
390 fragment_infos: impl IntoIterator<Item = (FragmentId, EdgeBuilderFragmentInfo)>,
391 ) {
392 for (fragment_id, after) in fragment_infos {
393 let status = self
394 .fragments
395 .remove(&fragment_id)
396 .expect("changed fragment must already be registered");
397 let FragmentStatus::Existing(before) = status else {
398 panic!("fragment {fragment_id} can only be changed once");
399 };
400 assert_eq!(
401 before.distribution_type, after.distribution_type,
402 "fragment distribution type cannot change"
403 );
404 for actor_id in before.actors.keys() {
405 assert!(
406 !after.actors.contains_key(actor_id),
407 "changed fragment {fragment_id} must use entirely fresh actor ids: retained {actor_id}"
408 );
409 }
410 self.fragments
411 .insert(fragment_id, FragmentStatus::Changed { before, after });
412 }
413 }
414
415 pub(crate) fn replace_existing_fragment_actors<'a>(
416 mut self,
417 fragments: impl IntoIterator<Item = &'a InflightFragmentInfo>,
418 partial_graph_id: PartialGraphId,
419 control_stream_manager: &ControlStreamManager,
420 ) -> Self {
421 self.replace_existing_fragment_actor_infos(fragments.into_iter().map(|fragment| {
422 (
423 fragment.fragment_id,
424 EdgeBuilderFragmentInfo::from_inflight(
425 fragment,
426 partial_graph_id,
427 control_stream_manager,
428 ),
429 )
430 }));
431 self
432 }
433
434 pub(crate) fn finish_fragments(self) -> FragmentEdgeBuilder<AddingRelations> {
435 FragmentEdgeBuilder {
436 fragments: self.fragments,
437 result: self.result,
438 actor_new_no_shuffle: self.actor_new_no_shuffle,
439 _state: PhantomData,
440 }
441 }
442}
443
444impl FragmentEdgeBuilder<AddingRelations> {
445 pub(crate) fn add_relations(
446 mut self,
447 relations: &FragmentDownstreamRelation,
448 ) -> MetaResult<Self> {
449 for fragment_id in self.fragments.keys().copied().collect::<Vec<_>>() {
450 let Some(relations) = relations.get(&fragment_id) else {
451 continue;
452 };
453 for relation in relations {
454 if self
455 .fragments
456 .contains_key(&relation.downstream_fragment_id)
457 {
458 self.add_edge_inner(fragment_id, relation)?;
459 }
460 }
461 }
462 Ok(self)
463 }
464
465 pub(crate) fn add_edge(
467 mut self,
468 fragment_id: FragmentId,
469 downstream: &DownstreamFragmentRelation,
470 ) -> MetaResult<Self> {
471 self.add_edge_inner(fragment_id, downstream)?;
472 Ok(self)
473 }
474
475 fn add_edge_inner(
476 &mut self,
477 fragment_id: FragmentId,
478 downstream: &DownstreamFragmentRelation,
479 ) -> MetaResult<()> {
480 let Some(fragment_status) = self.fragments.get(&fragment_id) else {
481 if self
482 .fragments
483 .contains_key(&downstream.downstream_fragment_id)
484 {
485 return Err(anyhow!(
486 "cannot find fragment {} with downstream {:?}",
487 fragment_id,
488 downstream
489 )
490 .into());
491 } else {
492 return Ok(());
494 }
495 };
496 let Some(downstream_status) = self.fragments.get(&downstream.downstream_fragment_id) else {
497 return Ok(());
500 };
501 match (fragment_status, downstream_status) {
502 (FragmentStatus::Existing(_), FragmentStatus::Existing(_)) => {}
507 (
508 FragmentStatus::Existing(fragment),
509 FragmentStatus::Changed {
510 before,
511 after: downstream_fragment,
512 },
513 ) => {
514 let (dispatchers, upstreams, no_shuffle_map) =
515 Self::compose_edge(fragment_id, fragment, downstream, downstream_fragment);
516 Self::add_no_shuffle_mapping(
517 &mut self.actor_new_no_shuffle,
518 fragment_id,
519 downstream.downstream_fragment_id,
520 no_shuffle_map,
521 );
522 Self::add_upstreams(
523 &mut self.result.upstreams,
524 downstream.downstream_fragment_id,
525 upstreams,
526 );
527 Self::add_dispatcher_updates(
528 &mut self.result.dispatcher_updates,
529 dispatchers,
530 before.actors.keys().copied(),
531 );
532 }
533 (
534 FragmentStatus::Changed {
535 before,
536 after: fragment,
537 },
538 FragmentStatus::Existing(downstream_fragment),
539 ) => {
540 let (dispatchers, upstreams, no_shuffle_map) =
541 Self::compose_edge(fragment_id, fragment, downstream, downstream_fragment);
542 Self::add_no_shuffle_mapping(
543 &mut self.actor_new_no_shuffle,
544 fragment_id,
545 downstream.downstream_fragment_id,
546 no_shuffle_map,
547 );
548 Self::add_dispatchers(&mut self.result.dispatchers, fragment_id, dispatchers);
549 Self::add_merge_updates(
550 &mut self.result.merge_updates,
551 fragment_id,
552 downstream.downstream_fragment_id,
553 downstream_fragment,
554 upstreams,
555 before.actors.keys().copied(),
556 );
557 }
558 (FragmentStatus::New(_), FragmentStatus::Changed { .. })
559 | (FragmentStatus::Changed { .. }, FragmentStatus::New(_)) => {
560 return Err(anyhow!(
561 "an edge cannot connect new and changed fragments: {} -> {}",
562 fragment_id,
563 downstream.downstream_fragment_id,
564 )
565 .into());
566 }
567 (FragmentStatus::Existing(fragment), FragmentStatus::New(downstream_fragment))
568 | (FragmentStatus::New(fragment), FragmentStatus::Existing(downstream_fragment))
569 | (FragmentStatus::New(fragment), FragmentStatus::New(downstream_fragment))
570 | (
571 FragmentStatus::Changed {
572 after: fragment, ..
573 },
574 FragmentStatus::Changed {
575 after: downstream_fragment,
576 ..
577 },
578 ) => {
579 let (dispatchers, upstreams, no_shuffle_map) =
580 Self::compose_edge(fragment_id, fragment, downstream, downstream_fragment);
581 Self::add_no_shuffle_mapping(
582 &mut self.actor_new_no_shuffle,
583 fragment_id,
584 downstream.downstream_fragment_id,
585 no_shuffle_map,
586 );
587 Self::add_dispatchers(&mut self.result.dispatchers, fragment_id, dispatchers);
588 Self::add_upstreams(
589 &mut self.result.upstreams,
590 downstream.downstream_fragment_id,
591 upstreams,
592 );
593 }
594 }
595 Ok(())
596 }
597
598 fn add_no_shuffle_mapping(
599 actor_new_no_shuffle: &mut ActorNewNoShuffle,
600 fragment_id: FragmentId,
601 downstream_fragment_id: FragmentId,
602 no_shuffle_map: Option<HashMap<ActorId, ActorId>>,
603 ) {
604 if let Some(no_shuffle_map) = no_shuffle_map {
605 actor_new_no_shuffle
606 .entry(fragment_id)
607 .or_default()
608 .insert(downstream_fragment_id, no_shuffle_map);
609 }
610 }
611
612 fn add_dispatchers(
613 fragment_dispatchers: &mut FragmentActorDispatchers,
614 fragment_id: FragmentId,
615 dispatchers: HashMap<ActorId, PbDispatcher>,
616 ) {
617 for (actor_id, dispatcher) in dispatchers {
618 fragment_dispatchers
619 .entry(fragment_id)
620 .or_default()
621 .entry(actor_id)
622 .or_default()
623 .push(dispatcher);
624 }
625 }
626
627 fn add_upstreams(
628 fragment_upstreams: &mut HashMap<FragmentId, HashMap<ActorId, ActorUpstreams>>,
629 fragment_id: FragmentId,
630 upstreams: HashMap<ActorId, ActorUpstreams>,
631 ) {
632 let target = fragment_upstreams.entry(fragment_id).or_default();
633 for (actor_id, actor_upstreams) in upstreams {
634 target.entry(actor_id).or_default().extend(actor_upstreams);
635 }
636 }
637
638 fn add_dispatcher_updates(
639 dispatcher_updates: &mut Vec<DispatcherUpdate>,
640 dispatchers: HashMap<ActorId, PbDispatcher>,
641 removed_downstream_actor_id: impl IntoIterator<Item = ActorId>,
642 ) {
643 let removed_downstream_actor_id: Vec<_> = removed_downstream_actor_id.into_iter().collect();
644 for (actor_id, dispatcher) in dispatchers {
645 dispatcher_updates.push(DispatcherUpdate {
646 actor_id,
647 dispatcher_id: dispatcher.dispatcher_id,
648 hash_mapping: dispatcher.hash_mapping,
649 added_downstream_actor_id: dispatcher.downstream_actor_id,
650 removed_downstream_actor_id: removed_downstream_actor_id.clone(),
651 });
652 }
653 }
654
655 fn add_merge_updates(
656 merge_updates: &mut HashMap<FragmentId, Vec<MergeUpdate>>,
657 fragment_id: FragmentId,
658 downstream_fragment_id: FragmentId,
659 downstream_fragment: &EdgeBuilderFragmentInfo,
660 upstreams: HashMap<ActorId, ActorUpstreams>,
661 removed_upstream_actor_id: impl IntoIterator<Item = ActorId>,
662 ) {
663 let removed_upstream_actor_id: Vec<_> = removed_upstream_actor_id.into_iter().collect();
664 let fragment_merge_updates = merge_updates.entry(downstream_fragment_id).or_default();
665 for actor_id in downstream_fragment.actors.keys() {
666 let added_upstream_actors = upstreams
667 .get(actor_id)
668 .and_then(|upstreams| upstreams.get(&fragment_id))
669 .into_iter()
670 .flat_map(|upstreams| upstreams.values().cloned())
671 .collect();
672 fragment_merge_updates.push(MergeUpdate {
673 actor_id: *actor_id,
674 upstream_fragment_id: fragment_id,
675 new_upstream_fragment_id: None,
676 added_upstream_actors,
677 removed_upstream_actor_id: removed_upstream_actor_id.clone(),
678 });
679 }
680 }
681
682 fn compose_edge_dispatchers(
683 fragment: &EdgeBuilderFragmentInfo,
684 downstream: &DownstreamFragmentRelation,
685 downstream_fragment: &EdgeBuilderFragmentInfo,
686 ) -> (
687 HashMap<ActorId, PbDispatcher>,
688 Option<HashMap<ActorId, ActorId>>,
689 ) {
690 compose_dispatchers(
691 fragment.distribution_type,
692 &fragment.actors,
693 downstream.downstream_fragment_id,
694 downstream_fragment.distribution_type,
695 &downstream_fragment.actors,
696 downstream.dispatcher_type,
697 downstream.dist_key_indices.clone(),
698 downstream.output_mapping.clone(),
699 )
700 }
701
702 fn compose_edge(
703 fragment_id: FragmentId,
704 fragment: &EdgeBuilderFragmentInfo,
705 downstream: &DownstreamFragmentRelation,
706 downstream_fragment: &EdgeBuilderFragmentInfo,
707 ) -> ComposedEdge {
708 let (dispatchers, no_shuffle_map) =
709 Self::compose_edge_dispatchers(fragment, downstream, downstream_fragment);
710 let mut upstreams: HashMap<ActorId, ActorUpstreams> = HashMap::new();
711 for (&actor_id, dispatcher) in &dispatchers {
712 let actor_location = &fragment.actor_location[&actor_id];
713 for &downstream_actor in &dispatcher.downstream_actor_id {
714 upstreams
715 .entry(downstream_actor)
716 .or_default()
717 .entry(fragment_id)
718 .or_default()
719 .insert(
720 actor_id,
721 ActorInfo {
722 actor_id,
723 host: Some(actor_location.clone()),
724 partial_graph_id: fragment.partial_graph_id,
725 },
726 );
727 }
728 }
729 (dispatchers, upstreams, no_shuffle_map)
730 }
731
732 pub(crate) fn replace_upstream(
733 mut self,
734 fragment_id: FragmentId,
735 original_upstream_fragment_id: FragmentId,
736 new_upstream_fragment_id: FragmentId,
737 ) -> Self {
738 let fragment_merge_updates = self.result.merge_updates.entry(fragment_id).or_default();
739 if let Some(fragment_upstreams) = self.result.upstreams.get_mut(&fragment_id) {
740 fragment_upstreams.retain(|&actor_id, actor_upstreams| {
741 if let Some(new_upstreams) = actor_upstreams.remove(&new_upstream_fragment_id) {
742 fragment_merge_updates.push(MergeUpdate {
743 actor_id,
744 upstream_fragment_id: original_upstream_fragment_id,
745 new_upstream_fragment_id: Some(new_upstream_fragment_id),
746 added_upstream_actors: new_upstreams.into_values().collect(),
747 removed_upstream_actor_id: vec![],
748 })
749 } else if cfg!(debug_assertions) {
750 panic!("cannot find new upstreams for actor {} in fragment {} to new_upstream {}. Current upstreams {:?}", actor_id, fragment_id, new_upstream_fragment_id, actor_upstreams);
751 } else {
752 warn!(%actor_id, %fragment_id, %new_upstream_fragment_id, ?actor_upstreams, "cannot find new upstreams for actor");
753 }
754 !actor_upstreams.is_empty()
755 })
756 } else if cfg!(debug_assertions) {
757 panic!(
758 "cannot find new upstreams for fragment {} to new_upstream {} to replace {}. Current upstreams: {:?}",
759 fragment_id,
760 new_upstream_fragment_id,
761 original_upstream_fragment_id,
762 self.result.upstreams
763 );
764 } else {
765 warn!(%fragment_id, %new_upstream_fragment_id, %original_upstream_fragment_id, upstreams = ?self.result.upstreams, "cannot find new upstreams to replace");
766 }
767 self
768 }
769
770 pub(crate) fn build(self) -> (FragmentEdgeBuildResult, ActorNewNoShuffle) {
772 (self.result, self.actor_new_no_shuffle)
773 }
774}
775
776#[cfg(test)]
777mod tests {
778 use risingwave_meta_model::DispatcherType;
779 use risingwave_pb::stream_plan::PbDispatchOutputMapping;
780
781 use super::*;
782
783 #[derive(Clone, Copy, Debug)]
784 enum TestStatus {
785 Existing,
786 New,
787 Changed,
788 }
789
790 fn actor(id: u32) -> ActorId {
791 ActorId::new(id)
792 }
793
794 fn fragment(id: u32) -> FragmentId {
795 FragmentId::new(id)
796 }
797
798 fn edge_info(
799 distribution_type: DistributionType,
800 actors: impl IntoIterator<Item = (u32, Option<Bitmap>)>,
801 ) -> EdgeBuilderFragmentInfo {
802 let actors: HashMap<_, _> = actors
803 .into_iter()
804 .map(|(id, bitmap)| (actor(id), bitmap))
805 .collect();
806 let actor_location = actors
807 .keys()
808 .map(|actor_id| {
809 (
810 *actor_id,
811 HostAddress {
812 host: format!("actor-{actor_id}"),
813 port: 1234,
814 },
815 )
816 })
817 .collect();
818 EdgeBuilderFragmentInfo {
819 distribution_type,
820 actors,
821 actor_location,
822 partial_graph_id: PartialGraphId::new(1),
823 }
824 }
825
826 fn single_info(actor_id: u32) -> EdgeBuilderFragmentInfo {
827 edge_info(DistributionType::Single, [(actor_id, None)])
828 }
829
830 fn relation(target: FragmentId, dispatcher_type: DispatcherType) -> DownstreamFragmentRelation {
831 DownstreamFragmentRelation {
832 downstream_fragment_id: target,
833 dispatcher_type,
834 dist_key_indices: vec![],
835 output_mapping: PbDispatchOutputMapping::default(),
836 }
837 }
838
839 fn build_status_pair(
840 source: TestStatus,
841 target: TestStatus,
842 ) -> MetaResult<FragmentEdgeBuildResult> {
843 let source_fragment = fragment(1);
844 let target_fragment = fragment(2);
845 let existing = [(source_fragment, source, 1), (target_fragment, target, 11)]
846 .into_iter()
847 .filter(|(_, status, _)| !matches!(status, TestStatus::New))
848 .map(|(fragment_id, _, actor_id)| (fragment_id, single_info(actor_id)));
849 let mut builder = FragmentEdgeBuilder::new();
850 builder.add_existing_fragment_infos(existing);
851 for (fragment_id, status, actor_id) in
852 [(source_fragment, source, 2), (target_fragment, target, 12)]
853 {
854 match status {
855 TestStatus::Existing => {}
856 TestStatus::New => {
857 builder.add_new_fragment_infos([(fragment_id, single_info(actor_id))]);
858 }
859 TestStatus::Changed => {
860 builder.replace_existing_fragment_actor_infos([(
861 fragment_id,
862 single_info(actor_id),
863 )]);
864 }
865 }
866 }
867 let (result, _) = builder
868 .finish_fragments()
869 .add_relations(&HashMap::from([(
870 source_fragment,
871 vec![relation(target_fragment, DispatcherType::Broadcast)],
872 )]))?
873 .build();
874 Ok(result)
875 }
876
877 #[test]
878 fn test_endpoint_status_matrix() {
879 for source in [TestStatus::Existing, TestStatus::New, TestStatus::Changed] {
880 for target in [TestStatus::Existing, TestStatus::New, TestStatus::Changed] {
881 let result = build_status_pair(source, target);
882 if matches!(
883 (source, target),
884 (TestStatus::New, TestStatus::Changed) | (TestStatus::Changed, TestStatus::New)
885 ) {
886 assert!(result.is_err(), "source={source:?}, target={target:?}");
887 continue;
888 }
889 let result = result.unwrap();
890 assert_eq!(
891 result.dispatchers.contains_key(&fragment(1)),
892 matches!(source, TestStatus::New | TestStatus::Changed)
893 || matches!((source, target), (TestStatus::Existing, TestStatus::New)),
894 "source={source:?}, target={target:?}"
895 );
896 assert_eq!(
897 result.upstreams.contains_key(&fragment(2)),
898 matches!(target, TestStatus::New | TestStatus::Changed)
899 || matches!((source, target), (TestStatus::New, TestStatus::Existing)),
900 "source={source:?}, target={target:?}"
901 );
902 assert_eq!(
903 !result.dispatcher_updates.is_empty(),
904 matches!(
905 (source, target),
906 (TestStatus::Existing, TestStatus::Changed)
907 ),
908 "source={source:?}, target={target:?}"
909 );
910 assert_eq!(
911 result
912 .merge_updates
913 .values()
914 .any(|updates| !updates.is_empty()),
915 matches!(
916 (source, target),
917 (TestStatus::Changed, TestStatus::Existing)
918 ),
919 "source={source:?}, target={target:?}"
920 );
921 }
922 }
923 }
924
925 #[test]
926 fn test_hash_and_broadcast_dispatcher_updates() {
927 for dispatcher_type in [DispatcherType::Hash, DispatcherType::Broadcast] {
928 let source = fragment(1);
929 let target = fragment(2);
930 let bitmap = || Some(Bitmap::from_iter([true, true]));
931 let mut builder = FragmentEdgeBuilder::new();
932 builder.add_existing_fragment_infos([
933 (source, edge_info(DistributionType::Single, [(1, None)])),
934 (target, edge_info(DistributionType::Hash, [(11, bitmap())])),
935 ]);
936 builder.replace_existing_fragment_actor_infos([(
937 target,
938 edge_info(DistributionType::Hash, [(12, bitmap())]),
939 )]);
940 let (result, _) = builder
941 .finish_fragments()
942 .add_edge(source, &relation(target, dispatcher_type))
943 .unwrap()
944 .build();
945
946 let update = &result.dispatcher_updates[0];
947 assert_eq!(update.added_downstream_actor_id, vec![actor(12)]);
948 assert_eq!(update.removed_downstream_actor_id, vec![actor(11)]);
949 assert_eq!(
950 update.hash_mapping.is_some(),
951 dispatcher_type == DispatcherType::Hash
952 );
953 }
954 }
955
956 #[test]
957 fn test_no_shuffle_changed_fragments() {
958 let source = fragment(1);
959 let target = fragment(2);
960 let mut builder = FragmentEdgeBuilder::new();
961 builder.add_existing_fragment_infos([(source, single_info(1)), (target, single_info(11))]);
962 builder.replace_existing_fragment_actor_infos([
963 (source, single_info(2)),
964 (target, single_info(12)),
965 ]);
966 let (result, actor_new_no_shuffle) = builder
967 .finish_fragments()
968 .add_edge(source, &relation(target, DispatcherType::NoShuffle))
969 .unwrap()
970 .build();
971
972 assert!(result.dispatcher_updates.is_empty());
973 assert!(result.merge_updates.is_empty());
974 assert_eq!(actor_new_no_shuffle[&source][&target][&actor(2)], actor(12));
975 }
976
977 #[test]
978 fn test_no_shuffle_updates_remove_before_actor_superset() {
979 let source = fragment(1);
980 let target = fragment(2);
981 let actors = |left, right| {
982 [
983 (left, Some(Bitmap::from_iter([true, false]))),
984 (right, Some(Bitmap::from_iter([false, true]))),
985 ]
986 };
987 let relations =
988 HashMap::from([(source, vec![relation(target, DispatcherType::NoShuffle)])]);
989
990 let mut builder = FragmentEdgeBuilder::new();
991 builder.add_existing_fragment_infos([
992 (source, edge_info(DistributionType::Hash, actors(1, 2))),
993 (target, edge_info(DistributionType::Hash, actors(11, 12))),
994 ]);
995 builder.replace_existing_fragment_actor_infos([(
996 target,
997 edge_info(DistributionType::Hash, actors(13, 14)),
998 )]);
999 let (result, _) = builder
1000 .finish_fragments()
1001 .add_relations(&relations)
1002 .unwrap()
1003 .build();
1004
1005 assert_eq!(result.dispatcher_updates.len(), 2);
1006 for update in result.dispatcher_updates {
1007 assert_eq!(update.added_downstream_actor_id.len(), 1);
1008 let mut removed = update.removed_downstream_actor_id;
1009 removed.sort_unstable();
1010 assert_eq!(removed, vec![actor(11), actor(12)]);
1011 }
1012
1013 let mut builder = FragmentEdgeBuilder::new();
1014 builder.add_existing_fragment_infos([
1015 (source, edge_info(DistributionType::Hash, actors(1, 2))),
1016 (target, edge_info(DistributionType::Hash, actors(11, 12))),
1017 ]);
1018 builder.replace_existing_fragment_actor_infos([(
1019 source,
1020 edge_info(DistributionType::Hash, actors(3, 4)),
1021 )]);
1022 let (result, _) = builder
1023 .finish_fragments()
1024 .add_relations(&relations)
1025 .unwrap()
1026 .build();
1027
1028 let merge_updates = &result.merge_updates[&target];
1029 assert_eq!(merge_updates.len(), 2);
1030 for update in merge_updates {
1031 assert_eq!(update.added_upstream_actors.len(), 1);
1032 let mut removed = update.removed_upstream_actor_id.clone();
1033 removed.sort_unstable();
1034 assert_eq!(removed, vec![actor(1), actor(2)]);
1035 }
1036 }
1037
1038 #[test]
1039 fn test_replace_upstream() {
1040 let old_source = fragment(1);
1041 let new_source = fragment(2);
1042 let target = fragment(3);
1043 let mut builder = FragmentEdgeBuilder::new();
1044 builder.add_existing_fragment_infos([(target, single_info(11))]);
1045 builder.add_new_fragment_infos([(new_source, single_info(2))]);
1046 let (result, _) = builder
1047 .finish_fragments()
1048 .add_edge(new_source, &relation(target, DispatcherType::Broadcast))
1049 .unwrap()
1050 .replace_upstream(target, old_source, new_source)
1051 .build();
1052
1053 let update = &result.merge_updates[&target][0];
1054 assert_eq!(update.upstream_fragment_id, old_source);
1055 assert_eq!(update.new_upstream_fragment_id, Some(new_source));
1056 assert_eq!(update.added_upstream_actors[0].actor_id, actor(2));
1057 }
1058
1059 #[test]
1060 fn test_new_to_existing_upstreams_are_consumed_for_add_mutation() {
1061 let source = fragment(1);
1062 let target = fragment(2);
1063 let mut builder = FragmentEdgeBuilder::new();
1064 builder.add_existing_fragment_infos([(target, single_info(11))]);
1065 builder.add_new_fragment_infos([(source, single_info(1))]);
1066 let (mut result, _) = builder
1067 .finish_fragments()
1068 .add_edge(source, &relation(target, DispatcherType::Broadcast))
1069 .unwrap()
1070 .build();
1071 let source_actor = StreamActor {
1072 actor_id: actor(1),
1073 fragment_id: source,
1074 vnode_bitmap: None,
1075 mview_definition: Default::default(),
1076 expr_context: None,
1077 config_override: Default::default(),
1078 };
1079 let node = StreamNode::default();
1080 let worker_id: WorkerId = 1.into();
1081 let actors_to_create = result.collect_actors_to_create(std::iter::once((
1082 source,
1083 &node,
1084 std::iter::once((&source_actor, worker_id)),
1085 [],
1086 )));
1087
1088 assert_eq!(actors_to_create[&worker_id][&source].1[0].2.len(), 1);
1089 let upstream_actors = result.take_common_upstream_actors(source, target).unwrap();
1090 assert_eq!(upstream_actors.len(), 1);
1091 assert_eq!(upstream_actors[0].actor_id, actor(1));
1092 result.apply_to_add_mutation(&mut AddMutation::default());
1093 }
1094
1095 #[test]
1096 fn test_take_common_upstream_actors_from_hash_edge() {
1097 let source = fragment(1);
1098 let target = fragment(2);
1099 let actors = |left, right| {
1100 [
1101 (left, Some(Bitmap::from_iter([true, false]))),
1102 (right, Some(Bitmap::from_iter([false, true]))),
1103 ]
1104 };
1105 let mut builder = FragmentEdgeBuilder::new();
1106 builder.add_existing_fragment_infos([(
1107 target,
1108 edge_info(DistributionType::Hash, actors(11, 12)),
1109 )]);
1110 builder.add_new_fragment_infos([(source, edge_info(DistributionType::Hash, actors(1, 2)))]);
1111 let (mut result, _) = builder
1112 .finish_fragments()
1113 .add_edge(source, &relation(target, DispatcherType::Hash))
1114 .unwrap()
1115 .build();
1116
1117 let upstream_actors = result.take_common_upstream_actors(source, target).unwrap();
1118 assert_eq!(
1119 upstream_actors
1120 .into_iter()
1121 .map(|actor| actor.actor_id)
1122 .collect::<Vec<_>>(),
1123 vec![actor(1), actor(2)]
1124 );
1125 assert!(!result.upstreams.contains_key(&target));
1126 }
1127
1128 #[test]
1129 #[should_panic(expected = "remaining_upstreams=1")]
1130 fn test_add_mutation_rejects_uncollected_actor_upstreams() {
1131 let source = fragment(1);
1132 let target = fragment(2);
1133 let mut builder = FragmentEdgeBuilder::new();
1134 builder.add_existing_fragment_infos([(source, single_info(1))]);
1135 builder.add_new_fragment_infos([(target, single_info(11))]);
1136 let (result, _) = builder
1137 .finish_fragments()
1138 .add_edge(source, &relation(target, DispatcherType::Broadcast))
1139 .unwrap()
1140 .build();
1141
1142 result.apply_to_add_mutation(&mut AddMutation::default());
1143 }
1144
1145 #[test]
1146 fn test_update_mutation_consumes_edge_updates_before_validation() {
1147 let source = fragment(1);
1148 let mut result = build_status_pair(TestStatus::Changed, TestStatus::Existing).unwrap();
1149 let source_actor = StreamActor {
1150 actor_id: actor(2),
1151 fragment_id: source,
1152 vnode_bitmap: None,
1153 mview_definition: Default::default(),
1154 expr_context: None,
1155 config_override: Default::default(),
1156 };
1157 let node = StreamNode::default();
1158 result.collect_actors_to_create(std::iter::once((
1159 source,
1160 &node,
1161 std::iter::once((&source_actor, 1.into())),
1162 [],
1163 )));
1164 let mut mutation = UpdateMutation::default();
1165
1166 result.apply_to_update_mutation(&mut mutation);
1167
1168 assert_eq!(mutation.merge_update.len(), 1);
1169 }
1170
1171 #[test]
1172 fn test_attach_new_relation_dispatcher_to_add_mutation() {
1173 let source = fragment(1);
1174 let target = fragment(2);
1175 let mut builder = FragmentEdgeBuilder::new();
1176 builder.add_existing_fragment_infos([(source, single_info(1))]);
1177 builder.add_new_fragment_infos([(target, single_info(11))]);
1178 let (mut result, _) = builder
1179 .finish_fragments()
1180 .add_edge(source, &relation(target, DispatcherType::Broadcast))
1181 .unwrap()
1182 .build();
1183 let target_actor = StreamActor {
1184 actor_id: actor(11),
1185 fragment_id: target,
1186 vnode_bitmap: None,
1187 mview_definition: Default::default(),
1188 expr_context: None,
1189 config_override: Default::default(),
1190 };
1191 let node = StreamNode::default();
1192 result.collect_actors_to_create(std::iter::once((
1193 target,
1194 &node,
1195 std::iter::once((&target_actor, 1.into())),
1196 [],
1197 )));
1198 let mut mutation = AddMutation::default();
1199
1200 result.apply_to_add_mutation(&mut mutation);
1201
1202 assert_eq!(mutation.actor_dispatchers[&actor(1)].dispatchers.len(), 1);
1203 assert_eq!(
1204 mutation.actor_dispatchers[&actor(1)].dispatchers[0].downstream_actor_id,
1205 vec![actor(11)]
1206 );
1207 }
1208
1209 #[test]
1210 fn test_add_relations_only_uses_registered_fragments() {
1211 let external_source = fragment(1);
1212 let source = fragment(2);
1213 let target = fragment(3);
1214 let mut builder = FragmentEdgeBuilder::new();
1215 builder.add_new_fragment_infos([(source, single_info(2)), (target, single_info(3))]);
1216 let (result, _) = builder
1217 .finish_fragments()
1218 .add_relations(&HashMap::from([
1219 (
1220 external_source,
1221 vec![relation(source, DispatcherType::Broadcast)],
1222 ),
1223 (source, vec![relation(target, DispatcherType::Broadcast)]),
1224 ]))
1225 .unwrap()
1226 .build();
1227
1228 assert!(!result.dispatchers.contains_key(&external_source));
1229 assert!(!result.upstreams.contains_key(&source));
1230 assert!(result.dispatchers.contains_key(&source));
1231 assert!(result.upstreams.contains_key(&target));
1232 }
1233
1234 #[test]
1235 fn test_add_edge_rejects_unregistered_upstream() {
1236 let source = fragment(1);
1237 let target = fragment(2);
1238 let mut builder = FragmentEdgeBuilder::new();
1239 builder.add_existing_fragment_infos([(target, single_info(2))]);
1240
1241 let result = builder
1242 .finish_fragments()
1243 .add_edge(source, &relation(target, DispatcherType::Broadcast));
1244
1245 assert!(result.is_err());
1246 }
1247
1248 #[test]
1249 #[should_panic(expected = "must use entirely fresh actor ids")]
1250 fn test_changed_fragment_rejects_retained_actor_ids() {
1251 let fragment = fragment(1);
1252 let mut builder = FragmentEdgeBuilder::new();
1253 builder.add_existing_fragment_infos([(fragment, single_info(1))]);
1254 builder.replace_existing_fragment_actor_infos([(fragment, single_info(1))]);
1255 }
1256}