1use 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 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 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 pending_non_checkpoint_barriers: Vec<u64>,
120 },
121 ConsumingLogStore {
125 tracking_job: TrackingJob,
126 info: CreatingJobInfo,
127 log_store_progress_tracker: CreateMviewLogStoreProgressTracker,
128 pending_barriers: VecDeque<BarrierInfo>,
129 },
130 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>, ) -> 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 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}