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