1use std::assert_matches;
16use std::collections::hash_map::Entry;
17use std::collections::{HashMap, HashSet};
18use std::mem::take;
19use std::sync::atomic::AtomicU32;
20
21use risingwave_common::bail;
22use risingwave_common::bitmap::Bitmap;
23use risingwave_common::catalog::TableId;
24use risingwave_common::hash::VnodeCountCompat;
25use risingwave_common::id::JobId;
26use risingwave_common::util::epoch::Epoch;
27use risingwave_meta_model::fragment::DistributionType;
28use risingwave_meta_model::{DispatcherType, WorkerId, streaming_job};
29use risingwave_pb::common::WorkerNode;
30use risingwave_pb::hummock::HummockVersionStats;
31use risingwave_pb::source::{ConnectorSplit, ConnectorSplits};
32use risingwave_pb::stream_plan::barrier_mutation::{Mutation, PbMutation};
33use risingwave_pb::stream_plan::{
34 AddMutation, PbDropSubscriptionsMutation, PbStartFragmentBackfillMutation,
35 PbSubscriptionUpstreamInfo, PbUpdateMutation, PbUpstreamSinkInfo,
36};
37use tracing::warn;
38
39use crate::barrier::cdc_progress::CdcTableBackfillTracker;
40use crate::barrier::checkpoint::{
41 BatchRefreshJobCheckpointControl, BatchRefreshLogicalFragments, CreatingStreamingJobControl,
42 DatabaseCheckpointControl, IndependentCheckpointJob,
43};
44use crate::barrier::command::{
45 CreateStreamingJobCommandInfo, PostCollectCommand, ReschedulePlan, ThrottleConfigMap,
46};
47use crate::barrier::context::CreateIndependentStreamingJobCommandInfo;
48use crate::barrier::edge_builder::FragmentEdgeBuilder;
49use crate::barrier::info::{
50 BarrierInfo, CreateStreamingJobStatus, InflightDatabaseInfo, InflightStreamingJobInfo,
51 SubscriberType,
52};
53use crate::barrier::partial_graph::{PartialGraphBarrierInfo, PartialGraphManager};
54use crate::barrier::rpc::to_partial_graph_id;
55use crate::barrier::{
56 BarrierKind, Command, CreateStreamingJobType, IndependentStreamingJobType, TracedEpoch,
57};
58use crate::controller::fragment::{InflightActorInfo, InflightFragmentInfo};
59use crate::controller::scale::{
60 ComponentFragmentAligner, EnsembleActorTemplate, LoadedFragment, NoShuffleEnsemble,
61 build_no_shuffle_fragment_graph_edges, find_no_shuffle_graphs,
62};
63use crate::model::{
64 ActorId, ActorNewNoShuffle, FragmentDownstreamRelation, FragmentId, StreamActor, StreamContext,
65 StreamJobActorsToCreate, StreamJobFragmentsToCreate,
66};
67use crate::notification::NotifierStarter;
68use crate::stream::cdc::parallel_cdc_table_backfill_fragment;
69use crate::stream::{
70 GlobalActorIdGen, RefreshCycleActors, ReplaceJobSplitPlan, SourceManager, SplitAssignment,
71 fill_snapshot_backfill_epoch,
72};
73use crate::{MetaError, MetaResult};
74
75pub(in crate::barrier) struct BarrierWorkerState {
77 in_flight_prev_epoch: TracedEpoch,
82
83 pending_non_checkpoint_barriers: Vec<u64>,
85
86 is_paused: bool,
88}
89
90impl BarrierWorkerState {
91 pub(super) fn new() -> Self {
92 Self {
93 in_flight_prev_epoch: TracedEpoch::new(Epoch::now()),
94 pending_non_checkpoint_barriers: vec![],
95 is_paused: false,
96 }
97 }
98
99 pub fn recovery(in_flight_prev_epoch: TracedEpoch, is_paused: bool) -> Self {
100 Self {
101 in_flight_prev_epoch,
102 pending_non_checkpoint_barriers: vec![],
103 is_paused,
104 }
105 }
106
107 pub fn is_paused(&self) -> bool {
108 self.is_paused
109 }
110
111 fn set_is_paused(&mut self, is_paused: bool) {
112 if self.is_paused != is_paused {
113 tracing::info!(
114 currently_paused = self.is_paused,
115 newly_paused = is_paused,
116 "update paused state"
117 );
118 self.is_paused = is_paused;
119 }
120 }
121
122 pub fn in_flight_prev_epoch(&self) -> &TracedEpoch {
123 &self.in_flight_prev_epoch
124 }
125
126 pub fn next_barrier_info(
128 &mut self,
129 is_checkpoint: bool,
130 curr_epoch: TracedEpoch,
131 ) -> BarrierInfo {
132 assert!(
133 self.in_flight_prev_epoch.value() < curr_epoch.value(),
134 "curr epoch regress. {} > {}",
135 self.in_flight_prev_epoch.value(),
136 curr_epoch.value()
137 );
138 let prev_epoch = self.in_flight_prev_epoch.clone();
139 self.in_flight_prev_epoch = curr_epoch.clone();
140 self.pending_non_checkpoint_barriers
141 .push(prev_epoch.value().0);
142 let kind = if is_checkpoint {
143 let epochs = take(&mut self.pending_non_checkpoint_barriers);
144 BarrierKind::Checkpoint(epochs)
145 } else {
146 BarrierKind::Barrier
147 };
148 BarrierInfo {
149 prev_epoch,
150 curr_epoch,
151 kind,
152 }
153 }
154}
155
156pub(super) struct ApplyCommandInfo {
157 pub jobs_to_wait: HashSet<JobId>,
158}
159
160type ApplyCommandResult = (
163 Option<Mutation>,
164 HashSet<TableId>,
165 Option<StreamJobActorsToCreate>,
166 HashMap<WorkerId, HashSet<ActorId>>,
167 PostCollectCommand,
168);
169
170pub(crate) struct RenderResult {
172 pub stream_actors: HashMap<FragmentId, Vec<StreamActor>>,
174 pub actor_location: HashMap<ActorId, WorkerId>,
176}
177
178pub(crate) fn resolve_no_shuffle_ensembles(
187 fragments: &StreamJobFragmentsToCreate,
188 upstream_fragment_downstreams: &FragmentDownstreamRelation,
189) -> MetaResult<Vec<NoShuffleEnsemble>> {
190 let mut new_no_shuffle: HashMap<_, HashSet<_>> = HashMap::new();
192
193 for (upstream_fid, relations) in &fragments.downstreams {
195 for rel in relations {
196 if rel.dispatcher_type == DispatcherType::NoShuffle {
197 new_no_shuffle
198 .entry(*upstream_fid)
199 .or_default()
200 .insert(rel.downstream_fragment_id);
201 }
202 }
203 }
204
205 for (upstream_fid, relations) in upstream_fragment_downstreams {
207 for rel in relations {
208 if rel.dispatcher_type == DispatcherType::NoShuffle {
209 new_no_shuffle
210 .entry(*upstream_fid)
211 .or_default()
212 .insert(rel.downstream_fragment_id);
213 }
214 }
215 }
216
217 let mut ensembles = if new_no_shuffle.is_empty() {
218 Vec::new()
219 } else {
220 let no_shuffle_edges: Vec<(FragmentId, FragmentId)> = new_no_shuffle
222 .iter()
223 .flat_map(|(upstream_fid, downstream_fids)| {
224 downstream_fids
225 .iter()
226 .map(move |downstream_fid| (*upstream_fid, *downstream_fid))
227 })
228 .collect();
229
230 let all_fragment_ids: Vec<FragmentId> = no_shuffle_edges
231 .iter()
232 .flat_map(|(u, d)| [*u, *d])
233 .collect::<HashSet<_>>()
234 .into_iter()
235 .collect();
236
237 let (fwd, bwd) = build_no_shuffle_fragment_graph_edges(no_shuffle_edges);
238 find_no_shuffle_graphs(&all_fragment_ids, &fwd, &bwd)?
239 };
240
241 let covered: HashSet<FragmentId> = ensembles
243 .iter()
244 .flat_map(|e| e.component_fragments())
245 .collect();
246 for fragment_id in fragments.inner.fragments.keys() {
247 if !covered.contains(fragment_id) {
248 ensembles.push(NoShuffleEnsemble::singleton(*fragment_id));
249 }
250 }
251
252 Ok(ensembles)
253}
254
255pub(super) fn render_actors(
266 fragments: &StreamJobFragmentsToCreate,
267 database_info: &InflightDatabaseInfo,
268 definition: &str,
269 ctx: &StreamContext,
270 streaming_job_model: &streaming_job::Model,
271 actor_id_counter: &AtomicU32,
272 worker_map: &HashMap<WorkerId, WorkerNode>,
273 ensembles: &[NoShuffleEnsemble],
274 database_resource_group: &str,
275) -> MetaResult<RenderResult> {
276 let mut actor_assignments: HashMap<FragmentId, HashMap<ActorId, (WorkerId, Option<Bitmap>)>> =
279 HashMap::new();
280
281 for ensemble in ensembles {
282 let existing_fragment_ids: Vec<FragmentId> = ensemble
287 .component_fragments()
288 .filter(|fragment_id| !fragments.inner.fragments.contains_key(fragment_id))
289 .collect();
290
291 let actor_template = if let Some(&first_existing) = existing_fragment_ids.first() {
292 let template = EnsembleActorTemplate::from_existing_inflight_fragment(
293 database_info.fragment(first_existing),
294 );
295
296 for &other_fragment_id in &existing_fragment_ids[1..] {
299 let other = EnsembleActorTemplate::from_existing_inflight_fragment(
300 database_info.fragment(other_fragment_id),
301 );
302 template.assert_aligned_with(&other, first_existing, other_fragment_id);
303 }
304
305 template
306 } else {
307 let first_component = ensemble
309 .component_fragments()
310 .next()
311 .expect("ensemble must have at least one component");
312 let fragment = &fragments.inner.fragments[&first_component];
313 let distribution_type: DistributionType = fragment.distribution_type.into();
314 let vnode_count = fragment.vnode_count();
315
316 for fragment_id in ensemble.component_fragments() {
318 let f = &fragments.inner.fragments[&fragment_id];
319 assert_eq!(
320 vnode_count,
321 f.vnode_count(),
322 "component fragments {} and {} in the same no-shuffle ensemble have \
323 different vnode counts: {} vs {}",
324 first_component,
325 fragment_id,
326 vnode_count,
327 f.vnode_count(),
328 );
329 }
330
331 EnsembleActorTemplate::render_new(
332 streaming_job_model,
333 worker_map,
334 None,
335 database_resource_group.to_owned(),
336 distribution_type,
337 vnode_count,
338 )?
339 };
340
341 for fragment_id in ensemble.component_fragments() {
343 if !fragments.inner.fragments.contains_key(&fragment_id) {
344 continue; }
346 let fragment = &fragments.inner.fragments[&fragment_id];
347 let distribution_type: DistributionType = fragment.distribution_type.into();
348 let aligner =
349 ComponentFragmentAligner::new_persistent(&actor_template, actor_id_counter);
350 let assignments = aligner.align_component_actor(distribution_type);
351 actor_assignments.insert(fragment_id, assignments);
352 }
353 }
354
355 let mut result_stream_actors: HashMap<FragmentId, Vec<StreamActor>> = HashMap::new();
357 let mut result_actor_location: HashMap<ActorId, WorkerId> = HashMap::new();
358
359 for (fragment_id, assignments) in &actor_assignments {
360 let mut actors = Vec::with_capacity(assignments.len());
361 for (&actor_id, (worker_id, vnode_bitmap)) in assignments {
362 result_actor_location.insert(actor_id, *worker_id);
363 actors.push(StreamActor {
364 actor_id,
365 fragment_id: *fragment_id,
366 vnode_bitmap: vnode_bitmap.clone(),
367 mview_definition: definition.to_owned(),
368 expr_context: Some(ctx.to_expr_context()),
369 config_override: ctx.config_override.clone(),
370 });
371 }
372 result_stream_actors.insert(*fragment_id, actors);
373 }
374
375 Ok(RenderResult {
376 stream_actors: result_stream_actors,
377 actor_location: result_actor_location,
378 })
379}
380
381impl DatabaseCheckpointControl {
382 fn take_pending_independent_job_subscriptions_to_drop(
383 &mut self,
384 ) -> Vec<PbSubscriptionUpstreamInfo> {
385 take(&mut self.pending_independent_job_subscriptions_to_drop)
386 .into_iter()
387 .filter(|info| {
388 let upstream_job_id = info.upstream_mv_table_id.as_job_id();
389 if !self.database_info.contains_job(upstream_job_id) {
390 return false;
393 }
394 assert_matches!(
395 self.database_info
396 .unregister_subscriber(upstream_job_id, info.subscriber_id),
397 Some(SubscriberType::SnapshotBackfill)
398 );
399 true
400 })
401 .collect()
402 }
403
404 fn collect_base_info(&self) -> (HashSet<TableId>, HashMap<WorkerId, HashSet<ActorId>>) {
406 let table_ids_to_commit = self.database_info.existing_table_ids().collect();
407 let node_actors =
408 InflightFragmentInfo::actor_ids_to_collect(self.database_info.fragment_infos());
409 (table_ids_to_commit, node_actors)
410 }
411
412 fn apply_simple_command(
416 &self,
417 mutation: Option<Mutation>,
418 command_name: &'static str,
419 ) -> ApplyCommandResult {
420 let (table_ids, node_actors) = self.collect_base_info();
421 (
422 mutation,
423 table_ids,
424 None,
425 node_actors,
426 PostCollectCommand::Command(command_name.to_owned()),
427 )
428 }
429
430 pub(super) fn apply_command(
433 &mut self,
434 command: Option<Command>,
435 notifier: &mut Option<NotifierStarter>,
436 barrier_info: BarrierInfo,
437 partial_graph_manager: &mut PartialGraphManager,
438 hummock_version_stats: &HummockVersionStats,
439 worker_nodes: &HashMap<WorkerId, WorkerNode>,
440 ) -> MetaResult<ApplyCommandInfo> {
441 debug_assert!(
442 !matches!(
443 command,
444 Some(Command::RescheduleIntent {
445 reschedule_plan: None,
446 ..
447 })
448 ),
449 "reschedule intent must be resolved before apply"
450 );
451 if matches!(
452 command,
453 Some(Command::RescheduleIntent {
454 reschedule_plan: None,
455 ..
456 })
457 ) {
458 bail!("reschedule intent must be resolved before apply");
459 }
460
461 fn resolve_source_splits(
466 info: &CreateStreamingJobCommandInfo,
467 render_result: &RenderResult,
468 actor_no_shuffle: &ActorNewNoShuffle,
469 database_info: &InflightDatabaseInfo,
470 ) -> MetaResult<SplitAssignment> {
471 let fragment_actor_ids: HashMap<FragmentId, Vec<ActorId>> = render_result
472 .stream_actors
473 .iter()
474 .map(|(fragment_id, actors)| {
475 (
476 *fragment_id,
477 actors.iter().map(|a| a.actor_id).collect::<Vec<_>>(),
478 )
479 })
480 .collect();
481 let mut resolved = SourceManager::resolve_fragment_to_actor_splits(
482 &info.stream_job_fragments,
483 &info.init_split_assignment,
484 &fragment_actor_ids,
485 )?;
486 resolved.extend(SourceManager::resolve_backfill_splits(
487 &info.stream_job_fragments,
488 actor_no_shuffle,
489 |fragment_id, actor_id| {
490 database_info
491 .fragment(fragment_id)
492 .actors
493 .get(&actor_id)
494 .map(|info| info.splits.clone())
495 },
496 )?);
497 Ok(resolved)
498 }
499
500 let mut notify_database_graph = command.is_some();
501 let mut throttle_config: Option<ThrottleConfigMap> = None;
502
503 let (
507 mutation,
508 mut table_ids_to_commit,
509 mut actors_to_create,
510 mut node_actors,
511 post_collect_command,
512 ) = match command {
513 None => self.apply_simple_command(None, "barrier"),
514 Some(Command::CreateStreamingJob {
515 mut info,
516 job_type:
517 CreateStreamingJobType::Independent {
518 mut snapshot_backfill_info,
519 kind: IndependentStreamingJobType::SnapshotBackfill { since_epoch },
520 },
521 cross_db_snapshot_backfill_info,
522 }) => {
523 notify_database_graph = false;
524 let ensembles = resolve_no_shuffle_ensembles(
525 &info.stream_job_fragments,
526 &info.upstream_fragment_downstreams,
527 )?;
528 let actors = render_actors(
529 &info.stream_job_fragments,
530 &self.database_info,
531 &info.definition,
532 &info.stream_job_fragments.inner.ctx,
533 &info.streaming_job_model,
534 partial_graph_manager
535 .control_stream_manager()
536 .env
537 .actor_id_generator(),
538 worker_nodes,
539 &ensembles,
540 &info.database_resource_group,
541 )?;
542 {
543 assert!(!self.state.is_paused());
544 let (snapshot_epoch, since_timestamp_upstream_log_epochs) =
545 if let Some(since_epoch) = &since_epoch {
546 let (snapshot_epoch, log_epochs) =
547 since_epoch.resolved.as_ref().ok_or_else(|| {
548 MetaError::from(anyhow::anyhow!(
549 "since_timestamp epoch has not been resolved for snapshot backfill"
550 ))
551 })?;
552 (
553 *snapshot_epoch,
554 Some((
555 log_epochs,
556 to_partial_graph_id(self.database_id, None),
557 barrier_info.prev_epoch(),
558 )),
559 )
560 } else {
561 (barrier_info.prev_epoch(), None)
562 };
563 for snapshot_backfill_epoch in snapshot_backfill_info
565 .upstream_mv_table_id_to_backfill_epoch
566 .values_mut()
567 {
568 assert_eq!(
569 snapshot_backfill_epoch.replace(snapshot_epoch),
570 None,
571 "must not set previously"
572 );
573 }
574 for fragment in info.stream_job_fragments.inner.fragments.values_mut() {
575 fill_snapshot_backfill_epoch(
576 &mut fragment.nodes,
577 Some(&snapshot_backfill_info),
578 &cross_db_snapshot_backfill_info,
579 )?;
580 }
581 let job_id = info.stream_job_fragments.stream_job_id();
582 let snapshot_backfill_upstream_tables = snapshot_backfill_info
583 .upstream_mv_table_id_to_backfill_epoch
584 .keys()
585 .cloned()
586 .collect();
587 let (mut edges, actor_new_no_shuffle) = self.database_info.build_edge(
589 Some((&info, true)),
590 None,
591 None,
592 partial_graph_manager.control_stream_manager(),
593 &actors.stream_actors,
594 &actors.actor_location,
595 )?;
596 let resolved_split_assignment = resolve_source_splits(
598 &info,
599 &actors,
600 &actor_new_no_shuffle,
601 &self.database_info,
602 )?;
603 let Entry::Vacant(entry) =
604 self.independent_checkpoint_job_controls.entry(job_id)
605 else {
606 panic!("duplicated creating snapshot backfill job {job_id}");
607 };
608
609 let term_id = self.term_id.as_str();
610 let job = CreatingStreamingJobControl::new(
611 entry,
612 CreateIndependentStreamingJobCommandInfo {
613 info: info.clone(),
614 snapshot_backfill_info: snapshot_backfill_info.clone(),
615 cross_db_snapshot_backfill_info,
616 resolved_split_assignment: resolved_split_assignment.clone(),
617 kind: IndependentStreamingJobType::SnapshotBackfill {
618 since_epoch: None,
619 },
620 },
621 notifier.as_mut(),
622 snapshot_backfill_upstream_tables,
623 snapshot_epoch,
624 since_timestamp_upstream_log_epochs,
625 hummock_version_stats,
626 term_id,
627 partial_graph_manager,
628 &mut edges,
629 &resolved_split_assignment,
630 &actors,
631 )?;
632
633 if let Some(fragment_infos) = job.fragment_infos() {
634 self.database_info.shared_actor_infos.upsert(
635 self.database_id,
636 fragment_infos.values().map(|f| (f, job_id)),
637 );
638 }
639
640 for upstream_mv_table_id in snapshot_backfill_info
641 .upstream_mv_table_id_to_backfill_epoch
642 .keys()
643 {
644 self.database_info.register_subscriber(
645 upstream_mv_table_id.as_job_id(),
646 info.streaming_job.id().as_subscriber_id(),
647 SubscriberType::SnapshotBackfill,
648 );
649 }
650
651 let create_job_type = CreateStreamingJobType::Independent {
652 snapshot_backfill_info,
653 kind: IndependentStreamingJobType::SnapshotBackfill { since_epoch },
654 };
655 let mutation = Command::create_streaming_job_to_mutation(
656 &info,
657 &create_job_type,
658 [],
659 self.state.is_paused(),
660 edges,
661 None,
662 &resolved_split_assignment,
663 &actors.stream_actors,
664 )?;
665
666 let (table_ids, node_actors) = self.collect_base_info();
667 (
668 Some(mutation),
669 table_ids,
670 None,
671 node_actors,
672 PostCollectCommand::barrier(),
673 )
674 }
675 }
676 Some(Command::CreateStreamingJob {
677 mut info,
678 job_type:
679 CreateStreamingJobType::Independent {
680 mut snapshot_backfill_info,
681 kind:
682 IndependentStreamingJobType::BatchRefresh {
683 refresh_interval_sec,
684 },
685 },
686 cross_db_snapshot_backfill_info,
687 }) => {
688 notify_database_graph = false;
689 {
690 if self.state.is_paused() {
691 bail!("cannot create batch refresh job while database barrier is paused");
692 }
693 let snapshot_epoch = barrier_info.prev_epoch();
694 let job_id = info.stream_job_fragments.stream_job_id();
695 let database_id = info.streaming_job.database_id();
696
697 for snapshot_backfill_epoch in snapshot_backfill_info
699 .upstream_mv_table_id_to_backfill_epoch
700 .values_mut()
701 {
702 assert_eq!(
703 snapshot_backfill_epoch.replace(snapshot_epoch),
704 None,
705 "must not set previously"
706 );
707 }
708 for fragment in info.stream_job_fragments.inner.fragments.values_mut() {
709 fill_snapshot_backfill_epoch(
710 &mut fragment.nodes,
711 Some(&snapshot_backfill_info),
712 &cross_db_snapshot_backfill_info,
713 )?;
714 }
715 let snapshot_backfill_upstream_tables: HashSet<TableId> =
716 snapshot_backfill_info
717 .upstream_mv_table_id_to_backfill_epoch
718 .keys()
719 .cloned()
720 .collect();
721
722 let logical = BatchRefreshLogicalFragments {
724 fragments: info
725 .stream_job_fragments
726 .inner
727 .fragments
728 .iter()
729 .map(|(&fragment_id, fragment)| {
730 (
731 fragment_id,
732 LoadedFragment {
733 fragment_id,
734 job_id,
735 fragment_type_mask: fragment.fragment_type_mask,
736 distribution_type: fragment.distribution_type.into(),
737 vnode_count: fragment.vnode_count(),
738 nodes: fragment.nodes.clone(),
739 state_table_ids: fragment
740 .state_table_ids
741 .iter()
742 .copied()
743 .collect(),
744 parallelism: None,
745 },
746 )
747 })
748 .collect(),
749 downstreams: info.stream_job_fragments.downstreams.clone(),
750 };
751
752 let term_id = self.term_id.as_str();
756 let Entry::Vacant(entry) =
757 self.independent_checkpoint_job_controls.entry(job_id)
758 else {
759 panic!("duplicated creating batch refresh job {job_id}");
760 };
761
762 let snapshot_backfill_info_clone = snapshot_backfill_info.clone();
763
764 let subscriber_id =
768 info.stream_job_fragments.stream_job_id().as_subscriber_id();
769 let mutation = Mutation::Add(AddMutation {
770 actor_dispatchers: Default::default(),
771 added_actors: Default::default(),
772 actor_splits: Default::default(),
773 pause: false,
774 subscriptions_to_add: snapshot_backfill_info_clone
775 .upstream_mv_table_id_to_backfill_epoch
776 .keys()
777 .map(|table_id| PbSubscriptionUpstreamInfo {
778 subscriber_id,
779 upstream_mv_table_id: *table_id,
780 })
781 .collect(),
782 backfill_nodes_to_pause: Default::default(),
783 actor_cdc_table_snapshot_splits: None,
784 new_upstream_sinks: Default::default(),
785 dropped_actors: Default::default(),
786 sink_log_store_flush: Default::default(),
787 });
788
789 let job = BatchRefreshJobCheckpointControl::new(
790 entry,
791 database_id,
792 job_id,
793 CreateIndependentStreamingJobCommandInfo {
794 info: info.clone(),
795 snapshot_backfill_info: snapshot_backfill_info_clone.clone(),
796 cross_db_snapshot_backfill_info,
797 resolved_split_assignment: Default::default(),
798 kind: IndependentStreamingJobType::BatchRefresh {
799 refresh_interval_sec,
800 },
801 },
802 notifier.as_mut(),
803 snapshot_backfill_upstream_tables,
804 snapshot_epoch,
805 hummock_version_stats,
806 term_id,
807 partial_graph_manager,
808 &logical,
809 worker_nodes,
810 refresh_interval_sec,
811 )?;
812
813 if let Some(fragment_infos) = job.fragment_infos() {
814 self.database_info.shared_actor_infos.upsert(
815 self.database_id,
816 fragment_infos.values().map(|f| (f, job_id)),
817 );
818 }
819
820 for upstream_mv_table_id in snapshot_backfill_info_clone
822 .upstream_mv_table_id_to_backfill_epoch
823 .keys()
824 {
825 self.database_info.register_subscriber(
826 upstream_mv_table_id.as_job_id(),
827 info.streaming_job.id().as_subscriber_id(),
828 SubscriberType::SnapshotBackfill,
829 );
830 }
831
832 let (table_ids, node_actors) = self.collect_base_info();
833 (
834 Some(mutation),
835 table_ids,
836 None,
837 node_actors,
838 PostCollectCommand::barrier(),
839 )
840 }
841 }
842 Some(Command::CreateStreamingJob {
843 mut info,
844 job_type,
845 cross_db_snapshot_backfill_info,
846 }) => {
847 let ensembles = resolve_no_shuffle_ensembles(
848 &info.stream_job_fragments,
849 &info.upstream_fragment_downstreams,
850 )?;
851 let actors = render_actors(
852 &info.stream_job_fragments,
853 &self.database_info,
854 &info.definition,
855 &info.stream_job_fragments.inner.ctx,
856 &info.streaming_job_model,
857 partial_graph_manager
858 .control_stream_manager()
859 .env
860 .actor_id_generator(),
861 worker_nodes,
862 &ensembles,
863 &info.database_resource_group,
864 )?;
865 for fragment in info.stream_job_fragments.inner.fragments.values_mut() {
866 fill_snapshot_backfill_epoch(
867 &mut fragment.nodes,
868 None,
869 &cross_db_snapshot_backfill_info,
870 )?;
871 }
872
873 let new_upstream_sink =
875 if let CreateStreamingJobType::SinkIntoTable(ref ctx) = job_type {
876 Some(ctx)
877 } else {
878 None
879 };
880
881 let (mut edges, actor_new_no_shuffle) = self.database_info.build_edge(
882 Some((&info, false)),
883 None,
884 new_upstream_sink,
885 partial_graph_manager.control_stream_manager(),
886 &actors.stream_actors,
887 &actors.actor_location,
888 )?;
889 let resolved_split_assignment = resolve_source_splits(
891 &info,
892 &actors,
893 &actor_new_no_shuffle,
894 &self.database_info,
895 )?;
896
897 let old_sink_job_id = info
898 .replace_sink
899 .as_ref()
900 .map(|old_sink_id| old_sink_id.as_job_id());
901 if old_sink_job_id.is_some()
902 && matches!(job_type, CreateStreamingJobType::Independent { .. })
903 {
904 bail!("replace sink must not use snapshot backfill");
905 }
906
907 let cdc_tracker = if let Some(splits) = &info.cdc_table_snapshot_splits {
909 let (fragment, _) =
910 parallel_cdc_table_backfill_fragment(info.stream_job_fragments.fragments())
911 .expect("should have parallel cdc fragment");
912 Some(CdcTableBackfillTracker::new(
913 fragment.fragment_id,
914 splits.clone(),
915 ))
916 } else {
917 None
918 };
919 self.database_info
920 .pre_apply_new_job(info.streaming_job.id(), cdc_tracker);
921 self.database_info.pre_apply_new_fragments(
922 info.stream_job_fragments
923 .new_fragment_info(
924 &actors.stream_actors,
925 &actors.actor_location,
926 &resolved_split_assignment,
927 )
928 .map(|(fragment_id, fragment_infos)| {
929 (fragment_id, info.streaming_job.id(), fragment_infos)
930 }),
931 );
932 if let CreateStreamingJobType::SinkIntoTable(ref ctx) = job_type {
933 let downstream_fragment_id = ctx.new_sink_downstream.downstream_fragment_id;
934 self.database_info.pre_apply_add_node_upstream(
935 downstream_fragment_id,
936 &PbUpstreamSinkInfo {
937 upstream_fragment_id: ctx.sink_fragment_id,
938 sink_output_schema: ctx.sink_output_fields.clone(),
939 project_exprs: ctx.project_exprs.clone(),
940 },
941 );
942 }
943
944 let (table_ids, node_actors) = self.collect_base_info();
945 let dropped_actors = if let Some(old_sink_job_id) = old_sink_job_id {
946 let Some(job) = self.database_info.post_apply_remove_job(old_sink_job_id)
947 else {
948 bail!(
949 "old sink job {} not found in barrier state",
950 old_sink_job_id
951 );
952 };
953 job.fragment_infos
954 .values()
955 .flat_map(|fragment| fragment.actors.keys().copied())
956 .collect()
957 } else {
958 vec![]
959 };
960
961 let actors_to_create = Some(Command::create_streaming_job_actors_to_create(
963 &info,
964 &mut edges,
965 &actors.stream_actors,
966 &actors.actor_location,
967 ));
968
969 let actor_cdc_table_snapshot_splits = self
971 .database_info
972 .assign_cdc_backfill_splits(info.stream_job_fragments.stream_job_id())?;
973
974 let is_currently_paused = self.state.is_paused();
976 let mutation = Command::create_streaming_job_to_mutation(
977 &info,
978 &job_type,
979 dropped_actors,
980 is_currently_paused,
981 edges,
982 actor_cdc_table_snapshot_splits,
983 &resolved_split_assignment,
984 &actors.stream_actors,
985 )?;
986
987 (
988 Some(mutation),
989 table_ids,
990 actors_to_create,
991 node_actors,
992 PostCollectCommand::CreateStreamingJob {
993 info,
994 job_type,
995 cross_db_snapshot_backfill_info,
996 resolved_split_assignment,
997 },
998 )
999 }
1000
1001 Some(Command::Flush) => self.apply_simple_command(None, "Flush"),
1002
1003 Some(Command::Pause) => {
1004 let prev_is_paused = self.state.is_paused();
1005 self.state.set_is_paused(true);
1006 let mutation = Command::pause_to_mutation(prev_is_paused);
1007 let (table_ids, node_actors) = self.collect_base_info();
1008 (
1009 mutation,
1010 table_ids,
1011 None,
1012 node_actors,
1013 PostCollectCommand::Command("Pause".to_owned()),
1014 )
1015 }
1016
1017 Some(Command::Resume) => {
1018 let prev_is_paused = self.state.is_paused();
1019 self.state.set_is_paused(false);
1020 let mutation = Command::resume_to_mutation(prev_is_paused);
1021 let (table_ids, node_actors) = self.collect_base_info();
1022 (
1023 mutation,
1024 table_ids,
1025 None,
1026 node_actors,
1027 PostCollectCommand::Command("Resume".to_owned()),
1028 )
1029 }
1030
1031 Some(Command::Throttle { mut config }) => {
1032 let mutation = self.database_info.pre_apply_throttle(&mut config);
1033 notify_database_graph = mutation.is_some();
1034 throttle_config = Some(config);
1035 self.apply_simple_command(mutation, "Throttle")
1036 }
1037
1038 Some(Command::DropStreamingJobs {
1039 streaming_job_ids,
1040 unregistered_state_table_ids: _,
1041 dropped_sink_fragment_by_targets,
1042 }) => {
1043 for (target_fragment, sink_fragments) in &dropped_sink_fragment_by_targets {
1045 self.database_info
1046 .pre_apply_drop_node_upstream(*target_fragment, sink_fragments);
1047 }
1048
1049 let (table_ids, node_actors) = self.collect_base_info();
1050
1051 let mut actors = Vec::new();
1052 for job_id in streaming_job_ids {
1053 let Some(job) = self.database_info.post_apply_remove_job(job_id) else {
1054 warn!(
1055 %job_id,
1056 "skip drop payload for streaming job that has already been removed from barrier worker"
1057 );
1058 continue;
1059 };
1060
1061 for fragment in job.fragment_infos.values() {
1062 actors.extend(fragment.actors.keys().copied());
1063 }
1064 }
1065
1066 let mutation = Some(Command::drop_streaming_jobs_to_mutation(
1067 &actors,
1068 &dropped_sink_fragment_by_targets,
1069 ));
1070 (
1071 mutation,
1072 table_ids,
1073 None,
1074 node_actors,
1075 PostCollectCommand::DropStreamingJobs,
1076 )
1077 }
1078
1079 Some(Command::RescheduleIntent {
1080 reschedule_plan, ..
1081 }) => {
1082 let ReschedulePlan {
1083 reschedules,
1084 fragment_actors,
1085 } = reschedule_plan
1086 .as_ref()
1087 .expect("reschedule intent should be resolved in global barrier worker");
1088
1089 for (fragment_id, reschedule) in reschedules {
1091 self.database_info.pre_apply_reschedule(
1092 *fragment_id,
1093 reschedule
1094 .added_actors
1095 .iter()
1096 .flat_map(|(node_id, actors): (&WorkerId, &Vec<ActorId>)| {
1097 actors.iter().map(|actor_id| {
1098 (
1099 *actor_id,
1100 InflightActorInfo {
1101 worker_id: *node_id,
1102 vnode_bitmap: reschedule
1103 .newly_created_actors
1104 .get(actor_id)
1105 .expect("should exist")
1106 .0
1107 .0
1108 .vnode_bitmap
1109 .clone(),
1110 splits: reschedule
1111 .actor_splits
1112 .get(actor_id)
1113 .cloned()
1114 .unwrap_or_default(),
1115 },
1116 )
1117 })
1118 })
1119 .collect(),
1120 reschedule
1121 .vnode_bitmap_updates
1122 .iter()
1123 .filter(|(actor_id, _)| {
1124 !reschedule.newly_created_actors.contains_key(*actor_id)
1125 })
1126 .map(|(actor_id, bitmap)| (*actor_id, bitmap.clone()))
1127 .collect(),
1128 reschedule.actor_splits.clone(),
1129 );
1130 }
1131
1132 let (table_ids, node_actors) = self.collect_base_info();
1133
1134 let actors_to_create = Some(Command::reschedule_actors_to_create(
1136 reschedules,
1137 fragment_actors,
1138 &self.database_info,
1139 partial_graph_manager.control_stream_manager(),
1140 ));
1141
1142 self.database_info
1144 .post_apply_reschedules(reschedules.iter().map(|(fragment_id, reschedule)| {
1145 (
1146 *fragment_id,
1147 reschedule.removed_actors.iter().cloned().collect(),
1148 )
1149 }));
1150
1151 let mutation = Command::reschedule_to_mutation(
1153 reschedules,
1154 fragment_actors,
1155 partial_graph_manager.control_stream_manager(),
1156 &mut self.database_info,
1157 )?;
1158
1159 let reschedules = reschedule_plan
1160 .expect("reschedule intent should be resolved in global barrier worker")
1161 .reschedules;
1162 (
1163 mutation,
1164 table_ids,
1165 actors_to_create,
1166 node_actors,
1167 PostCollectCommand::Reschedule { reschedules },
1168 )
1169 }
1170
1171 Some(Command::ReplaceStreamJob(plan)) => {
1172 let ensembles = resolve_no_shuffle_ensembles(
1173 &plan.new_fragments,
1174 &plan.upstream_fragment_downstreams,
1175 )?;
1176 let mut render_result = render_actors(
1177 &plan.new_fragments,
1178 &self.database_info,
1179 "", &plan.new_fragments.inner.ctx,
1181 &plan.streaming_job_model,
1182 partial_graph_manager
1183 .control_stream_manager()
1184 .env
1185 .actor_id_generator(),
1186 worker_nodes,
1187 &ensembles,
1188 &plan.database_resource_group,
1189 )?;
1190
1191 if let Some(sinks) = &plan.auto_refresh_schema_sinks {
1194 let actor_id_counter = partial_graph_manager
1195 .control_stream_manager()
1196 .env
1197 .actor_id_generator();
1198 for sink_ctx in sinks {
1199 let original_fragment_id = sink_ctx.original_fragment.fragment_id;
1200 let original_frag_info = self.database_info.fragment(original_fragment_id);
1201 let actor_template = EnsembleActorTemplate::from_existing_inflight_fragment(
1202 original_frag_info,
1203 );
1204 let new_aligner = ComponentFragmentAligner::new_persistent(
1205 &actor_template,
1206 actor_id_counter,
1207 );
1208 let distribution_type: DistributionType =
1209 sink_ctx.new_fragment.distribution_type.into();
1210 let actor_assignments =
1211 new_aligner.align_component_actor(distribution_type);
1212 let new_fragment_id = sink_ctx.new_fragment.fragment_id;
1213 let mut actors = Vec::with_capacity(actor_assignments.len());
1214 for (&actor_id, (worker_id, vnode_bitmap)) in &actor_assignments {
1215 render_result.actor_location.insert(actor_id, *worker_id);
1216 actors.push(StreamActor {
1217 actor_id,
1218 fragment_id: new_fragment_id,
1219 vnode_bitmap: vnode_bitmap.clone(),
1220 mview_definition: String::new(),
1221 expr_context: Some(sink_ctx.ctx.to_expr_context()),
1222 config_override: sink_ctx.ctx.config_override.clone(),
1223 });
1224 }
1225 render_result.stream_actors.insert(new_fragment_id, actors);
1226 }
1227 }
1228
1229 let (mut edges, actor_new_no_shuffle) = self.database_info.build_edge(
1231 None,
1232 Some(&plan),
1233 None,
1234 partial_graph_manager.control_stream_manager(),
1235 &render_result.stream_actors,
1236 &render_result.actor_location,
1237 )?;
1238
1239 let fragment_actor_ids: HashMap<FragmentId, Vec<ActorId>> = render_result
1241 .stream_actors
1242 .iter()
1243 .map(|(fragment_id, actors)| {
1244 (
1245 *fragment_id,
1246 actors.iter().map(|a| a.actor_id).collect::<Vec<_>>(),
1247 )
1248 })
1249 .collect();
1250 let resolved_split_assignment = match &plan.split_plan {
1251 ReplaceJobSplitPlan::Discovered(discovered) => {
1252 SourceManager::resolve_fragment_to_actor_splits(
1253 &plan.new_fragments,
1254 discovered,
1255 &fragment_actor_ids,
1256 )?
1257 }
1258 ReplaceJobSplitPlan::AlignFromPrevious => {
1259 SourceManager::resolve_replace_source_splits(
1260 &plan.new_fragments,
1261 &plan.replace_upstream,
1262 &actor_new_no_shuffle,
1263 |_fragment_id, actor_id| {
1264 self.database_info.fragment_infos().find_map(|fragment| {
1265 fragment
1266 .actors
1267 .get(&actor_id)
1268 .map(|info| info.splits.clone())
1269 })
1270 },
1271 )?
1272 }
1273 };
1274
1275 self.database_info.pre_apply_new_fragments(
1277 plan.new_fragments
1278 .new_fragment_info(
1279 &render_result.stream_actors,
1280 &render_result.actor_location,
1281 &resolved_split_assignment,
1282 )
1283 .map(|(fragment_id, new_fragment)| {
1284 (fragment_id, plan.streaming_job.id(), new_fragment)
1285 }),
1286 );
1287 for (fragment_id, replace_map) in &plan.replace_upstream {
1288 self.database_info
1289 .pre_apply_replace_node_upstream(*fragment_id, replace_map);
1290 }
1291 if let Some(sinks) = &plan.auto_refresh_schema_sinks {
1292 self.database_info
1293 .pre_apply_new_fragments(sinks.iter().map(|sink| {
1294 (
1295 sink.new_fragment.fragment_id,
1296 sink.original_sink.id.as_job_id(),
1297 sink.new_fragment_info(
1298 &render_result.stream_actors,
1299 &render_result.actor_location,
1300 ),
1301 )
1302 }));
1303 }
1304
1305 let (table_ids, node_actors) = self.collect_base_info();
1306
1307 let actors_to_create = Some(Command::replace_stream_job_actors_to_create(
1309 &plan,
1310 &mut edges,
1311 &self.database_info,
1312 &render_result.stream_actors,
1313 &render_result.actor_location,
1314 ));
1315
1316 let mutation = Some(Command::replace_stream_job_to_mutation(
1319 &plan,
1320 edges,
1321 &mut self.database_info,
1322 &resolved_split_assignment,
1323 )?);
1324
1325 {
1327 let mut fragment_ids_to_remove: Vec<_> = plan
1328 .old_fragments
1329 .fragments
1330 .values()
1331 .map(|f| f.fragment_id)
1332 .collect();
1333 if let Some(sinks) = &plan.auto_refresh_schema_sinks {
1334 fragment_ids_to_remove
1335 .extend(sinks.iter().map(|sink| sink.original_fragment.fragment_id));
1336 }
1337 self.database_info
1338 .post_apply_remove_fragments(fragment_ids_to_remove);
1339 }
1340
1341 (
1342 mutation,
1343 table_ids,
1344 actors_to_create,
1345 node_actors,
1346 PostCollectCommand::ReplaceStreamJob {
1347 plan,
1348 resolved_split_assignment,
1349 },
1350 )
1351 }
1352
1353 Some(Command::SourceChangeSplit(split_state)) => {
1354 self.database_info.pre_apply_split_assignments(
1356 split_state
1357 .split_assignment
1358 .iter()
1359 .map(|(&fragment_id, splits)| (fragment_id, splits.clone())),
1360 );
1361
1362 let mutation = Some(Command::source_change_split_to_mutation(
1363 &split_state.split_assignment,
1364 ));
1365 let (table_ids, node_actors) = self.collect_base_info();
1366 (
1367 mutation,
1368 table_ids,
1369 None,
1370 node_actors,
1371 PostCollectCommand::SourceChangeSplit {
1372 split_assignment: split_state.split_assignment,
1373 },
1374 )
1375 }
1376
1377 Some(Command::CreateSubscription {
1378 subscription_id,
1379 upstream_mv_table_id,
1380 retention_second,
1381 }) => {
1382 self.database_info.register_subscriber(
1383 upstream_mv_table_id.as_job_id(),
1384 subscription_id.as_subscriber_id(),
1385 SubscriberType::Subscription(retention_second),
1386 );
1387 let mutation = Some(Command::create_subscription_to_mutation(
1388 upstream_mv_table_id,
1389 subscription_id,
1390 ));
1391 let (table_ids, node_actors) = self.collect_base_info();
1392 (
1393 mutation,
1394 table_ids,
1395 None,
1396 node_actors,
1397 PostCollectCommand::CreateSubscription { subscription_id },
1398 )
1399 }
1400
1401 Some(Command::DropSubscription {
1402 subscription_id,
1403 upstream_mv_table_id,
1404 }) => {
1405 if self
1406 .database_info
1407 .unregister_subscriber(
1408 upstream_mv_table_id.as_job_id(),
1409 subscription_id.as_subscriber_id(),
1410 )
1411 .is_none()
1412 {
1413 warn!(%subscription_id, %upstream_mv_table_id, "no subscription to drop");
1414 }
1415 let mutation = Some(Command::drop_subscription_to_mutation(
1416 upstream_mv_table_id,
1417 subscription_id,
1418 ));
1419 let (table_ids, node_actors) = self.collect_base_info();
1420 (
1421 mutation,
1422 table_ids,
1423 None,
1424 node_actors,
1425 PostCollectCommand::Command("DropSubscription".to_owned()),
1426 )
1427 }
1428
1429 Some(Command::AlterSubscriptionRetention {
1430 subscription_id,
1431 upstream_mv_table_id,
1432 retention_second,
1433 }) => {
1434 self.database_info.update_subscription_retention(
1435 upstream_mv_table_id.as_job_id(),
1436 subscription_id.as_subscriber_id(),
1437 retention_second,
1438 );
1439 self.apply_simple_command(None, "AlterSubscriptionRetention")
1440 }
1441
1442 Some(Command::ConnectorPropsChange(config)) => {
1443 let mutation = Some(Command::connector_props_change_to_mutation(&config));
1444 let (table_ids, node_actors) = self.collect_base_info();
1445 (
1446 mutation,
1447 table_ids,
1448 None,
1449 node_actors,
1450 PostCollectCommand::ConnectorPropsChange(config),
1451 )
1452 }
1453
1454 Some(Command::Refresh {
1455 table_id,
1456 associated_source_id,
1457 staging_table_id,
1458 trigger_time,
1459 }) => {
1460 let mutation = Some(Command::refresh_to_mutation(table_id, associated_source_id));
1461 let actors = RefreshCycleActors::from_fragments(
1462 self.database_info.job_fragment_infos(table_id.as_job_id()),
1463 );
1464 let (table_ids, node_actors) = self.collect_base_info();
1465 (
1466 mutation,
1467 table_ids,
1468 None,
1469 node_actors,
1470 PostCollectCommand::RefreshStarted {
1471 table_id,
1472 database_id: self.database_info.database_id,
1473 associated_source_id,
1474 staging_table_id,
1475 trigger_time,
1476 actors,
1477 },
1478 )
1479 }
1480
1481 Some(Command::ListFinish {
1482 table_id: _,
1483 associated_source_id,
1484 }) => {
1485 let mutation = Some(Command::list_finish_to_mutation(associated_source_id));
1486 self.apply_simple_command(mutation, "ListFinish")
1487 }
1488
1489 Some(Command::LoadFinish {
1490 table_id: _,
1491 associated_source_id,
1492 }) => {
1493 let mutation = Some(Command::load_finish_to_mutation(associated_source_id));
1494 self.apply_simple_command(mutation, "LoadFinish")
1495 }
1496
1497 Some(Command::FinishRefresh {
1498 table_id,
1499 staging_table_id,
1500 trigger_time,
1501 }) => {
1502 let (table_ids, node_actors) = self.collect_base_info();
1503 (
1504 None,
1505 table_ids,
1506 None,
1507 node_actors,
1508 PostCollectCommand::FinishRefresh {
1509 table_id,
1510 staging_table_id,
1511 trigger_time,
1512 },
1513 )
1514 }
1515
1516 Some(Command::ResetSource { source_id }) => {
1517 let mutation = Some(Command::reset_source_to_mutation(source_id));
1518 self.apply_simple_command(mutation, "ResetSource")
1519 }
1520
1521 Some(Command::ResumeBackfill { target }) => {
1522 let mutation = Command::resume_backfill_to_mutation(&target, &self.database_info)?;
1523 let (table_ids, node_actors) = self.collect_base_info();
1524 (
1525 mutation,
1526 table_ids,
1527 None,
1528 node_actors,
1529 PostCollectCommand::ResumeBackfill { target },
1530 )
1531 }
1532
1533 Some(Command::InjectSourceOffsets {
1534 source_id,
1535 split_offsets,
1536 }) => {
1537 let mutation = Some(Command::inject_source_offsets_to_mutation(
1538 source_id,
1539 &split_offsets,
1540 ));
1541 self.apply_simple_command(mutation, "InjectSourceOffsets")
1542 }
1543 };
1544
1545 let mut finished_snapshot_backfill_jobs = HashSet::new();
1546 let mut mutation = match mutation {
1547 Some(mutation) => Some(mutation),
1548 None => {
1549 let mut finished_snapshot_backfill_job_info = HashMap::new();
1550 if barrier_info.kind.is_checkpoint() {
1551 for (&job_id, job) in &mut self.independent_checkpoint_job_controls {
1552 if let Some(IndependentCheckpointJob::CreatingStreamingJob(creating_job)) =
1553 job.running_mut()
1554 && creating_job.should_merge_to_upstream(partial_graph_manager)
1555 {
1556 if throttle_config
1560 .as_mut()
1561 .and_then(|config| creating_job.pre_apply_throttle(config))
1562 .is_some()
1563 {
1564 notify_database_graph = true;
1565 }
1566 let info = creating_job
1567 .start_consume_upstream(partial_graph_manager, &barrier_info)?;
1568 finished_snapshot_backfill_job_info
1569 .try_insert(job_id, info)
1570 .expect("non-duplicated");
1571 }
1572 }
1573 }
1574
1575 if !finished_snapshot_backfill_job_info.is_empty() {
1576 let actors_to_create = actors_to_create.get_or_insert_default();
1577 let mut mutation = PbUpdateMutation::default();
1578 for (job_id, info) in finished_snapshot_backfill_job_info {
1579 finished_snapshot_backfill_jobs.insert(job_id);
1580 mutation.subscriptions_to_drop.extend(
1581 info.snapshot_backfill_upstream_tables.iter().map(
1582 |upstream_table_id| PbSubscriptionUpstreamInfo {
1583 subscriber_id: job_id.as_subscriber_id(),
1584 upstream_mv_table_id: *upstream_table_id,
1585 },
1586 ),
1587 );
1588 for upstream_mv_table_id in &info.snapshot_backfill_upstream_tables {
1589 assert_matches!(
1590 self.database_info.unregister_subscriber(
1591 upstream_mv_table_id.as_job_id(),
1592 job_id.as_subscriber_id()
1593 ),
1594 Some(SubscriberType::SnapshotBackfill)
1595 );
1596 }
1597
1598 table_ids_to_commit.extend(
1599 info.fragment_infos
1600 .values()
1601 .flat_map(|fragment| fragment.state_table_ids.iter())
1602 .copied(),
1603 );
1604
1605 let actor_len = info
1606 .fragment_infos
1607 .values()
1608 .map(|fragment| fragment.actors.len() as u64)
1609 .sum();
1610 let id_gen = GlobalActorIdGen::new(
1611 partial_graph_manager
1612 .control_stream_manager()
1613 .env
1614 .actor_id_generator(),
1615 actor_len,
1616 );
1617 let mut next_local_actor_id = 0;
1618 let actor_mapping: HashMap<_, _> = info
1620 .fragment_infos
1621 .values()
1622 .flat_map(|fragment| fragment.actors.keys())
1623 .map(|old_actor_id| {
1624 let new_actor_id = id_gen.to_global_id(next_local_actor_id);
1625 next_local_actor_id += 1;
1626 (*old_actor_id, new_actor_id.as_global_id())
1627 })
1628 .collect();
1629 let actor_mapping = &actor_mapping;
1630 let database_partial_graph_id = to_partial_graph_id(self.database_id, None);
1633 let mut edge_builder = FragmentEdgeBuilder::new()
1634 .add_existing_fragments(
1635 info.upstream_fragment_downstreams.keys().map(
1636 |upstream_fragment_id| {
1637 self.database_info.fragment(*upstream_fragment_id)
1638 },
1639 ),
1640 database_partial_graph_id,
1641 partial_graph_manager.control_stream_manager(),
1642 )
1643 .add_existing_fragments(
1644 info.fragment_infos.values(),
1645 to_partial_graph_id(self.database_id, Some(job_id)),
1646 partial_graph_manager.control_stream_manager(),
1647 );
1648 let new_stream_actors: HashMap<_, _> = info
1649 .stream_actors
1650 .into_iter()
1651 .map(|(old_actor_id, mut actor)| {
1652 let new_actor_id = actor_mapping[&old_actor_id];
1653 actor.actor_id = new_actor_id;
1654 (new_actor_id, actor)
1655 })
1656 .collect();
1657 let new_fragment_info: HashMap<_, _> = info
1658 .fragment_infos
1659 .into_iter()
1660 .map(|(fragment_id, mut fragment)| {
1661 let actors = take(&mut fragment.actors);
1662 fragment.actors = actors
1663 .into_iter()
1664 .map(|(old_actor_id, actor)| {
1665 let new_actor_id = actor_mapping[&old_actor_id];
1666 (new_actor_id, actor)
1667 })
1668 .collect();
1669 (fragment_id, fragment)
1670 })
1671 .collect();
1672 mutation.actor_splits.extend(
1673 new_fragment_info
1674 .values()
1675 .flat_map(|fragment| &fragment.actors)
1676 .map(|(actor_id, actor)| {
1677 (
1678 *actor_id,
1679 ConnectorSplits {
1680 splits: actor
1681 .splits
1682 .iter()
1683 .map(ConnectorSplit::from)
1684 .collect(),
1685 },
1686 )
1687 }),
1688 );
1689 edge_builder = edge_builder.replace_existing_fragment_actors(
1691 new_fragment_info.values(),
1692 database_partial_graph_id,
1693 partial_graph_manager.control_stream_manager(),
1694 );
1695 let (mut edges, _) = edge_builder
1696 .finish_fragments()
1697 .add_relations(&info.upstream_fragment_downstreams)?
1698 .add_relations(&info.downstreams)?
1699 .build();
1700 let new_actors_to_create = edges.collect_actors_to_create(
1701 new_fragment_info.values().map(|fragment| {
1702 (
1703 fragment.fragment_id,
1704 &fragment.nodes,
1705 fragment.actors.iter().map(|(actor_id, actor)| {
1706 (&new_stream_actors[actor_id], actor.worker_id)
1707 }),
1708 [], )
1710 }),
1711 );
1712 edges.apply_to_update_mutation(&mut mutation);
1713 for (worker_id, worker_actors) in new_actors_to_create {
1714 node_actors.entry(worker_id).or_default().extend(
1715 worker_actors.values().flat_map(|(_, actors, _)| {
1716 actors.iter().map(|(actor, _, _)| actor.actor_id)
1717 }),
1718 );
1719 actors_to_create
1720 .entry(worker_id)
1721 .or_default()
1722 .extend(worker_actors);
1723 }
1724 self.database_info.add_existing(InflightStreamingJobInfo {
1725 job_id,
1726 fragment_infos: new_fragment_info,
1727 subscribers: Default::default(), status: CreateStreamingJobStatus::Created,
1729 cdc_table_backfill_tracker: None, });
1731 }
1732 Some(PbMutation::Update(mutation))
1733 } else {
1734 let fragment_ids = self.database_info.take_pending_backfill_nodes();
1735 if fragment_ids.is_empty() {
1736 None
1737 } else {
1738 Some(PbMutation::StartFragmentBackfill(
1739 PbStartFragmentBackfillMutation { fragment_ids },
1740 ))
1741 }
1742 }
1743 }
1744 };
1745
1746 if matches!(
1747 mutation,
1748 None | Some(PbMutation::Update(_)) | Some(PbMutation::DropSubscriptions(_))
1749 ) && !self
1750 .pending_independent_job_subscriptions_to_drop
1751 .is_empty()
1752 {
1753 let subscriptions_to_drop = self.take_pending_independent_job_subscriptions_to_drop();
1754 if !subscriptions_to_drop.is_empty() {
1755 match &mut mutation {
1756 None => {
1757 mutation =
1758 Some(PbMutation::DropSubscriptions(PbDropSubscriptionsMutation {
1759 info: subscriptions_to_drop,
1760 }));
1761 }
1762 Some(PbMutation::Update(update)) => {
1763 update.subscriptions_to_drop.extend(subscriptions_to_drop);
1764 }
1765 Some(PbMutation::DropSubscriptions(drop_subscriptions)) => {
1766 drop_subscriptions.info.extend(subscriptions_to_drop);
1767 }
1768 Some(_) => unreachable!("checked compatible mutation above"),
1769 }
1770 }
1771 }
1772
1773 for (job_id, job) in &mut self.independent_checkpoint_job_controls {
1775 let Some(job) = job.running_mut() else {
1776 continue;
1777 };
1778 if finished_snapshot_backfill_jobs.contains(job_id)
1779 && matches!(job, IndependentCheckpointJob::CreatingStreamingJob(_))
1780 {
1781 continue;
1782 }
1783 let throttle_mutation = throttle_config.as_mut().and_then(|config| {
1784 job.pre_apply_throttle(config)
1785 .map(|mutation| (mutation, notifier.as_mut()))
1786 });
1787 job.on_new_upstream_barrier(partial_graph_manager, &barrier_info, throttle_mutation)?;
1788 }
1789
1790 let database_notifier = if notify_database_graph {
1791 notifier.as_mut()
1792 } else {
1793 None
1794 };
1795 partial_graph_manager.inject_barrier(
1796 to_partial_graph_id(self.database_id, None),
1797 mutation,
1798 &node_actors,
1799 InflightFragmentInfo::existing_table_ids(self.database_info.fragment_infos()),
1800 InflightFragmentInfo::workers(self.database_info.fragment_infos()),
1801 actors_to_create,
1802 PartialGraphBarrierInfo::new(
1803 post_collect_command,
1804 barrier_info,
1805 database_notifier,
1806 table_ids_to_commit,
1807 ),
1808 )?;
1809
1810 if let Some(notifier) = notifier.take() {
1813 notifier.started();
1814 }
1815
1816 Ok(ApplyCommandInfo {
1817 jobs_to_wait: finished_snapshot_backfill_jobs,
1818 })
1819 }
1820}