1mod barrier_control;
16mod status;
17
18use std::cmp::max;
19use std::collections::{HashMap, HashSet, VecDeque, hash_map};
20use std::mem::take;
21use std::ops::Bound::{Excluded, Unbounded};
22use std::time::Duration;
23
24use risingwave_common::catalog::{DatabaseId, TableId};
25use risingwave_common::id::JobId;
26use risingwave_common::metrics::LabelGuardedIntGauge;
27use risingwave_common::util::epoch::Epoch;
28use risingwave_meta_model::WorkerId;
29use risingwave_pb::ddl_service::PbBackfillType;
30use risingwave_pb::hummock::HummockVersionStats;
31use risingwave_pb::id::{ActorId, FragmentId, PartialGraphId};
32use risingwave_pb::stream_plan::barrier::PbBarrierKind;
33use risingwave_pb::stream_plan::barrier_mutation::Mutation;
34use risingwave_pb::stream_plan::{AddMutation, StopMutation};
35use risingwave_pb::stream_service::BarrierCompleteResponse;
36use status::CreatingStreamingJobStatus;
37use tracing::{debug, info};
38
39use super::super::state::RenderResult;
40use super::IndependentCheckpointJobControl;
41use crate::MetaResult;
42use crate::barrier::backfill_order_control::get_nodes_with_backfill_dependencies;
43use crate::barrier::checkpoint::independent_job::creating_job::barrier_control::CreatingStreamingJobBarrierStats;
44use crate::barrier::checkpoint::independent_job::creating_job::status::CreateMviewLogStoreProgressTracker;
45use crate::barrier::command::{
46 PostCollectCommand, TableLogEpochs, ThrottleConfigMap, UpstreamTableLogEpochs,
47};
48use crate::barrier::context::CreateSnapshotBackfillJobCommandInfo;
49use crate::barrier::edge_builder::FragmentEdgeBuildResult;
50use crate::barrier::info::{BarrierInfo, InflightStreamingJobInfo};
51use crate::barrier::notifier::NotifierStarter;
52use crate::barrier::partial_graph::{
53 CollectedBarrier, PartialGraphBarrierInfo, PartialGraphManager, PartialGraphRecoverer,
54};
55use crate::barrier::progress::{CreateMviewProgressTracker, TrackingJob, collect_done_fragments};
56use crate::barrier::rpc::{build_locality_fragment_state_table_mapping, to_partial_graph_id};
57use crate::barrier::{
58 BackfillOrderState, BackfillProgress, BarrierKind, Command, FragmentBackfillProgress,
59 TracedEpoch,
60};
61use crate::controller::fragment::InflightFragmentInfo;
62use crate::manager::MetaOpts;
63use crate::model::{FragmentDownstreamRelation, StreamActor, StreamJobActorsToCreate};
64use crate::rpc::metrics::GLOBAL_META_METRICS;
65use crate::stream::source_manager::SplitAssignment;
66use crate::stream::{ExtendedFragmentBackfillOrder, build_actor_connector_splits};
67
68fn snapshot_backfill_max_pending_barrier_num(opts: &MetaOpts) -> usize {
69 opts.in_flight_barrier_nums
70 .saturating_mul(opts.snapshot_backfill_barrier_amplification_factor.max(1))
71}
72
73#[derive(Debug)]
74pub(crate) struct CreatingJobInfo {
75 pub fragment_infos: HashMap<FragmentId, InflightFragmentInfo>,
76 pub upstream_fragment_downstreams: FragmentDownstreamRelation,
77 pub downstreams: FragmentDownstreamRelation,
78 pub snapshot_backfill_upstream_tables: HashSet<TableId>,
79 pub stream_actors: HashMap<ActorId, StreamActor>,
80}
81
82#[derive(Debug)]
83pub(crate) struct CreatingStreamingJobControl {
84 job_id: JobId,
85 partial_graph_id: PartialGraphId,
86 snapshot_backfill_upstream_tables: HashSet<TableId>,
87 snapshot_epoch: u64,
88
89 node_actors: HashMap<WorkerId, HashSet<ActorId>>,
90 state_table_ids: HashSet<TableId>,
91
92 max_committed_epoch: Option<u64>,
93 status: CreatingStreamingJobStatus,
94 max_lagged_barrier_num: usize,
95 max_pending_barrier_num: usize,
96
97 upstream_lag: LabelGuardedIntGauge,
98}
99
100impl CreatingStreamingJobControl {
101 #[expect(clippy::too_many_arguments)]
102 pub(crate) fn new<'a>(
103 entry: hash_map::VacantEntry<'a, JobId, IndependentCheckpointJobControl>,
104 create_info: CreateSnapshotBackfillJobCommandInfo,
105 notifier: Option<&mut NotifierStarter>,
106 snapshot_backfill_upstream_tables: HashSet<TableId>,
107 snapshot_epoch: u64,
108 since_timestamp_upstream_log_epochs: Option<(&TableLogEpochs, PartialGraphId, u64)>,
109 version_stat: &HummockVersionStats,
110 partial_graph_manager: &mut PartialGraphManager,
111 edges: &mut FragmentEdgeBuildResult,
112 split_assignment: &SplitAssignment,
113 actors: &RenderResult,
114 ) -> MetaResult<&'a mut Self> {
115 let info = create_info.info.clone();
116 let job_id = info.stream_job_fragments.stream_job_id();
117 let database_id = info.streaming_job.database_id();
118 debug!(
119 %job_id,
120 definition = info.definition,
121 "new creating job"
122 );
123 let fragment_infos = info
124 .stream_job_fragments
125 .new_fragment_info(
126 &actors.stream_actors,
127 &actors.actor_location,
128 split_assignment,
129 )
130 .collect();
131 let snapshot_backfill_actors: HashSet<ActorId> =
132 InflightStreamingJobInfo::snapshot_backfill_actor_ids(&fragment_infos).collect();
133 let backfill_nodes_to_pause =
134 get_nodes_with_backfill_dependencies(&info.fragment_backfill_ordering)
135 .into_iter()
136 .collect();
137 let backfill_order_state = BackfillOrderState::new(
138 &info.fragment_backfill_ordering,
139 &fragment_infos,
140 info.locality_fragment_state_table_mapping.clone(),
141 );
142 let create_mview_tracker = CreateMviewProgressTracker::recover(
143 job_id,
144 &fragment_infos,
145 backfill_order_state,
146 version_stat,
147 );
148
149 let actors_to_create = Command::create_streaming_job_actors_to_create(
150 &info,
151 edges,
152 &actors.stream_actors,
153 &actors.actor_location,
154 );
155
156 let mut prev_epoch_fake_physical_time = 0;
157 let mut pending_non_checkpoint_barriers = vec![];
158
159 let (initial_barrier_info, log_store_barriers_to_inject) = if let Some((
160 upstream_log_epochs,
161 upstream_partial_graph_id,
162 new_upstream_barrier_prev_epoch,
163 )) =
164 since_timestamp_upstream_log_epochs
165 {
166 let (initial_barrier, barriers_to_inject) =
167 Self::resolve_since_timestamp_upstream_log_epochs(
168 upstream_log_epochs,
169 partial_graph_manager.pending_barrier_infos(upstream_partial_graph_id),
170 snapshot_epoch,
171 new_upstream_barrier_prev_epoch,
172 )?;
173 (initial_barrier, Some(barriers_to_inject))
174 } else {
175 (
176 CreatingStreamingJobStatus::new_fake_barrier(
177 &mut prev_epoch_fake_physical_time,
178 &mut pending_non_checkpoint_barriers,
179 PbBarrierKind::Checkpoint,
180 ),
181 None,
182 )
183 };
184
185 let added_actors: Vec<ActorId> = actors
186 .stream_actors
187 .values()
188 .flatten()
189 .map(|actor| actor.actor_id)
190 .collect();
191 let actor_splits = split_assignment
192 .values()
193 .flat_map(build_actor_connector_splits)
194 .collect();
195
196 assert!(
197 info.cdc_table_snapshot_splits.is_none(),
198 "should not have cdc backfill for snapshot backfill job"
199 );
200
201 let initial_mutation = Mutation::Add(AddMutation {
202 actor_dispatchers: Default::default(),
204 added_actors,
205 actor_splits,
206 pause: false,
208 subscriptions_to_add: Default::default(),
209 backfill_nodes_to_pause,
210 actor_cdc_table_snapshot_splits: None,
211 new_upstream_sinks: Default::default(),
212 dropped_actors: Default::default(),
213 sink_log_store_flush: Default::default(),
214 });
215
216 let node_actors = InflightFragmentInfo::actor_ids_to_collect(fragment_infos.values());
217 let state_table_ids =
218 InflightFragmentInfo::existing_table_ids(fragment_infos.values()).collect();
219
220 let partial_graph_id = to_partial_graph_id(database_id, Some(job_id));
221 let max_lagged_barrier_num = partial_graph_manager
222 .control_stream_manager()
223 .env
224 .opts
225 .snapshot_backfill_finish_max_lagged_barriers;
226 let opts = &partial_graph_manager.control_stream_manager().env.opts;
227 let max_pending_barrier_num = snapshot_backfill_max_pending_barrier_num(opts);
228
229 let IndependentCheckpointJobControl::CreatingStreamingJob(job) = entry.insert(
230 IndependentCheckpointJobControl::CreatingStreamingJob(Self {
231 partial_graph_id,
232 job_id,
233 snapshot_backfill_upstream_tables,
234 max_committed_epoch: None,
235 snapshot_epoch,
236 status: CreatingStreamingJobStatus::PlaceHolder, max_lagged_barrier_num,
238 max_pending_barrier_num,
239 upstream_lag: GLOBAL_META_METRICS
240 .snapshot_backfill_lag
241 .with_guarded_label_values(&[&format!("{}", job_id)]),
242 node_actors,
243 state_table_ids,
244 }),
245 ) else {
246 unreachable!()
247 };
248
249 let mut graph_adder = partial_graph_manager.add_partial_graph(
250 partial_graph_id,
251 CreatingStreamingJobBarrierStats::new(job_id, snapshot_epoch),
252 );
253
254 if let Err(e) = Self::inject_barrier(
255 partial_graph_id,
256 graph_adder.manager(),
257 &job.node_actors,
258 &job.state_table_ids,
259 false,
260 initial_barrier_info,
261 Some(actors_to_create),
262 Some(initial_mutation),
263 notifier,
264 Some(create_info),
265 ) {
266 graph_adder.failed();
267 job.status = CreatingStreamingJobStatus::Resetting(vec![]);
268 Err(e)
269 } else {
270 graph_adder.added();
271 let job_info = CreatingJobInfo {
272 fragment_infos,
273 upstream_fragment_downstreams: info.upstream_fragment_downstreams.clone(),
274 downstreams: info.stream_job_fragments.downstreams,
275 snapshot_backfill_upstream_tables: job.snapshot_backfill_upstream_tables.clone(),
276 stream_actors: actors
277 .stream_actors
278 .values()
279 .flatten()
280 .map(|actor| (actor.actor_id, actor.clone()))
281 .collect(),
282 };
283 if let Some(log_store_barriers_to_inject) = log_store_barriers_to_inject {
284 let upstream_lag = log_store_barriers_to_inject
285 .last()
286 .map(|info| info.prev_epoch().saturating_sub(snapshot_epoch))
287 .unwrap_or(0);
288 job.status = CreatingStreamingJobStatus::ConsumingLogStore {
289 tracking_job: TrackingJob::recovered(job_id, &job_info.fragment_infos),
290 info: job_info,
291 log_store_progress_tracker: CreateMviewLogStoreProgressTracker::new(
292 snapshot_backfill_actors.iter().cloned(),
293 upstream_lag,
294 ),
295 pending_barriers: log_store_barriers_to_inject.into(),
296 };
297 } else {
298 assert!(pending_non_checkpoint_barriers.is_empty());
299 job.status = CreatingStreamingJobStatus::ConsumingSnapshot {
300 prev_epoch_fake_physical_time,
301 pending_upstream_barriers: vec![],
302 version_stats: version_stat.clone(),
303 create_mview_tracker,
304 snapshot_backfill_actors,
305 snapshot_epoch,
306 info: job_info,
307 pending_non_checkpoint_barriers,
308 };
309 };
310 Ok(job)
311 }
312 }
313
314 pub(super) fn gen_fragment_backfill_progress(&self) -> Vec<FragmentBackfillProgress> {
315 match &self.status {
316 CreatingStreamingJobStatus::ConsumingSnapshot {
317 create_mview_tracker,
318 info,
319 ..
320 } => create_mview_tracker.collect_fragment_progress(&info.fragment_infos, true),
321 CreatingStreamingJobStatus::ConsumingLogStore { info, .. } => {
322 collect_done_fragments(self.job_id, &info.fragment_infos)
323 }
324 CreatingStreamingJobStatus::Finishing(_, _)
325 | CreatingStreamingJobStatus::Resetting(_)
326 | CreatingStreamingJobStatus::PlaceHolder => vec![],
327 }
328 }
329
330 fn resolve_upstream_log_epochs(
331 snapshot_backfill_upstream_tables: &HashSet<TableId>,
332 upstream_table_log_epochs: &UpstreamTableLogEpochs,
333 exclusive_start_log_epoch: u64,
334 upstream_barrier_info: &BarrierInfo,
335 ) -> MetaResult<Vec<BarrierInfo>> {
336 let table_id = snapshot_backfill_upstream_tables
337 .iter()
338 .next()
339 .expect("snapshot backfill job should have upstream");
340 let epochs_iter = if let Some(epochs) = upstream_table_log_epochs.get(table_id) {
341 let mut epochs_iter = epochs.iter();
342 loop {
343 let (_, checkpoint_epoch) =
344 epochs_iter.next().expect("not reach committed epoch yet");
345 if *checkpoint_epoch < exclusive_start_log_epoch {
346 continue;
347 }
348 assert_eq!(*checkpoint_epoch, exclusive_start_log_epoch);
349 break;
350 }
351 epochs_iter
352 } else {
353 assert_eq!(
355 upstream_barrier_info.prev_epoch(),
356 exclusive_start_log_epoch
357 );
358 static EMPTY_VEC: Vec<(Vec<u64>, u64)> = Vec::new();
359 EMPTY_VEC.iter()
360 };
361
362 let mut ret = vec![];
363 let mut prev_epoch = exclusive_start_log_epoch;
364 let mut pending_non_checkpoint_barriers = vec![];
365 for (non_checkpoint_epochs, checkpoint_epoch) in epochs_iter {
366 for (i, epoch) in non_checkpoint_epochs
367 .iter()
368 .chain([checkpoint_epoch])
369 .enumerate()
370 {
371 assert!(*epoch > prev_epoch);
372 pending_non_checkpoint_barriers.push(prev_epoch);
373 ret.push(BarrierInfo {
374 prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
375 curr_epoch: TracedEpoch::new(Epoch(*epoch)),
376 kind: if i == 0 {
377 BarrierKind::Checkpoint(take(&mut pending_non_checkpoint_barriers))
378 } else {
379 BarrierKind::Barrier
380 },
381 });
382 prev_epoch = *epoch;
383 }
384 }
385 ret.push(BarrierInfo {
386 prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
387 curr_epoch: TracedEpoch::new(Epoch(upstream_barrier_info.curr_epoch())),
388 kind: BarrierKind::Checkpoint(pending_non_checkpoint_barriers),
389 });
390 Ok(ret)
391 }
392
393 fn resolve_since_timestamp_upstream_log_epochs(
422 upstream_log_epochs: &TableLogEpochs,
423 pending_upstream_barriers: impl Iterator<Item = &BarrierInfo>,
424 snapshot_epoch: u64,
425 new_upstream_barrier_prev_epoch: u64,
426 ) -> MetaResult<(BarrierInfo, Vec<BarrierInfo>)> {
427 let mut initial_barrier = None;
428 let mut barriers = vec![];
429 fn emit_barrier(
430 initial_barrier: &mut Option<BarrierInfo>,
431 barriers: &mut Vec<BarrierInfo>,
432 barrier: BarrierInfo,
433 ) {
434 if initial_barrier.is_none() {
435 *initial_barrier = Some(barrier);
436 } else {
437 barriers.push(barrier);
438 }
439 }
440
441 let mut prev_epoch = snapshot_epoch;
442 let mut pending_non_checkpoint_barriers = vec![];
443 for (non_checkpoint_epochs, checkpoint_epoch) in upstream_log_epochs {
444 for (i, epoch) in non_checkpoint_epochs
445 .iter()
446 .chain([checkpoint_epoch])
447 .enumerate()
448 {
449 assert!(
450 *epoch > prev_epoch,
451 "changelog epochs should be strictly increasing"
452 );
453 pending_non_checkpoint_barriers.push(prev_epoch);
454 emit_barrier(
455 &mut initial_barrier,
456 &mut barriers,
457 BarrierInfo {
458 prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
459 curr_epoch: TracedEpoch::new(Epoch(*epoch)),
460 kind: if i == 0 {
461 BarrierKind::Checkpoint(take(&mut pending_non_checkpoint_barriers))
462 } else {
463 BarrierKind::Barrier
464 },
465 },
466 );
467 prev_epoch = *epoch;
468 }
469 }
470
471 let mut pending_upstream_barriers = pending_upstream_barriers.peekable();
472 pending_non_checkpoint_barriers.push(prev_epoch);
473 if pending_upstream_barriers.peek().is_none() {
474 assert!(
475 new_upstream_barrier_prev_epoch > prev_epoch,
476 "new upstream barrier prev epoch should be newer than the latest changelog epoch"
477 );
478 emit_barrier(
479 &mut initial_barrier,
480 &mut barriers,
481 BarrierInfo {
482 prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
483 curr_epoch: TracedEpoch::new(Epoch(new_upstream_barrier_prev_epoch)),
484 kind: BarrierKind::Checkpoint(pending_non_checkpoint_barriers),
485 },
486 );
487 } else {
488 let first_pending_barrier = pending_upstream_barriers
489 .peek()
490 .expect("first pending upstream barrier should exist after peek");
491 assert!(
492 first_pending_barrier.prev_epoch() > prev_epoch,
493 "first pending upstream barrier should be newer than the latest resolved changelog epoch"
494 );
495 emit_barrier(
496 &mut initial_barrier,
497 &mut barriers,
498 BarrierInfo {
499 prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
500 curr_epoch: TracedEpoch::new(Epoch(first_pending_barrier.prev_epoch())),
501 kind: BarrierKind::Checkpoint(take(&mut pending_non_checkpoint_barriers)),
502 },
503 );
504 prev_epoch = first_pending_barrier.prev_epoch();
505 for pending_barrier in pending_upstream_barriers {
506 assert_eq!(
507 pending_barrier.prev_epoch(),
508 prev_epoch,
509 "pending upstream barriers should continue from resolved changelog epochs"
510 );
511 pending_non_checkpoint_barriers.push(prev_epoch);
512 emit_barrier(
513 &mut initial_barrier,
514 &mut barriers,
515 BarrierInfo {
516 prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
517 curr_epoch: TracedEpoch::new(Epoch(pending_barrier.curr_epoch())),
518 kind: if pending_barrier.kind.is_checkpoint() {
519 BarrierKind::Checkpoint(take(&mut pending_non_checkpoint_barriers))
520 } else {
521 BarrierKind::Barrier
522 },
523 },
524 );
525 prev_epoch = pending_barrier.curr_epoch();
526 }
527 assert_eq!(
528 new_upstream_barrier_prev_epoch, prev_epoch,
529 "new upstream barrier prev epoch should match the latest pending log-store epoch"
530 );
531 }
532 let Some(initial_barrier) = initial_barrier else {
533 return Err(anyhow::anyhow!(
534 "missing lagging barriers for direct log-store start from snapshot epoch {}",
535 snapshot_epoch
536 )
537 .into());
538 };
539 assert!(initial_barrier.kind.is_checkpoint());
540 Ok((initial_barrier, barriers))
541 }
542
543 fn recover_consuming_snapshot(
544 job_id: JobId,
545 upstream_table_log_epochs: &UpstreamTableLogEpochs,
546 snapshot_epoch: u64,
547 committed_epoch: u64,
548 upstream_barrier_info: &BarrierInfo,
549 info: CreatingJobInfo,
550 backfill_order_state: BackfillOrderState,
551 version_stat: &HummockVersionStats,
552 ) -> MetaResult<(CreatingStreamingJobStatus, BarrierInfo)> {
553 let mut prev_epoch_fake_physical_time = Epoch(committed_epoch).physical_time();
554 let mut pending_non_checkpoint_barriers = vec![];
555 let create_mview_tracker = CreateMviewProgressTracker::recover(
556 job_id,
557 &info.fragment_infos,
558 backfill_order_state,
559 version_stat,
560 );
561 let barrier_info = CreatingStreamingJobStatus::new_fake_barrier(
562 &mut prev_epoch_fake_physical_time,
563 &mut pending_non_checkpoint_barriers,
564 PbBarrierKind::Initial,
565 );
566 Ok((
567 CreatingStreamingJobStatus::ConsumingSnapshot {
568 prev_epoch_fake_physical_time,
569 pending_upstream_barriers: Self::resolve_upstream_log_epochs(
570 &info.snapshot_backfill_upstream_tables,
571 upstream_table_log_epochs,
572 snapshot_epoch,
573 upstream_barrier_info,
574 )?,
575 version_stats: version_stat.clone(),
576 create_mview_tracker,
577 snapshot_backfill_actors: InflightStreamingJobInfo::snapshot_backfill_actor_ids(
578 &info.fragment_infos,
579 )
580 .collect(),
581 info,
582 snapshot_epoch,
583 pending_non_checkpoint_barriers,
584 },
585 barrier_info,
586 ))
587 }
588
589 fn recover_consuming_log_store(
590 job_id: JobId,
591 upstream_table_log_epochs: &UpstreamTableLogEpochs,
592 committed_epoch: u64,
593 upstream_barrier_info: &BarrierInfo,
594 info: CreatingJobInfo,
595 ) -> MetaResult<(CreatingStreamingJobStatus, BarrierInfo)> {
596 let mut pending_barriers: VecDeque<_> = Self::resolve_upstream_log_epochs(
597 &info.snapshot_backfill_upstream_tables,
598 upstream_table_log_epochs,
599 committed_epoch,
600 upstream_barrier_info,
601 )?
602 .into();
603 let mut first_barrier = pending_barriers
604 .pop_front()
605 .expect("resolved upstream log epochs should not be empty");
606 assert!(first_barrier.kind.is_checkpoint());
607 first_barrier.kind = BarrierKind::Initial;
608
609 Ok((
610 CreatingStreamingJobStatus::ConsumingLogStore {
611 tracking_job: TrackingJob::recovered(job_id, &info.fragment_infos),
612 log_store_progress_tracker: CreateMviewLogStoreProgressTracker::new(
613 InflightStreamingJobInfo::snapshot_backfill_actor_ids(&info.fragment_infos),
614 pending_barriers
615 .back()
616 .map(|info| info.prev_epoch() - committed_epoch)
617 .unwrap_or(0),
618 ),
619 pending_barriers,
620 info,
621 },
622 first_barrier,
623 ))
624 }
625
626 #[expect(clippy::too_many_arguments)]
627 pub(crate) fn recover(
628 database_id: DatabaseId,
629 job_id: JobId,
630 snapshot_backfill_upstream_tables: HashSet<TableId>,
631 upstream_table_log_epochs: &UpstreamTableLogEpochs,
632 snapshot_epoch: u64,
633 committed_epoch: u64,
634 upstream_barrier_info: &BarrierInfo,
635 fragment_infos: HashMap<FragmentId, InflightFragmentInfo>,
636 backfill_order: ExtendedFragmentBackfillOrder,
637 fragment_relations: &FragmentDownstreamRelation,
638 version_stat: &HummockVersionStats,
639 new_actors: StreamJobActorsToCreate,
640 initial_mutation: Mutation,
641 partial_graph_recoverer: &mut PartialGraphRecoverer<'_>,
642 ) -> MetaResult<Self> {
643 info!(
644 %job_id,
645 "recovered creating snapshot backfill job"
646 );
647
648 let node_actors = InflightFragmentInfo::actor_ids_to_collect(fragment_infos.values());
649 let state_table_ids: HashSet<_> =
650 InflightFragmentInfo::existing_table_ids(fragment_infos.values()).collect();
651
652 let mut upstream_fragment_downstreams: FragmentDownstreamRelation = Default::default();
653 for (upstream_fragment_id, downstreams) in fragment_relations {
654 if fragment_infos.contains_key(upstream_fragment_id) {
655 continue;
656 }
657 for downstream in downstreams {
658 if fragment_infos.contains_key(&downstream.downstream_fragment_id) {
659 upstream_fragment_downstreams
660 .entry(*upstream_fragment_id)
661 .or_default()
662 .push(downstream.clone());
663 }
664 }
665 }
666 let downstreams = fragment_infos
667 .keys()
668 .filter_map(|fragment_id| {
669 fragment_relations
670 .get(fragment_id)
671 .map(|relation| (*fragment_id, relation.clone()))
672 })
673 .collect();
674
675 let info = CreatingJobInfo {
676 fragment_infos,
677 upstream_fragment_downstreams,
678 downstreams,
679 snapshot_backfill_upstream_tables: snapshot_backfill_upstream_tables.clone(),
680 stream_actors: new_actors
681 .values()
682 .flat_map(|fragments| {
683 fragments.values().flat_map(|(_, actors, _)| {
684 actors
685 .iter()
686 .map(|(actor, _, _)| (actor.actor_id, actor.clone()))
687 })
688 })
689 .collect(),
690 };
691
692 let (status, first_barrier_info) = if committed_epoch < snapshot_epoch {
693 let locality_fragment_state_table_mapping =
694 build_locality_fragment_state_table_mapping(&info.fragment_infos);
695 let backfill_order_state = BackfillOrderState::recover_from_fragment_infos(
696 &backfill_order,
697 &info.fragment_infos,
698 locality_fragment_state_table_mapping,
699 );
700 Self::recover_consuming_snapshot(
701 job_id,
702 upstream_table_log_epochs,
703 snapshot_epoch,
704 committed_epoch,
705 upstream_barrier_info,
706 info,
707 backfill_order_state,
708 version_stat,
709 )?
710 } else {
711 Self::recover_consuming_log_store(
712 job_id,
713 upstream_table_log_epochs,
714 committed_epoch,
715 upstream_barrier_info,
716 info,
717 )?
718 };
719
720 let partial_graph_id = to_partial_graph_id(database_id, Some(job_id));
721 let max_lagged_barrier_num = partial_graph_recoverer
722 .control_stream_manager()
723 .env
724 .opts
725 .snapshot_backfill_finish_max_lagged_barriers;
726 let opts = &partial_graph_recoverer.control_stream_manager().env.opts;
727 let max_pending_barrier_num = snapshot_backfill_max_pending_barrier_num(opts);
728
729 partial_graph_recoverer.recover_graph(
730 partial_graph_id,
731 initial_mutation,
732 &first_barrier_info,
733 &node_actors,
734 state_table_ids.iter().copied(),
735 new_actors,
736 CreatingStreamingJobBarrierStats::new(job_id, snapshot_epoch),
737 )?;
738
739 Ok(Self {
740 job_id,
741 partial_graph_id,
742 snapshot_backfill_upstream_tables,
743 snapshot_epoch,
744 node_actors,
745 state_table_ids,
746 max_committed_epoch: Some(committed_epoch),
747 status,
748 max_lagged_barrier_num,
749 max_pending_barrier_num,
750 upstream_lag: GLOBAL_META_METRICS
751 .snapshot_backfill_lag
752 .with_guarded_label_values(&[&format!("{}", job_id)]),
753 })
754 }
755
756 pub(crate) fn gen_backfill_progress(&self) -> BackfillProgress {
757 let progress = match &self.status {
758 CreatingStreamingJobStatus::ConsumingSnapshot {
759 create_mview_tracker,
760 ..
761 } => {
762 if create_mview_tracker.is_finished() {
763 "Snapshot finished".to_owned()
764 } else {
765 let progress = create_mview_tracker.gen_backfill_progress();
766 format!("Snapshot [{}]", progress)
767 }
768 }
769 CreatingStreamingJobStatus::ConsumingLogStore {
770 log_store_progress_tracker,
771 ..
772 } => {
773 format!(
774 "LogStore [{}]",
775 log_store_progress_tracker.gen_backfill_progress()
776 )
777 }
778 CreatingStreamingJobStatus::Finishing(finish_epoch, ..) => {
779 let committed_epoch = self.max_committed_epoch.expect("should have committed");
780 let lag = Duration::from_millis(
781 Epoch(*finish_epoch).physical_time() - Epoch(committed_epoch).physical_time(),
782 );
783 format!("Finishing [epoch lag: {lag:?}]",)
784 }
785 CreatingStreamingJobStatus::Resetting(_) => "Resetting".to_owned(),
786 CreatingStreamingJobStatus::PlaceHolder => {
787 unreachable!()
788 }
789 };
790 BackfillProgress {
791 progress,
792 backfill_type: PbBackfillType::SnapshotBackfill,
793 }
794 }
795
796 pub(super) fn pinned_upstream_log_epoch(&self) -> (u64, HashSet<TableId>) {
797 (
798 max(self.max_committed_epoch.unwrap_or(0), self.snapshot_epoch),
799 self.snapshot_backfill_upstream_tables.clone(),
800 )
801 }
802
803 fn inject_barrier(
804 partial_graph_id: PartialGraphId,
805 partial_graph_manager: &mut PartialGraphManager,
806 node_actors: &HashMap<WorkerId, HashSet<ActorId>>,
807 state_table_ids: &HashSet<TableId>,
808 is_finishing: bool,
809 barrier_info: BarrierInfo,
810 new_actors: Option<StreamJobActorsToCreate>,
811 mutation: Option<Mutation>,
812 notifier: Option<&mut NotifierStarter>,
813 first_create_info: Option<CreateSnapshotBackfillJobCommandInfo>,
814 ) -> MetaResult<()> {
815 let (table_ids_to_sync, nodes_to_sync_table) = if !is_finishing {
816 (Some(state_table_ids), Some(node_actors.keys().copied()))
817 } else {
818 (None, None)
819 };
820 partial_graph_manager.inject_barrier(
821 partial_graph_id,
822 mutation,
823 None,
824 node_actors,
825 table_ids_to_sync.into_iter().flatten().copied(),
826 nodes_to_sync_table.into_iter().flatten(),
827 new_actors,
828 PartialGraphBarrierInfo::new(
829 first_create_info.map_or_else(
830 PostCollectCommand::barrier,
831 CreateSnapshotBackfillJobCommandInfo::into_post_collect,
832 ),
833 barrier_info,
834 notifier,
835 state_table_ids.clone(),
836 ),
837 )?;
838 Ok(())
839 }
840
841 pub(crate) fn start_consume_upstream(
842 &mut self,
843 partial_graph_manager: &mut PartialGraphManager,
844 barrier_info: &BarrierInfo,
845 ) -> MetaResult<CreatingJobInfo> {
846 info!(
847 job_id = %self.job_id,
848 prev_epoch = barrier_info.prev_epoch(),
849 "start consuming upstream"
850 );
851 let info = self.status.start_consume_upstream(barrier_info);
852 Self::inject_barrier(
853 self.partial_graph_id,
854 partial_graph_manager,
855 &self.node_actors,
856 &self.state_table_ids,
857 true,
858 barrier_info.clone(),
859 None,
860 Some(Mutation::Stop(StopMutation {
861 actors: info
863 .fragment_infos
864 .values()
865 .flat_map(|info| info.actors.keys().copied())
866 .collect(),
867 dropped_sink_fragments: vec![], })),
869 None, None,
871 )?;
872 Ok(info)
873 }
874
875 pub(crate) fn on_new_upstream_barrier(
876 &mut self,
877 partial_graph_manager: &mut PartialGraphManager,
878 barrier_info: &BarrierInfo,
879 mutation: Option<(Mutation, Option<&mut NotifierStarter>)>,
880 ) -> MetaResult<()> {
881 let progress_epoch = if let Some(max_committed_epoch) = self.max_committed_epoch {
882 max(max_committed_epoch, self.snapshot_epoch)
883 } else {
884 self.snapshot_epoch
885 };
886 self.upstream_lag.set(
887 barrier_info
888 .prev_epoch
889 .value()
890 .0
891 .saturating_sub(progress_epoch) as _,
892 );
893 let (mut mutation, mut notifier) = match mutation {
894 Some((mutation, notifier)) => (Some(mutation), notifier),
895 None => (None, None),
896 };
897 for (barrier_to_inject, mutation) in self.status.on_new_upstream_epoch(
898 partial_graph_manager,
899 self.partial_graph_id,
900 self.max_pending_barrier_num,
901 barrier_info,
902 mutation.take(),
903 ) {
904 Self::inject_barrier(
905 self.partial_graph_id,
906 partial_graph_manager,
907 &self.node_actors,
908 &self.state_table_ids,
909 false,
910 barrier_to_inject,
911 None,
912 mutation,
913 notifier.take(),
914 None,
915 )?;
916 }
917 Ok(())
918 }
919
920 pub(crate) fn pre_apply_throttle(
921 &mut self,
922 config: &mut ThrottleConfigMap,
923 ) -> Option<Mutation> {
924 self.status.pre_apply_throttle(config)
925 }
926
927 pub(crate) fn collect(&mut self, collected_barrier: CollectedBarrier<'_>) -> bool {
929 let pending_barrier_num = collected_barrier.pending_barrier_num;
930 self.status.update_progress(
931 collected_barrier
932 .resps
933 .values()
934 .flat_map(|resp| &resp.create_mview_progress),
935 );
936 self.is_ready_to_merge() && pending_barrier_num <= self.max_lagged_barrier_num
937 }
938
939 fn is_ready_to_merge(&self) -> bool {
940 if let CreatingStreamingJobStatus::ConsumingLogStore {
941 log_store_progress_tracker,
942 pending_barriers,
943 ..
944 } = &self.status
945 && pending_barriers.is_empty()
946 && log_store_progress_tracker.is_finished()
947 {
948 true
949 } else {
950 false
951 }
952 }
953
954 pub(crate) fn should_merge_to_upstream(
955 &self,
956 partial_graph_manager: &PartialGraphManager,
957 ) -> bool {
958 if !self.is_ready_to_merge() {
959 return false;
960 }
961
962 partial_graph_manager.pending_barrier_num(self.partial_graph_id)
965 <= self.max_lagged_barrier_num
966 }
967}
968
969impl CreatingStreamingJobControl {
970 pub(crate) fn start_completing(
971 &mut self,
972 partial_graph_manager: &mut PartialGraphManager,
973 min_upstream_inflight_epoch: Option<u64>,
974 upstream_committed_epoch: u64,
975 ) -> Option<(
976 u64,
977 HashMap<WorkerId, BarrierCompleteResponse>,
978 PartialGraphBarrierInfo,
979 bool,
980 )> {
981 if upstream_committed_epoch < self.snapshot_epoch {
983 return None;
984 }
985 let (finished_at_epoch, epoch_end_bound) = match &self.status {
986 CreatingStreamingJobStatus::Finishing(finish_at_epoch, _) => {
987 let epoch_end_bound = min_upstream_inflight_epoch
988 .map(|upstream_epoch| {
989 if upstream_epoch < *finish_at_epoch {
990 Excluded(upstream_epoch)
991 } else {
992 Unbounded
993 }
994 })
995 .unwrap_or(Unbounded);
996 (Some(*finish_at_epoch), epoch_end_bound)
997 }
998 CreatingStreamingJobStatus::ConsumingSnapshot { .. }
999 | CreatingStreamingJobStatus::ConsumingLogStore { .. } => (
1000 None,
1001 min_upstream_inflight_epoch
1002 .map(Excluded)
1003 .unwrap_or(Unbounded),
1004 ),
1005 CreatingStreamingJobStatus::Resetting(..) => {
1006 return None;
1007 }
1008 CreatingStreamingJobStatus::PlaceHolder => {
1009 unreachable!()
1010 }
1011 };
1012 partial_graph_manager
1013 .start_completing(
1014 self.partial_graph_id,
1015 epoch_end_bound,
1016 |non_checkpoint_epoch, _, _| {
1017 if let Some(finish_at_epoch) = finished_at_epoch {
1018 assert!(non_checkpoint_epoch.prev < finish_at_epoch);
1019 }
1020 },
1021 )
1022 .map(|(epoch, resps, info)| {
1023 let is_finish_epoch = if let Some(finish_at_epoch) = finished_at_epoch {
1024 assert!(!info.post_collect_command.should_checkpoint());
1025 if epoch == finish_at_epoch {
1026 self.ack_completed(partial_graph_manager, epoch);
1028 true
1029 } else {
1030 false
1031 }
1032 } else {
1033 false
1034 };
1035 (epoch, resps, info, is_finish_epoch)
1036 })
1037 }
1038
1039 pub(super) fn ack_completed(
1040 &mut self,
1041 partial_graph_manager: &mut PartialGraphManager,
1042 completed_epoch: u64,
1043 ) {
1044 match &self.status {
1045 CreatingStreamingJobStatus::ConsumingSnapshot { .. }
1046 | CreatingStreamingJobStatus::ConsumingLogStore { .. }
1047 | CreatingStreamingJobStatus::Finishing(_, _) => {
1048 partial_graph_manager.ack_completed(self.partial_graph_id, completed_epoch);
1049 if let Some(prev_max_committed_epoch) =
1050 self.max_committed_epoch.replace(completed_epoch)
1051 {
1052 assert!(completed_epoch > prev_max_committed_epoch);
1053 }
1054 }
1055 CreatingStreamingJobStatus::Resetting(_) => {
1056 }
1059 CreatingStreamingJobStatus::PlaceHolder => {
1060 unreachable!()
1061 }
1062 }
1063 }
1064
1065 pub(crate) fn fragment_infos(&self) -> Option<&HashMap<FragmentId, InflightFragmentInfo>> {
1066 self.status.fragment_infos()
1067 }
1068
1069 pub fn into_tracking_job(self) -> TrackingJob {
1070 match self.status {
1071 CreatingStreamingJobStatus::ConsumingSnapshot { .. }
1072 | CreatingStreamingJobStatus::ConsumingLogStore { .. }
1073 | CreatingStreamingJobStatus::Resetting(..)
1074 | CreatingStreamingJobStatus::PlaceHolder => {
1075 unreachable!("expect finish")
1076 }
1077 CreatingStreamingJobStatus::Finishing(_, tracking_job) => tracking_job,
1078 }
1079 }
1080
1081 pub(super) fn on_partial_graph_reset(mut self) {
1082 match &mut self.status {
1083 CreatingStreamingJobStatus::Resetting(notifiers) => {
1084 for notifier in notifiers.drain(..) {
1085 notifier.notify_collected();
1086 }
1087 }
1088 CreatingStreamingJobStatus::ConsumingSnapshot { .. }
1089 | CreatingStreamingJobStatus::ConsumingLogStore { .. }
1090 | CreatingStreamingJobStatus::Finishing(_, _) => {
1091 panic!(
1092 "should be resetting when receiving reset partial graph resp, but at {:?}",
1093 self.status
1094 )
1095 }
1096 CreatingStreamingJobStatus::PlaceHolder => {
1097 unreachable!()
1098 }
1099 }
1100 }
1101
1102 pub(super) fn drop(
1106 &mut self,
1107 notifier: Option<&mut NotifierStarter>,
1108 partial_graph_manager: &mut PartialGraphManager,
1109 ) -> bool {
1110 match &mut self.status {
1111 CreatingStreamingJobStatus::Resetting(existing_notifiers) => {
1112 existing_notifiers.extend(notifier.map(NotifierStarter::add_notify));
1113 true
1114 }
1115 CreatingStreamingJobStatus::ConsumingSnapshot { .. }
1116 | CreatingStreamingJobStatus::ConsumingLogStore { .. } => {
1117 partial_graph_manager.reset_partial_graphs([self.partial_graph_id]);
1118 self.status = CreatingStreamingJobStatus::Resetting(
1119 notifier
1120 .map(NotifierStarter::add_notify)
1121 .into_iter()
1122 .collect(),
1123 );
1124 true
1125 }
1126 CreatingStreamingJobStatus::Finishing(_, _) => false,
1127 CreatingStreamingJobStatus::PlaceHolder => {
1128 unreachable!()
1129 }
1130 }
1131 }
1132
1133 pub(crate) fn reset(self) -> bool {
1134 match self.status {
1135 CreatingStreamingJobStatus::ConsumingSnapshot { .. }
1136 | CreatingStreamingJobStatus::ConsumingLogStore { .. }
1137 | CreatingStreamingJobStatus::Finishing(_, _) => false,
1138 CreatingStreamingJobStatus::Resetting(notifiers) => {
1139 for notifier in notifiers {
1140 notifier.notify_collected();
1141 }
1142 true
1143 }
1144 CreatingStreamingJobStatus::PlaceHolder => {
1145 unreachable!()
1146 }
1147 }
1148 }
1149}
1150
1151#[cfg(test)]
1152mod tests {
1153 use super::*;
1154
1155 #[test]
1156 fn test_snapshot_backfill_max_pending_barrier_num() {
1157 let mut opts = MetaOpts::test(false);
1158 opts.in_flight_barrier_nums = 10;
1159
1160 opts.snapshot_backfill_barrier_amplification_factor = 0;
1161 assert_eq!(snapshot_backfill_max_pending_barrier_num(&opts), 10);
1162
1163 opts.snapshot_backfill_barrier_amplification_factor = 1;
1164 assert_eq!(snapshot_backfill_max_pending_barrier_num(&opts), 10);
1165
1166 opts.snapshot_backfill_barrier_amplification_factor = 10;
1167 assert_eq!(snapshot_backfill_max_pending_barrier_num(&opts), 100);
1168
1169 opts.in_flight_barrier_nums = usize::MAX;
1170 assert_eq!(snapshot_backfill_max_pending_barrier_num(&opts), usize::MAX);
1171 }
1172
1173 #[test]
1174 fn test_resolve_since_timestamp_upstream_log_epochs() {
1175 let upstream_log_epochs = vec![(vec![45, 50], 55)];
1176
1177 let (initial_barrier, barriers) =
1178 CreatingStreamingJobControl::resolve_since_timestamp_upstream_log_epochs(
1179 &upstream_log_epochs,
1180 [].iter(),
1181 40,
1182 60,
1183 )
1184 .unwrap();
1185
1186 assert_eq!(
1187 (initial_barrier.prev_epoch(), initial_barrier.curr_epoch()),
1188 (40, 45)
1189 );
1190 assert!(initial_barrier.kind.is_checkpoint());
1191 assert_eq!(
1192 barriers
1193 .iter()
1194 .map(|barrier| (barrier.prev_epoch(), barrier.curr_epoch()))
1195 .collect::<Vec<_>>(),
1196 vec![(45, 50), (50, 55), (55, 60)]
1197 );
1198 assert_eq!(
1199 barriers
1200 .iter()
1201 .map(|barrier| match &barrier.kind {
1202 BarrierKind::Checkpoint(epochs) => Some(epochs.clone()),
1203 _ => None,
1204 })
1205 .collect::<Vec<_>>(),
1206 vec![None, None, Some(vec![45, 50, 55])]
1207 );
1208 }
1209
1210 #[test]
1211 fn test_resolve_since_timestamp_upstream_log_epochs_with_pending_barriers() {
1212 let upstream_log_epochs = vec![(vec![45, 50], 55)];
1213 let pending_upstream_barriers = [
1214 BarrierInfo {
1215 prev_epoch: TracedEpoch::new(Epoch(60)),
1216 curr_epoch: TracedEpoch::new(Epoch(65)),
1217 kind: BarrierKind::Barrier,
1218 },
1219 BarrierInfo {
1220 prev_epoch: TracedEpoch::new(Epoch(65)),
1221 curr_epoch: TracedEpoch::new(Epoch(70)),
1222 kind: BarrierKind::Checkpoint(vec![60, 65]),
1223 },
1224 ];
1225
1226 let (initial_barrier, barriers) =
1227 CreatingStreamingJobControl::resolve_since_timestamp_upstream_log_epochs(
1228 &upstream_log_epochs,
1229 pending_upstream_barriers.iter(),
1230 40,
1231 70,
1232 )
1233 .unwrap();
1234
1235 assert_eq!(
1236 (initial_barrier.prev_epoch(), initial_barrier.curr_epoch()),
1237 (40, 45)
1238 );
1239 assert!(initial_barrier.kind.is_checkpoint());
1240 assert_eq!(
1241 barriers
1242 .iter()
1243 .map(|barrier| (barrier.prev_epoch(), barrier.curr_epoch()))
1244 .collect::<Vec<_>>(),
1245 vec![(45, 50), (50, 55), (55, 60), (60, 65), (65, 70)]
1246 );
1247 assert_eq!(
1248 barriers
1249 .iter()
1250 .map(|barrier| match &barrier.kind {
1251 BarrierKind::Checkpoint(epochs) => Some(epochs.clone()),
1252 _ => None,
1253 })
1254 .collect::<Vec<_>>(),
1255 vec![None, None, Some(vec![45, 50, 55]), None, Some(vec![60, 65])]
1256 );
1257 }
1258
1259 #[test]
1260 fn test_resolve_since_timestamp_upstream_log_epochs_with_gap_before_pending_barriers() {
1261 let upstream_log_epochs = vec![(vec![61, 62, 63, 64], 65)];
1262 let pending_upstream_barriers = [
1263 BarrierInfo {
1264 prev_epoch: TracedEpoch::new(Epoch(66)),
1265 curr_epoch: TracedEpoch::new(Epoch(67)),
1266 kind: BarrierKind::Barrier,
1267 },
1268 BarrierInfo {
1269 prev_epoch: TracedEpoch::new(Epoch(67)),
1270 curr_epoch: TracedEpoch::new(Epoch(68)),
1271 kind: BarrierKind::Barrier,
1272 },
1273 BarrierInfo {
1274 prev_epoch: TracedEpoch::new(Epoch(68)),
1275 curr_epoch: TracedEpoch::new(Epoch(69)),
1276 kind: BarrierKind::Barrier,
1277 },
1278 BarrierInfo {
1279 prev_epoch: TracedEpoch::new(Epoch(69)),
1280 curr_epoch: TracedEpoch::new(Epoch(70)),
1281 kind: BarrierKind::Barrier,
1282 },
1283 ];
1284
1285 let (initial_barrier, barriers) =
1286 CreatingStreamingJobControl::resolve_since_timestamp_upstream_log_epochs(
1287 &upstream_log_epochs,
1288 pending_upstream_barriers.iter(),
1289 60,
1290 70,
1291 )
1292 .unwrap();
1293
1294 assert_eq!(
1295 (initial_barrier.prev_epoch(), initial_barrier.curr_epoch()),
1296 (60, 61)
1297 );
1298 assert!(initial_barrier.kind.is_checkpoint());
1299 assert_eq!(
1300 barriers
1301 .iter()
1302 .map(|barrier| (barrier.prev_epoch(), barrier.curr_epoch()))
1303 .collect::<Vec<_>>(),
1304 vec![
1305 (61, 62),
1306 (62, 63),
1307 (63, 64),
1308 (64, 65),
1309 (65, 66),
1310 (66, 67),
1311 (67, 68),
1312 (68, 69),
1313 (69, 70)
1314 ]
1315 );
1316 assert_eq!(
1317 barriers
1318 .iter()
1319 .map(|barrier| match &barrier.kind {
1320 BarrierKind::Checkpoint(epochs) => Some(epochs.clone()),
1321 _ => None,
1322 })
1323 .collect::<Vec<_>>(),
1324 vec![
1325 None,
1326 None,
1327 None,
1328 None,
1329 Some(vec![61, 62, 63, 64, 65]),
1330 None,
1331 None,
1332 None,
1333 None
1334 ]
1335 );
1336 }
1337
1338 #[test]
1339 fn test_resolve_since_timestamp_upstream_log_epochs_without_pending_barriers() {
1340 let upstream_log_epochs = vec![(vec![61, 62, 63, 64], 65)];
1341
1342 let (initial_barrier, barriers) =
1343 CreatingStreamingJobControl::resolve_since_timestamp_upstream_log_epochs(
1344 &upstream_log_epochs,
1345 [].iter(),
1346 60,
1347 66,
1348 )
1349 .unwrap();
1350
1351 assert_eq!(
1352 (initial_barrier.prev_epoch(), initial_barrier.curr_epoch()),
1353 (60, 61)
1354 );
1355 assert!(initial_barrier.kind.is_checkpoint());
1356 assert_eq!(
1357 barriers
1358 .iter()
1359 .map(|barrier| (barrier.prev_epoch(), barrier.curr_epoch()))
1360 .collect::<Vec<_>>(),
1361 vec![(61, 62), (62, 63), (63, 64), (64, 65), (65, 66)]
1362 );
1363 assert_eq!(
1364 barriers
1365 .iter()
1366 .map(|barrier| match &barrier.kind {
1367 BarrierKind::Checkpoint(epochs) => Some(epochs.clone()),
1368 _ => None,
1369 })
1370 .collect::<Vec<_>>(),
1371 vec![None, None, None, None, Some(vec![61, 62, 63, 64, 65])]
1372 );
1373 }
1374}