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