Skip to main content

risingwave_meta/barrier/checkpoint/independent_job/creating_job/
status.rs

1// Copyright 2026 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::collections::hash_map::Entry;
16use std::collections::{HashMap, HashSet, VecDeque};
17use std::mem::{replace, take};
18use std::time::Duration;
19
20use itertools::Itertools;
21use risingwave_common::hash::ActorId;
22use risingwave_common::util::epoch::Epoch;
23use risingwave_pb::hummock::HummockVersionStats;
24use risingwave_pb::id::{FragmentId, PartialGraphId};
25use risingwave_pb::stream_plan::barrier::PbBarrierKind;
26use risingwave_pb::stream_plan::barrier_mutation::Mutation;
27use risingwave_pb::stream_plan::{PbStreamNode, StartFragmentBackfillMutation};
28use risingwave_pb::stream_service::barrier_complete_response::{
29    CreateMviewProgress, PbCreateMviewProgress,
30};
31use tracing::warn;
32
33use crate::barrier::checkpoint::independent_job::creating_job::CreatingJobInfo;
34use crate::barrier::notifier::Notifier;
35use crate::barrier::partial_graph::PartialGraphManager;
36use crate::barrier::progress::{CreateMviewProgressTracker, TrackingJob};
37use crate::barrier::{BarrierInfo, BarrierKind, TracedEpoch};
38use crate::controller::fragment::InflightFragmentInfo;
39
40#[derive(Debug)]
41pub(super) struct CreateMviewLogStoreProgressTracker {
42    /// `actor_id` -> `pending_epoch_lag`
43    ongoing_actors: HashMap<ActorId, u64>,
44    finished_actors: HashSet<ActorId>,
45}
46
47impl CreateMviewLogStoreProgressTracker {
48    pub(super) fn new(actors: impl Iterator<Item = ActorId>, pending_barrier_lag: u64) -> Self {
49        Self {
50            ongoing_actors: HashMap::from_iter(actors.map(|actor| (actor, pending_barrier_lag))),
51            finished_actors: HashSet::new(),
52        }
53    }
54
55    pub(super) fn gen_backfill_progress(&self) -> String {
56        let sum = self.ongoing_actors.values().sum::<u64>() as f64;
57        let count = if self.ongoing_actors.is_empty() {
58            1
59        } else {
60            self.ongoing_actors.len()
61        } as f64;
62        let avg = sum / count;
63        let avg_lag_time = Duration::from_millis(Epoch(avg as _).physical_time());
64        format!(
65            "actor: {}/{}, avg lag {:?}",
66            self.finished_actors.len(),
67            self.ongoing_actors.len() + self.finished_actors.len(),
68            avg_lag_time
69        )
70    }
71
72    fn update(&mut self, progress: impl IntoIterator<Item = &PbCreateMviewProgress>) {
73        for progress in progress {
74            match self.ongoing_actors.entry(progress.backfill_actor_id) {
75                Entry::Occupied(mut entry) => {
76                    if progress.done {
77                        entry.remove_entry();
78                        assert!(
79                            self.finished_actors.insert(progress.backfill_actor_id),
80                            "non-duplicate"
81                        );
82                    } else {
83                        *entry.get_mut() = progress.pending_epoch_lag as _;
84                    }
85                }
86                Entry::Vacant(_) => {
87                    if cfg!(debug_assertions) {
88                        panic!(
89                            "reporting progress on non-inflight actor: {:?} {:?}",
90                            progress, self
91                        );
92                    } else {
93                        warn!(?progress, progress_tracker = ?self, "reporting progress on non-inflight actor");
94                    }
95                }
96            }
97        }
98    }
99
100    pub(super) fn is_finished(&self) -> bool {
101        self.ongoing_actors.is_empty()
102    }
103}
104
105#[derive(Debug)]
106pub(super) enum CreatingStreamingJobStatus {
107    /// The creating job is consuming upstream snapshot.
108    /// Will transit to `ConsumingLogStore` on `update_progress` when
109    /// the snapshot has been fully consumed after `update_progress`.
110    ConsumingSnapshot {
111        prev_epoch_fake_physical_time: u64,
112        pending_upstream_barriers: Vec<BarrierInfo>,
113        version_stats: HummockVersionStats,
114        create_mview_tracker: CreateMviewProgressTracker,
115        snapshot_backfill_actors: HashSet<ActorId>,
116        snapshot_epoch: u64,
117        info: CreatingJobInfo,
118        /// The `prev_epoch` of pending non checkpoint barriers
119        pending_non_checkpoint_barriers: Vec<u64>,
120    },
121    /// The creating job is consuming log store.
122    ///
123    /// Will transit to `Finishing` on `on_new_upstream_epoch` when `start_consume_upstream` is `true`.
124    ConsumingLogStore {
125        tracking_job: TrackingJob,
126        info: CreatingJobInfo,
127        log_store_progress_tracker: CreateMviewLogStoreProgressTracker,
128        pending_barriers: VecDeque<BarrierInfo>,
129    },
130    /// All backfill actors have started consuming upstream, and the job
131    /// will be finished when all previously injected barriers have been collected
132    /// Store the `prev_epoch` that will finish at.
133    Finishing(u64, TrackingJob),
134    Resetting(Vec<Notifier>),
135    PlaceHolder,
136}
137
138impl CreatingStreamingJobStatus {
139    pub(super) fn update_progress(
140        &mut self,
141        create_mview_progress: impl IntoIterator<Item = &CreateMviewProgress>,
142    ) {
143        match self {
144            &mut Self::ConsumingSnapshot {
145                ref mut create_mview_tracker,
146                ref version_stats,
147                ref mut prev_epoch_fake_physical_time,
148                ref mut pending_upstream_barriers,
149                ref mut pending_non_checkpoint_barriers,
150                ref snapshot_epoch,
151                ..
152            } => {
153                for progress in create_mview_progress {
154                    create_mview_tracker.apply_progress(progress, version_stats);
155                }
156                if create_mview_tracker.is_finished() {
157                    pending_non_checkpoint_barriers.push(*snapshot_epoch);
158
159                    let prev_epoch = Epoch::from_physical_time(*prev_epoch_fake_physical_time);
160                    let pending_barriers: VecDeque<_> = [BarrierInfo {
161                        curr_epoch: TracedEpoch::new(Epoch(*snapshot_epoch)),
162                        prev_epoch: TracedEpoch::new(prev_epoch),
163                        kind: BarrierKind::Checkpoint(take(pending_non_checkpoint_barriers)),
164                    }]
165                    .into_iter()
166                    .chain(pending_upstream_barriers.drain(..))
167                    .collect();
168
169                    let CreatingStreamingJobStatus::ConsumingSnapshot {
170                        create_mview_tracker,
171                        info,
172                        snapshot_epoch,
173                        snapshot_backfill_actors,
174                        ..
175                    } = replace(self, CreatingStreamingJobStatus::PlaceHolder)
176                    else {
177                        unreachable!()
178                    };
179
180                    let tracking_job = create_mview_tracker.into_tracking_job();
181
182                    *self = CreatingStreamingJobStatus::ConsumingLogStore {
183                        tracking_job,
184                        info,
185                        log_store_progress_tracker: CreateMviewLogStoreProgressTracker::new(
186                            snapshot_backfill_actors.iter().cloned(),
187                            pending_barriers
188                                .back()
189                                .map(|barrier_info| {
190                                    barrier_info.prev_epoch().saturating_sub(snapshot_epoch)
191                                })
192                                .unwrap_or(0),
193                        ),
194                        pending_barriers,
195                    };
196                }
197            }
198            CreatingStreamingJobStatus::ConsumingLogStore {
199                log_store_progress_tracker,
200                ..
201            } => {
202                log_store_progress_tracker.update(create_mview_progress);
203            }
204            CreatingStreamingJobStatus::Finishing(..)
205            | CreatingStreamingJobStatus::Resetting(..) => {}
206            CreatingStreamingJobStatus::PlaceHolder => {
207                unreachable!()
208            }
209        }
210    }
211
212    pub(super) fn start_consume_upstream(&mut self, barrier_info: &BarrierInfo) -> CreatingJobInfo {
213        match self {
214            CreatingStreamingJobStatus::ConsumingSnapshot { .. } => {
215                unreachable!(
216                    "should not start consuming upstream for a job that are consuming snapshot"
217                )
218            }
219            CreatingStreamingJobStatus::ConsumingLogStore { .. } => {
220                let prev_epoch = barrier_info.prev_epoch();
221                {
222                    assert!(barrier_info.kind.is_checkpoint());
223                    let CreatingStreamingJobStatus::ConsumingLogStore {
224                        info, tracking_job, ..
225                    } = replace(self, CreatingStreamingJobStatus::PlaceHolder)
226                    else {
227                        unreachable!()
228                    };
229                    *self = CreatingStreamingJobStatus::Finishing(prev_epoch, tracking_job);
230                    info
231                }
232            }
233            CreatingStreamingJobStatus::Finishing { .. } => {
234                unreachable!("should not start consuming upstream for a job again")
235            }
236            CreatingStreamingJobStatus::Resetting(..) => {
237                unreachable!("unlikely to start consume upstream when resetting")
238            }
239            CreatingStreamingJobStatus::PlaceHolder => {
240                unreachable!()
241            }
242        }
243    }
244
245    pub(super) fn on_new_upstream_epoch(
246        &mut self,
247        partial_graph_manager: &PartialGraphManager,
248        partial_graph_id: PartialGraphId,
249        max_pending_barrier_num: usize,
250        barrier_info: &BarrierInfo,
251        mutation: Option<Mutation>, // mutation to be set for the first barrier to inject
252    ) -> Vec<(BarrierInfo, Option<Mutation>)> {
253        let resolve_initial_barrier_num_to_inject = || {
254            max_pending_barrier_num
255                .saturating_sub(partial_graph_manager.pending_barrier_num(partial_graph_id))
256        };
257        match self {
258            CreatingStreamingJobStatus::ConsumingSnapshot {
259                pending_upstream_barriers,
260                prev_epoch_fake_physical_time,
261                pending_non_checkpoint_barriers,
262                create_mview_tracker,
263                ..
264            } => {
265                let mutation = mutation.or_else(|| {
266                    let pending_backfill_nodes = create_mview_tracker
267                        .take_pending_backfill_nodes()
268                        .collect_vec();
269                    if pending_backfill_nodes.is_empty() {
270                        None
271                    } else {
272                        Some(Mutation::StartFragmentBackfill(
273                            StartFragmentBackfillMutation {
274                                fragment_ids: pending_backfill_nodes,
275                            },
276                        ))
277                    }
278                });
279                let barrier_num_to_inject = resolve_initial_barrier_num_to_inject();
280                pending_upstream_barriers.push(barrier_info.clone());
281                // Mutation barriers must be forwarded even when the partial graph has reached the
282                // configured pending-barrier limit.
283                if barrier_num_to_inject == 0 && mutation.is_none() {
284                    return vec![];
285                }
286                vec![(
287                    CreatingStreamingJobStatus::new_fake_barrier(
288                        prev_epoch_fake_physical_time,
289                        pending_non_checkpoint_barriers,
290                        match barrier_info.kind {
291                            BarrierKind::Barrier => PbBarrierKind::Barrier,
292                            BarrierKind::Checkpoint(_) => PbBarrierKind::Checkpoint,
293                            BarrierKind::Initial => {
294                                unreachable!("upstream new epoch should not be initial")
295                            }
296                        },
297                    ),
298                    mutation,
299                )]
300            }
301            CreatingStreamingJobStatus::ConsumingLogStore {
302                pending_barriers, ..
303            } => drain_pending_barriers(
304                pending_barriers,
305                barrier_info.clone(),
306                resolve_initial_barrier_num_to_inject(),
307            )
308            .into_iter()
309            .map(|barrier_info| (barrier_info, None))
310            .collect(),
311            CreatingStreamingJobStatus::Finishing { .. }
312            | CreatingStreamingJobStatus::Resetting(..) => vec![],
313            CreatingStreamingJobStatus::PlaceHolder => {
314                unreachable!()
315            }
316        }
317    }
318
319    pub(super) fn new_fake_barrier(
320        prev_epoch_fake_physical_time: &mut u64,
321        pending_non_checkpoint_barriers: &mut Vec<u64>,
322        kind: PbBarrierKind,
323    ) -> BarrierInfo {
324        super::super::new_fake_barrier(
325            prev_epoch_fake_physical_time,
326            pending_non_checkpoint_barriers,
327            kind,
328        )
329    }
330
331    pub(super) fn fragment_infos(&self) -> Option<&HashMap<FragmentId, InflightFragmentInfo>> {
332        match self {
333            CreatingStreamingJobStatus::ConsumingSnapshot { info, .. }
334            | CreatingStreamingJobStatus::ConsumingLogStore { info, .. } => {
335                Some(&info.fragment_infos)
336            }
337            CreatingStreamingJobStatus::Finishing(..)
338            | CreatingStreamingJobStatus::Resetting(..) => None,
339            CreatingStreamingJobStatus::PlaceHolder => {
340                unreachable!()
341            }
342        }
343    }
344
345    pub(super) fn pre_apply_throttle<'a>(
346        &mut self,
347        fragment_nodes: impl IntoIterator<Item = (FragmentId, &'a PbStreamNode)>,
348    ) {
349        let fragment_infos = match self {
350            CreatingStreamingJobStatus::ConsumingSnapshot { info, .. }
351            | CreatingStreamingJobStatus::ConsumingLogStore { info, .. } => {
352                &mut info.fragment_infos
353            }
354            CreatingStreamingJobStatus::Finishing(..)
355            | CreatingStreamingJobStatus::Resetting(..) => return,
356            CreatingStreamingJobStatus::PlaceHolder => {
357                unreachable!()
358            }
359        };
360
361        for (fragment_id, stream_node) in fragment_nodes {
362            if let Some(fragment_info) = fragment_infos.get_mut(&fragment_id) {
363                fragment_info.nodes = stream_node.clone();
364            }
365        }
366    }
367}
368
369fn drain_pending_barriers(
370    pending_barriers: &mut VecDeque<BarrierInfo>,
371    new_upstream_barrier: BarrierInfo,
372    barrier_num_to_inject: usize,
373) -> Vec<BarrierInfo> {
374    pending_barriers.push_back(new_upstream_barrier);
375    let barrier_count = pending_barriers.len().min(barrier_num_to_inject);
376    pending_barriers.drain(..barrier_count).collect()
377}
378
379#[cfg(test)]
380mod tests {
381    use super::*;
382
383    fn barrier(prev_epoch: u64, curr_epoch: u64) -> BarrierInfo {
384        BarrierInfo {
385            prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
386            curr_epoch: TracedEpoch::new(Epoch(curr_epoch)),
387            kind: BarrierKind::Barrier,
388        }
389    }
390
391    fn epochs(barriers: &[BarrierInfo]) -> Vec<(u64, u64)> {
392        barriers
393            .iter()
394            .map(|barrier| (barrier.prev_epoch(), barrier.curr_epoch()))
395            .collect()
396    }
397
398    #[test]
399    fn test_drain_pending_barriers_with_available_capacity() {
400        let mut pending_barriers = VecDeque::from([barrier(1, 2), barrier(2, 3), barrier(3, 4)]);
401
402        let injected = drain_pending_barriers(&mut pending_barriers, barrier(4, 5), 0);
403        assert!(injected.is_empty());
404        assert_eq!(
405            epochs(pending_barriers.make_contiguous()),
406            vec![(1, 2), (2, 3), (3, 4), (4, 5)]
407        );
408
409        let injected = drain_pending_barriers(&mut pending_barriers, barrier(5, 6), 2);
410        assert_eq!(epochs(&injected), vec![(1, 2), (2, 3)]);
411        assert_eq!(
412            epochs(pending_barriers.make_contiguous()),
413            vec![(3, 4), (4, 5), (5, 6)]
414        );
415
416        let injected = drain_pending_barriers(&mut pending_barriers, barrier(6, 7), 2);
417        assert_eq!(epochs(&injected), vec![(3, 4), (4, 5)]);
418        assert_eq!(
419            epochs(pending_barriers.make_contiguous()),
420            vec![(5, 6), (6, 7)]
421        );
422
423        let injected = drain_pending_barriers(&mut pending_barriers, barrier(7, 8), 2);
424        assert_eq!(epochs(&injected), vec![(5, 6), (6, 7)]);
425        assert_eq!(epochs(pending_barriers.make_contiguous()), vec![(7, 8)]);
426    }
427
428    #[test]
429    fn test_drain_pending_barriers_without_backlog() {
430        let mut pending_barriers = VecDeque::new();
431
432        let injected = drain_pending_barriers(&mut pending_barriers, barrier(1, 2), 100);
433
434        assert_eq!(epochs(&injected), vec![(1, 2)]);
435        assert!(pending_barriers.is_empty());
436    }
437
438    #[tokio::test]
439    async fn test_resetting_skips_barrier_capacity_lookup() {
440        let mut status = CreatingStreamingJobStatus::Resetting(vec![]);
441        let partial_graph_manager =
442            PartialGraphManager::uninitialized(crate::manager::MetaSrvEnv::for_test().await);
443
444        let injected = status.on_new_upstream_epoch(
445            &partial_graph_manager,
446            PartialGraphId::new(1),
447            10,
448            &barrier(1, 2),
449            None,
450        );
451
452        assert!(injected.is_empty());
453    }
454}