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