Skip to main content

risingwave_meta/barrier/checkpoint/
state.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
76/// The latest state of `GlobalBarrierWorker` after injecting the latest barrier.
77pub(in crate::barrier) struct BarrierWorkerState {
78    /// The last sent `prev_epoch`
79    ///
80    /// There's no need to persist this field. On recovery, we will restore this from the latest
81    /// committed snapshot in `HummockManager`.
82    in_flight_prev_epoch: TracedEpoch,
83
84    /// The `prev_epoch` of pending non checkpoint barriers
85    pending_non_checkpoint_barriers: Vec<u64>,
86
87    /// Whether the cluster is paused.
88    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    /// Returns the `BarrierInfo` for the next barrier, and updates the state.
128    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
161/// Result tuple of `apply_command`: mutation, table IDs to commit, actors to create,
162/// node actors, and post-collect command.
163type ApplyCommandResult = (
164    Option<Mutation>,
165    HashSet<TableId>,
166    Option<StreamJobActorsToCreate>,
167    HashMap<WorkerId, HashSet<ActorId>>,
168    PostCollectCommand,
169);
170
171/// Result of actor rendering for a create/replace streaming job.
172pub(crate) struct RenderResult {
173    /// Rendered actors grouped by fragment.
174    pub stream_actors: HashMap<FragmentId, Vec<StreamActor>>,
175    /// Worker placement for each actor.
176    pub actor_location: HashMap<ActorId, WorkerId>,
177}
178
179/// Derive `NoShuffle` edges from fragment downstream relations and resolve ensembles.
180///
181/// This scans both the internal downstream relations (`fragments.downstreams`) and
182/// the cross-boundary upstream-to-new-fragment relations (`upstream_fragment_downstreams`)
183/// to find all `NoShuffle` edges. It then runs BFS to find connected components (ensembles)
184/// and categorizes them into:
185/// - Ensembles whose entry fragments include existing (non-new) fragments
186/// - Ensembles whose entry fragments are all newly created
187pub(crate) fn resolve_no_shuffle_ensembles(
188    fragments: &StreamJobFragmentsToCreate,
189    upstream_fragment_downstreams: &FragmentDownstreamRelation,
190) -> MetaResult<Vec<NoShuffleEnsemble>> {
191    // Derive FragmentNewNoShuffle from the two downstream relation maps.
192    let mut new_no_shuffle: HashMap<_, HashSet<_>> = HashMap::new();
193
194    // Internal edges (new → new) and edges from new → existing downstream (replace job).
195    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    // Cross-boundary edges: existing upstream → new downstream.
207    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        // Flatten into directed edge pairs for BFS.
222        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    // Add standalone fragments (not covered by any ensemble) as single-fragment ensembles.
243    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
256/// Render actors for a create or replace streaming job.
257///
258/// This determines the parallelism for each no-shuffle ensemble (either from an existing
259/// inflight upstream or computed fresh), and produces `StreamActor` instances with worker
260/// placements and actor-level no-shuffle mappings.
261///
262/// The process follows three steps:
263/// 1. For each ensemble, resolve `EnsembleActorTemplate` (from existing or fresh).
264/// 2. For each new component fragment, allocate actor IDs and compute worker/vnode assignments.
265/// 3. Expand the simple assignments into full `StreamActor` structures.
266pub(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    // Step 2: Render actors for each ensemble.
278    // For each new fragment, produce a simple assignment: actor_id -> (worker_id, vnode_bitmap).
279    let mut actor_assignments: HashMap<FragmentId, HashMap<ActorId, (WorkerId, Option<Bitmap>)>> =
280        HashMap::new();
281
282    for ensemble in ensembles {
283        // Determine the EnsembleActorTemplate for this ensemble.
284        //
285        // Check if any component fragment in the ensemble already exists (i.e. is inflight).
286        // If so, derive the actor assignment from an existing fragment. Otherwise render fresh.
287        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            // Sanity check: all existing fragments in the same ensemble must be aligned —
298            // same actor count and same worker placement per vnode.
299            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            // All fragments are new — render from scratch.
309            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            // Assert all component fragments in this ensemble share the same vnode count.
318            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        // Render each new component fragment in this ensemble.
343        for fragment_id in ensemble.component_fragments() {
344            if !fragments.inner.fragments.contains_key(&fragment_id) {
345                continue; // Skip existing fragments.
346            }
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    // Step 3: Expand simple assignments into full StreamActor structures.
357    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                    // The upstream job and its materialize executors have already been removed, so
418                    // there is no subscriber left to unregister or notify.
419                    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    /// Collect table IDs to commit and actor IDs to collect from current fragment infos.
432    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    /// Helper for the simplest command variants: those that only need a
440    /// pre-computed mutation and a command name, with no actors to create
441    /// and no additional side effects on `self`.
442    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    /// Returns the inflight actor infos that have included the newly added actors in the given command. The dropped actors
458    /// will be removed from the state after the info get resolved.
459    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        /// Resolve source splits for a create streaming job command.
489        ///
490        /// Combines source fragment split resolution and backfill split alignment
491        /// into one step, looking up existing upstream actor splits from the inflight database info.
492        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        // Each variant handles its own pre-apply, edge building, mutation generation,
531        // collect base info, and post-apply. The match produces values consumed by the
532        // common snapshot-backfill-merging code that follows.
533        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                    // set snapshot epoch of upstream table for snapshot backfill
591                    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                    // Build edges first (needed for no-shuffle mapping used in split resolution)
615                    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                    // Phase 2: Resolve source-level DiscoveredSplits to actor-level SplitAssignment
624                    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                    // 1. Fill snapshot backfill epochs.
846                    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                    // 2. Build BatchRefreshLogicalFragments (after epoch filling).
871                    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                    // 3. Create BatchRefreshJobCheckpointControl. `new()` handles actor
901                    //    rendering, the partial-graph initial barrier, and produces the
902                    //    database-graph mutation for the main barrier.
903                    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                    // Database-graph `Add` mutation: batch refresh has no actors in the
913                    // database graph; it only needs to register snapshot-backfill
914                    // subscribers on the upstream MV tables.
915                    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                    // Register permanent subscriber (never unregistered until MV is dropped)
969                    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                // Build edges
1022                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                // Phase 2: Resolve source-level DiscoveredSplits to actor-level SplitAssignment
1038                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                // Pre-apply: add new job and fragments
1056                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                // Actors to create
1110                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                // CDC table snapshot splits
1118                let actor_cdc_table_snapshot_splits = self
1119                    .database_info
1120                    .assign_cdc_backfill_splits(info.stream_job_fragments.stream_job_id())?;
1121
1122                // Mutation
1123                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                // pre_apply: drop node upstream for sink targets
1192                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                // Pre-apply: reschedule fragments
1238                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                // Actors to create
1283                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                // Post-apply: remove old actors
1291                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                // Mutation
1300                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                    "", // replace jobs don't need mview definition
1328                    &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                // Render actors for auto_refresh_schema_sinks.
1340                // Each sink's new_fragment inherits parallelism from its original_fragment.
1341                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                // Build edges first (needed for no-shuffle mapping used in split resolution)
1378                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                // Phase 2: Resolve splits to actor-level assignment.
1388                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                // Pre-apply: add new fragments and replace upstream
1424                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                // Actors to create
1456                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                // Mutation (must be generated before removing old fragments,
1465                // because it reads actor info from database_info)
1466                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                // Post-apply: remove old fragments
1474                {
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                // Pre-apply: split assignments
1503                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                            // The independent actors will stop on this barrier. Apply throttle to
1705                            // the in-memory plan used to create the database-graph actors, and let
1706                            // the database barrier own the collection notification.
1707                            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                        // mapping from old_actor_id to new_actor_id
1767                        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                        // Capture the old actor layouts before remapping the inflight fragment
1779                        // information to fresh actor IDs below.
1780                        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                        // new actors belong to the database partial graph
1838                        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                                    [], // no initial subscriber for backfilling job
1857                                )
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(), // no initial subscribers for newly created snapshot backfill
1876                            status: CreateStreamingJobStatus::Created,
1877                            cdc_table_backfill_tracker: None, // no cdc table backfill for snapshot backfill
1878                        });
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        // Forward barrier to independent job controls
1922        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        // Publish the collection receivers only after all parts of a scheduled command have been
1959        // dispatched successfully. Periodic barriers do not have a notifier.
1960        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}