1use std::collections::hash_map::Entry;
16use std::collections::{HashMap, HashSet};
17use std::future::{Future, poll_fn};
18use std::ops::Bound::{Excluded, Unbounded};
19use std::sync::atomic::AtomicU32;
20use std::task::Poll;
21
22use anyhow::anyhow;
23use fail::fail_point;
24use itertools::Itertools;
25use risingwave_common::catalog::{DatabaseId, TableId};
26use risingwave_common::id::JobId;
27use risingwave_common::metrics::{LabelGuardedHistogram, LabelGuardedIntGauge};
28use risingwave_common::util::epoch::{Epoch, EpochPair};
29use risingwave_common::util::stream_graph_visitor::visit_stream_node_cont;
30use risingwave_meta_model::WorkerId;
31use risingwave_pb::common::WorkerNode;
32use risingwave_pb::hummock::HummockVersionStats;
33use risingwave_pb::id::{FragmentId, PartialGraphId};
34use risingwave_pb::stream_plan::DispatcherType as PbDispatcherType;
35use risingwave_pb::stream_plan::stream_node::NodeBody;
36use risingwave_pb::stream_service::BarrierCompleteResponse;
37use risingwave_pb::stream_service::streaming_control_stream_response::ResetPartialGraphResponse;
38use tracing::{debug, warn};
39
40use crate::barrier::cdc_progress::CdcProgress;
41use crate::barrier::checkpoint::independent_job::{
42 BatchRefreshJobTriggerContext, IndependentCheckpointJobControl,
43};
44use crate::barrier::checkpoint::recovery::{
45 DatabaseRecoveringState, DatabaseStatusAction, EnterInitializing, EnterRunning,
46 RecoveringStateAction,
47};
48use crate::barrier::checkpoint::state::{ApplyCommandInfo, BarrierWorkerState};
49use crate::barrier::complete_task::{BarrierCompleteOutput, CompleteBarrierTask};
50use crate::barrier::info::{InflightDatabaseInfo, SharedActorInfos};
51use crate::barrier::notifier::Notifier;
52use crate::barrier::partial_graph::{CollectedBarrier, PartialGraphManager, PartialGraphStat};
53use crate::barrier::progress::TrackingJob;
54use crate::barrier::rpc::{from_partial_graph_id, to_partial_graph_id};
55use crate::barrier::schedule::{NewBarrier, PeriodicBarriers};
56use crate::barrier::utils::{BarrierItemCollector, collect_independent_job_commit_epoch_info};
57use crate::barrier::{
58 BackfillProgress, Command, CreateStreamingJobType, FragmentBackfillProgress, Reschedule,
59};
60use crate::controller::fragment::InflightFragmentInfo;
61use crate::controller::scale::{build_no_shuffle_fragment_graph_edges, find_no_shuffle_graphs};
62use crate::manager::MetaSrvEnv;
63
64fn fragment_has_online_unreschedulable_scan(fragment: &InflightFragmentInfo) -> bool {
65 let mut has_unreschedulable_scan = false;
66 visit_stream_node_cont(&fragment.nodes, |node| {
67 if let Some(NodeBody::StreamScan(stream_scan)) = node.node_body.as_ref() {
68 let scan_type = stream_scan.stream_scan_type();
69 if !scan_type.is_reschedulable(true) {
70 has_unreschedulable_scan = true;
71 return false;
72 }
73 }
74 true
75 });
76 has_unreschedulable_scan
77}
78
79fn collect_fragment_upstream_fragment_ids(
80 fragment: &InflightFragmentInfo,
81 upstream_fragment_ids: &mut HashSet<FragmentId>,
82) {
83 visit_stream_node_cont(&fragment.nodes, |node| {
84 if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref() {
85 upstream_fragment_ids.insert(merge.upstream_fragment_id);
86 }
87 true
88 });
89}
90
91use crate::model::ActorId;
92use crate::rpc::metrics::GLOBAL_META_METRICS;
93use crate::{MetaError, MetaResult};
94
95pub(crate) struct CheckpointControl {
96 pub(crate) env: MetaSrvEnv,
97 pub(super) databases: HashMap<DatabaseId, DatabaseCheckpointControlStatus>,
98 pub(super) hummock_version_stats: HummockVersionStats,
99 pub(crate) in_flight_barrier_nums: usize,
101}
102
103impl CheckpointControl {
104 pub fn new(env: MetaSrvEnv) -> Self {
105 Self {
106 in_flight_barrier_nums: env.opts.in_flight_barrier_nums,
107 env,
108 databases: Default::default(),
109 hummock_version_stats: Default::default(),
110 }
111 }
112
113 pub(crate) fn recover(
114 databases: HashMap<DatabaseId, DatabaseCheckpointControl>,
115 failed_databases: HashMap<DatabaseId, HashSet<PartialGraphId>>, hummock_version_stats: HummockVersionStats,
117 env: MetaSrvEnv,
118 ) -> Self {
119 env.shared_actor_infos()
120 .retain_databases(databases.keys().chain(failed_databases.keys()).cloned());
121 Self {
122 in_flight_barrier_nums: env.opts.in_flight_barrier_nums,
123 env,
124 databases: databases
125 .into_iter()
126 .map(|(database_id, control)| {
127 (
128 database_id,
129 DatabaseCheckpointControlStatus::Running(control),
130 )
131 })
132 .chain(failed_databases.into_iter().map(
133 |(database_id, resetting_partial_graphs)| {
134 (
135 database_id,
136 DatabaseCheckpointControlStatus::Recovering(
137 DatabaseRecoveringState::new_resetting(
138 database_id,
139 resetting_partial_graphs,
140 ),
141 ),
142 )
143 },
144 ))
145 .collect(),
146 hummock_version_stats,
147 }
148 }
149
150 pub(crate) fn ack_completed(
151 &mut self,
152 partial_graph_manager: &mut PartialGraphManager,
153 output: BarrierCompleteOutput,
154 ) {
155 self.hummock_version_stats = output.hummock_version_stats;
156 for (database_id, (command_prev_epoch, independent_job_epochs)) in output.epochs_to_ack {
157 self.databases
158 .get_mut(&database_id)
159 .expect("should exist")
160 .expect_running("should have wait for completing command before enter recovery")
161 .ack_completed(
162 partial_graph_manager,
163 command_prev_epoch,
164 independent_job_epochs,
165 );
166 }
167 }
168
169 pub(crate) fn next_complete_barrier_task(
170 &mut self,
171 periodic_barriers: &mut PeriodicBarriers,
172 partial_graph_manager: &mut PartialGraphManager,
173 ) -> Option<CompleteBarrierTask> {
174 let mut task = None;
175 for database in self.databases.values_mut() {
176 let Some(database) = database.running_state_mut() else {
177 continue;
178 };
179 database.next_complete_barrier_task(
180 periodic_barriers,
181 partial_graph_manager,
182 &mut task,
183 &self.hummock_version_stats,
184 );
185 }
186 task
187 }
188
189 pub(crate) fn barrier_collected(
190 &mut self,
191 partial_graph_id: PartialGraphId,
192 collected_barrier: CollectedBarrier<'_>,
193 periodic_barriers: &mut PeriodicBarriers,
194 ) -> MetaResult<()> {
195 let (database_id, _) = from_partial_graph_id(partial_graph_id);
196 let database_status = self.databases.get_mut(&database_id).expect("should exist");
197 match database_status {
198 DatabaseCheckpointControlStatus::Running(database) => {
199 database.barrier_collected(partial_graph_id, collected_barrier, periodic_barriers)
200 }
201 DatabaseCheckpointControlStatus::Recovering(_) => {
202 if cfg!(debug_assertions) {
203 panic!(
204 "receive collected barrier {:?} on recovering database {} from partial graph {}",
205 collected_barrier, database_id, partial_graph_id
206 );
207 } else {
208 warn!(?collected_barrier, %partial_graph_id, "ignore collected barrier on recovering database");
209 }
210 Ok(())
211 }
212 }
213 }
214
215 pub(crate) fn recovering_databases(&self) -> impl Iterator<Item = DatabaseId> + '_ {
216 self.databases.iter().filter_map(|(database_id, database)| {
217 database.running_state().is_none().then_some(*database_id)
218 })
219 }
220
221 pub(crate) fn running_databases(&self) -> impl Iterator<Item = DatabaseId> + '_ {
222 self.databases.iter().filter_map(|(database_id, database)| {
223 database.running_state().is_some().then_some(*database_id)
224 })
225 }
226
227 pub(crate) fn database_info(&self, database_id: DatabaseId) -> Option<&InflightDatabaseInfo> {
228 self.databases
229 .get(&database_id)
230 .and_then(|database| database.running_state())
231 .map(|database| &database.database_info)
232 }
233
234 pub(crate) fn handle_new_barrier(
236 &mut self,
237 new_barrier: NewBarrier,
238 partial_graph_manager: &mut PartialGraphManager,
239 worker_nodes: &HashMap<WorkerId, WorkerNode>,
240 ) -> MetaResult<()> {
241 let NewBarrier {
242 database_id,
243 command,
244 span,
245 checkpoint,
246 } = new_barrier;
247
248 if let Some((mut command, notifier)) = command {
249 if let &mut Command::CreateStreamingJob {
250 ref mut cross_db_snapshot_backfill_info,
251 ref info,
252 ..
253 } = &mut command
254 {
255 for (table_id, snapshot_epoch) in
256 &mut cross_db_snapshot_backfill_info.upstream_mv_table_id_to_backfill_epoch
257 {
258 for database in self.databases.values() {
259 if let Some(database) = database.running_state()
260 && database.database_info.contains_job(table_id.as_job_id())
261 {
262 if let Some(committed_epoch) = database.committed_epoch {
263 *snapshot_epoch = Some(committed_epoch);
264 }
265 break;
266 }
267 }
268 if snapshot_epoch.is_none() {
269 let table_id = *table_id;
270 warn!(
271 ?cross_db_snapshot_backfill_info,
272 ?table_id,
273 ?info,
274 "database of cross db upstream table not found"
275 );
276 let err: MetaError =
277 anyhow!("database of cross db upstream table {} not found", table_id)
278 .into();
279 notifier.notify_start_failed(err);
280
281 return Ok(());
282 }
283 }
284 }
285
286 let database = match self.databases.entry(database_id) {
287 Entry::Occupied(entry) => entry
288 .into_mut()
289 .expect_running("should not have command when not running"),
290 Entry::Vacant(entry) => match &command {
291 Command::CreateStreamingJob { info, job_type, .. } => {
292 let CreateStreamingJobType::Normal = job_type else {
293 if cfg!(debug_assertions) {
294 panic!(
295 "unexpected first job of type {job_type:?} with info {info:?}"
296 );
297 } else {
298 notifier.notify_start_failed(anyhow!("unexpected job_type {job_type:?} for first job {} in database {database_id}", info.streaming_job.id()).into());
299 return Ok(());
300 }
301 };
302 let new_database = DatabaseCheckpointControl::new(
303 database_id,
304 self.env.shared_actor_infos().clone(),
305 );
306 let adder = partial_graph_manager.add_partial_graph(
307 to_partial_graph_id(database_id, None),
308 DatabaseCheckpointControlMetrics::new(database_id),
309 );
310 adder.added();
311 entry
312 .insert(DatabaseCheckpointControlStatus::Running(new_database))
313 .expect_running("just initialized as running")
314 }
315 Command::Flush
316 | Command::Pause
317 | Command::Resume
318 | Command::DropStreamingJobs { .. }
319 | Command::DropSubscription { .. } => {
320 notifier.start().started();
321 warn!(?command, "skip command for empty database");
322 return Ok(());
323 }
324 Command::RescheduleIntent { .. }
325 | Command::ReplaceStreamJob(_)
326 | Command::SourceChangeSplit(_)
327 | Command::Throttle { .. }
328 | Command::CreateSubscription { .. }
329 | Command::AlterSubscriptionRetention { .. }
330 | Command::ConnectorPropsChange(_)
331 | Command::Refresh { .. }
332 | Command::ListFinish { .. }
333 | Command::LoadFinish { .. }
334 | Command::ResetSource { .. }
335 | Command::ResumeBackfill { .. }
336 | Command::InjectSourceOffsets { .. } => {
337 if cfg!(debug_assertions) {
338 panic!(
339 "new database graph info can only be created for normal creating streaming job, but get command: {} {:?}",
340 database_id, command
341 )
342 } else {
343 warn!(%database_id, ?command, "database not exist when handling command");
344 notifier.notify_start_failed(anyhow!("database {database_id} not exist when handling command {command:?}").into());
345 return Ok(());
346 }
347 }
348 },
349 };
350
351 database.handle_new_barrier(
352 Some((command, notifier)),
353 checkpoint,
354 span,
355 partial_graph_manager,
356 &self.hummock_version_stats,
357 worker_nodes,
358 )
359 } else {
360 let database = match self.databases.entry(database_id) {
361 Entry::Occupied(entry) => entry.into_mut(),
362 Entry::Vacant(_) => {
363 return Ok(());
366 }
367 };
368 let Some(database) = database.running_state_mut() else {
369 return Ok(());
371 };
372 if partial_graph_manager.pending_barrier_num(database.partial_graph_id)
373 >= self.in_flight_barrier_nums
374 {
375 return Ok(());
377 }
378 database.handle_new_barrier(
379 None,
380 checkpoint,
381 span,
382 partial_graph_manager,
383 &self.hummock_version_stats,
384 worker_nodes,
385 )
386 }
387 }
388
389 pub(crate) fn gen_backfill_progress(&self) -> HashMap<JobId, BackfillProgress> {
390 let mut progress = HashMap::new();
391 for status in self.databases.values() {
392 let Some(database_checkpoint_control) = status.running_state() else {
393 continue;
394 };
395 progress.extend(
397 database_checkpoint_control
398 .database_info
399 .gen_backfill_progress(),
400 );
401 for (job_id, job) in &database_checkpoint_control.independent_checkpoint_job_controls {
403 if let Some(p) = job.gen_backfill_progress() {
404 progress.insert(*job_id, p);
405 }
406 }
407 }
408 progress
409 }
410
411 pub(crate) fn gen_fragment_backfill_progress(&self) -> Vec<FragmentBackfillProgress> {
412 let mut progress = Vec::new();
413 for status in self.databases.values() {
414 let Some(database_checkpoint_control) = status.running_state() else {
415 continue;
416 };
417 progress.extend(
418 database_checkpoint_control
419 .database_info
420 .gen_fragment_backfill_progress(),
421 );
422 for job in database_checkpoint_control
423 .independent_checkpoint_job_controls
424 .values()
425 {
426 progress.extend(job.gen_fragment_backfill_progress());
427 }
428 }
429 progress
430 }
431
432 pub(crate) fn gen_cdc_progress(&self) -> HashMap<JobId, CdcProgress> {
433 let mut progress = HashMap::new();
434 for status in self.databases.values() {
435 let Some(database_checkpoint_control) = status.running_state() else {
436 continue;
437 };
438 progress.extend(database_checkpoint_control.database_info.gen_cdc_progress());
440 }
441 progress
442 }
443
444 pub(crate) fn databases_failed_at_worker_err(
445 &mut self,
446 worker_id: WorkerId,
447 ) -> impl Iterator<Item = DatabaseId> + '_ {
448 self.databases
449 .iter_mut()
450 .filter_map(
451 move |(database_id, database_status)| match database_status {
452 DatabaseCheckpointControlStatus::Running(control) => {
453 if !control.is_valid_after_worker_err(worker_id) {
454 Some(*database_id)
455 } else {
456 None
457 }
458 }
459 DatabaseCheckpointControlStatus::Recovering(state) => {
460 if !state.is_valid_after_worker_err(worker_id) {
461 Some(*database_id)
462 } else {
463 None
464 }
465 }
466 },
467 )
468 }
469
470 pub(crate) fn get_batch_refresh_trigger_info(
473 &self,
474 database_id: DatabaseId,
475 job_id: JobId,
476 ) -> u64 {
477 let database = self
478 .databases
479 .get(&database_id)
480 .and_then(|s| s.running_state())
481 .expect("database should be running for batch refresh trigger");
482 database.get_batch_refresh_trigger_info(job_id)
483 }
484
485 pub(crate) fn start_batch_refresh_run(
486 &mut self,
487 database_id: DatabaseId,
488 job_id: JobId,
489 context: &BatchRefreshJobTriggerContext,
490 worker_nodes: &HashMap<WorkerId, WorkerNode>,
491 actor_id_counter: &AtomicU32,
492 partial_graph_manager: &mut PartialGraphManager,
493 ) -> MetaResult<bool> {
494 let database = self
495 .databases
496 .get_mut(&database_id)
497 .and_then(|s| s.running_state_mut())
498 .expect("database should be running");
499 database.start_batch_refresh_run(
500 job_id,
501 context,
502 worker_nodes,
503 actor_id_counter,
504 partial_graph_manager,
505 )
506 }
507
508 pub(crate) fn apply_batch_refresh_fragment_infos(
509 &mut self,
510 database_id: DatabaseId,
511 job_id: JobId,
512 ) {
513 let database = self
514 .databases
515 .get_mut(&database_id)
516 .and_then(|s| s.running_state_mut())
517 .expect("database should be running");
518 let br_job = match database
519 .independent_checkpoint_job_controls
520 .get(&job_id)
521 .expect("job should exist")
522 {
523 IndependentCheckpointJobControl::BatchRefresh(job) => job,
524 _ => panic!("expected batch refresh job"),
525 };
526 if let Some(fragment_infos) = br_job.fragment_infos() {
527 database
528 .database_info
529 .shared_actor_infos
530 .upsert(database_id, fragment_infos.values().map(|f| (f, job_id)));
531 }
532 }
533}
534
535pub(crate) enum CheckpointControlEvent<'a> {
536 EnteringInitializing(DatabaseStatusAction<'a, EnterInitializing>),
537 EnteringRunning(DatabaseStatusAction<'a, EnterRunning>),
538 BatchRefreshTrigger {
541 database_id: DatabaseId,
542 job_id: JobId,
543 },
544}
545
546impl CheckpointControl {
547 pub(crate) fn on_partial_graph_reset(
548 &mut self,
549 partial_graph_id: PartialGraphId,
550 reset_resps: HashMap<WorkerId, ResetPartialGraphResponse>,
551 ) {
552 let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
553 match self.databases.get_mut(&database_id).expect("should exist") {
554 DatabaseCheckpointControlStatus::Running(database) => {
555 if let Some(independent_job_id) = independent_job_id {
556 match database
557 .independent_checkpoint_job_controls
558 .remove(&independent_job_id)
559 {
560 Some(independent_job) => {
561 independent_job.on_partial_graph_reset();
562 }
563 None => {
564 if cfg!(debug_assertions) {
565 panic!(
566 "receive reset partial graph resp on non-existing independent job {independent_job_id} in database {database_id}"
567 )
568 }
569 warn!(
570 %database_id,
571 %independent_job_id,
572 "ignore reset partial graph resp on non-existing independent job on running database"
573 );
574 }
575 }
576 } else {
577 unreachable!("should not receive reset database resp when database running")
578 }
579 }
580 DatabaseCheckpointControlStatus::Recovering(state) => {
581 state.on_partial_graph_reset(partial_graph_id, reset_resps);
582 }
583 }
584 }
585
586 pub(crate) fn on_partial_graph_initialized(
587 &mut self,
588 partial_graph_id: PartialGraphId,
589 partial_graph_manager: &mut PartialGraphManager,
590 ) -> MetaResult<()> {
591 let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
592 match self.databases.get_mut(&database_id).expect("should exist") {
593 DatabaseCheckpointControlStatus::Running(database) => {
594 let Some(independent_job_id) = independent_job_id else {
595 unreachable!("database partial graph should not initialize when running")
596 };
597 let job = database
598 .independent_checkpoint_job_controls
599 .get_mut(&independent_job_id)
600 .expect("independent job should exist");
601 match job {
602 IndependentCheckpointJobControl::BatchRefresh(job) => {
603 job.on_log_store_initialized(partial_graph_manager)
604 }
605 IndependentCheckpointJobControl::CreatingStreamingJob(_) => {
606 unreachable!("creating streaming job should not initialize when running")
607 }
608 }
609 }
610 DatabaseCheckpointControlStatus::Recovering(state) => {
611 state.partial_graph_initialized(partial_graph_id);
612 Ok(())
613 }
614 }
615 }
616
617 pub(crate) fn next_event(
618 &mut self,
619 ) -> impl Future<Output = CheckpointControlEvent<'_>> + Send + '_ {
620 let mut this = Some(self);
621 poll_fn(move |cx| {
622 let Some(this_mut) = this.as_mut() else {
623 unreachable!("should not be polled after poll ready")
624 };
625 for (&database_id, database_status) in &mut this_mut.databases {
626 match database_status {
627 DatabaseCheckpointControlStatus::Running(database) => {
628 if let Some(committed_epoch) = database.committed_epoch {
630 for (job_id, job) in &database.independent_checkpoint_job_controls {
631 if let IndependentCheckpointJobControl::BatchRefresh(br_job) = job
632 && br_job.should_start_refresh(committed_epoch)
633 {
634 let job_id = *job_id;
635 let _ = this.take().expect("checked Some");
636 return Poll::Ready(
637 CheckpointControlEvent::BatchRefreshTrigger {
638 database_id,
639 job_id,
640 },
641 );
642 }
643 }
644 }
645 }
646 DatabaseCheckpointControlStatus::Recovering(state) => {
647 let poll_result = state.poll_next_event(cx);
648 if let Poll::Ready(action) = poll_result {
649 let this = this.take().expect("checked Some");
650 return Poll::Ready(match action {
651 RecoveringStateAction::EnterInitializing(reset_workers) => {
652 CheckpointControlEvent::EnteringInitializing(
653 this.new_database_status_action(
654 database_id,
655 EnterInitializing(reset_workers),
656 ),
657 )
658 }
659 RecoveringStateAction::EnterRunning => {
660 CheckpointControlEvent::EnteringRunning(
661 this.new_database_status_action(database_id, EnterRunning),
662 )
663 }
664 });
665 }
666 }
667 }
668 }
669 Poll::Pending
670 })
671 }
672}
673
674pub(crate) enum DatabaseCheckpointControlStatus {
675 Running(DatabaseCheckpointControl),
676 Recovering(DatabaseRecoveringState),
677}
678
679impl DatabaseCheckpointControlStatus {
680 fn running_state(&self) -> Option<&DatabaseCheckpointControl> {
681 match self {
682 DatabaseCheckpointControlStatus::Running(state) => Some(state),
683 DatabaseCheckpointControlStatus::Recovering(_) => None,
684 }
685 }
686
687 fn running_state_mut(&mut self) -> Option<&mut DatabaseCheckpointControl> {
688 match self {
689 DatabaseCheckpointControlStatus::Running(state) => Some(state),
690 DatabaseCheckpointControlStatus::Recovering(_) => None,
691 }
692 }
693
694 fn expect_running(&mut self, reason: &'static str) -> &mut DatabaseCheckpointControl {
695 match self {
696 DatabaseCheckpointControlStatus::Running(state) => state,
697 DatabaseCheckpointControlStatus::Recovering(_) => {
698 panic!("should be at running: {}", reason)
699 }
700 }
701 }
702}
703
704pub(in crate::barrier) struct DatabaseCheckpointControlMetrics {
705 barrier_latency: LabelGuardedHistogram,
706 in_flight_barrier_nums: LabelGuardedIntGauge,
707 all_barrier_nums: LabelGuardedIntGauge,
708}
709
710impl DatabaseCheckpointControlMetrics {
711 pub(in crate::barrier) fn new(database_id: DatabaseId) -> Self {
712 let database_id_str = database_id.to_string();
713 let barrier_latency = GLOBAL_META_METRICS
714 .barrier_latency
715 .with_guarded_label_values(&[&database_id_str]);
716 let in_flight_barrier_nums = GLOBAL_META_METRICS
717 .in_flight_barrier_nums
718 .with_guarded_label_values(&[&database_id_str]);
719 let all_barrier_nums = GLOBAL_META_METRICS
720 .all_barrier_nums
721 .with_guarded_label_values(&[&database_id_str]);
722 Self {
723 barrier_latency,
724 in_flight_barrier_nums,
725 all_barrier_nums,
726 }
727 }
728}
729
730impl PartialGraphStat for DatabaseCheckpointControlMetrics {
731 fn observe_barrier_latency(&self, _epoch: EpochPair, barrier_latency_secs: f64) {
732 self.barrier_latency.observe(barrier_latency_secs);
733 }
734
735 fn observe_barrier_num(&self, inflight_barrier_num: usize, collected_barrier_num: usize) {
736 self.in_flight_barrier_nums.set(inflight_barrier_num as _);
737 self.all_barrier_nums
738 .set((inflight_barrier_num + collected_barrier_num) as _);
739 }
740}
741
742pub(in crate::barrier) struct DatabaseCheckpointControl {
744 pub(super) database_id: DatabaseId,
745 partial_graph_id: PartialGraphId,
746 pub(super) state: BarrierWorkerState,
747
748 finishing_jobs_collector:
749 BarrierItemCollector<JobId, (Vec<BarrierCompleteResponse>, TrackingJob), ()>,
750 completing_barrier: Option<EpochPair>,
752
753 committed_epoch: Option<u64>,
754
755 last_committed_barrier_time: Option<LabelGuardedIntGauge>,
758
759 pub(super) database_info: InflightDatabaseInfo,
760 pub independent_checkpoint_job_controls: HashMap<JobId, IndependentCheckpointJobControl>,
761}
762
763impl DatabaseCheckpointControl {
764 fn new(database_id: DatabaseId, shared_actor_infos: SharedActorInfos) -> Self {
765 Self {
766 database_id,
767 partial_graph_id: to_partial_graph_id(database_id, None),
768 state: BarrierWorkerState::new(),
769 finishing_jobs_collector: BarrierItemCollector::new(false),
770 completing_barrier: None,
771 committed_epoch: None,
772 last_committed_barrier_time: None,
773 database_info: InflightDatabaseInfo::empty(database_id, shared_actor_infos),
774 independent_checkpoint_job_controls: Default::default(),
775 }
776 }
777
778 pub(crate) fn recovery(
779 database_id: DatabaseId,
780 state: BarrierWorkerState,
781 committed_epoch: u64,
782 database_info: InflightDatabaseInfo,
783 independent_checkpoint_job_controls: HashMap<JobId, IndependentCheckpointJobControl>,
784 ) -> Self {
785 Self {
786 database_id,
787 partial_graph_id: to_partial_graph_id(database_id, None),
788 state,
789 finishing_jobs_collector: BarrierItemCollector::new(false),
790 completing_barrier: None,
791 committed_epoch: Some(committed_epoch),
792 last_committed_barrier_time: None,
793 database_info,
794 independent_checkpoint_job_controls,
795 }
796 }
797
798 pub(crate) fn is_valid_after_worker_err(&self, worker_id: WorkerId) -> bool {
799 !self.database_info.contains_worker(worker_id as _)
800 && self
801 .independent_checkpoint_job_controls
802 .values()
803 .all(|job| {
804 job.fragment_infos()
805 .map(|fragment_infos| {
806 !InflightFragmentInfo::contains_worker(
807 fragment_infos.values(),
808 worker_id,
809 )
810 })
811 .unwrap_or(true)
812 })
813 }
814
815 fn enqueue_command(&mut self, epoch: EpochPair, independent_jobs_to_wait: HashSet<JobId>) {
817 let prev_epoch = epoch.prev;
818 tracing::trace!(prev_epoch, ?independent_jobs_to_wait, "enqueue command");
819 if !independent_jobs_to_wait.is_empty() {
820 self.finishing_jobs_collector
821 .enqueue(epoch, independent_jobs_to_wait, ());
822 }
823 }
824
825 fn barrier_collected(
828 &mut self,
829 partial_graph_id: PartialGraphId,
830 collected_barrier: CollectedBarrier<'_>,
831 periodic_barriers: &mut PeriodicBarriers,
832 ) -> MetaResult<()> {
833 let prev_epoch = collected_barrier.epoch.prev;
834 tracing::trace!(
835 prev_epoch,
836 partial_graph_id = %partial_graph_id,
837 "barrier collected"
838 );
839 let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
840 assert_eq!(self.database_id, database_id);
841 if let Some(independent_job_id) = independent_job_id {
842 let job = self
843 .independent_checkpoint_job_controls
844 .get_mut(&independent_job_id)
845 .expect("should exist");
846 let should_force_checkpoint = job.collect(collected_barrier);
847 if should_force_checkpoint {
848 periodic_barriers.force_checkpoint_in_next_barrier(self.database_id);
849 }
850 }
851 Ok(())
852 }
853}
854
855impl DatabaseCheckpointControl {
856 fn collect_backfill_pinned_upstream_log_epoch(
858 &self,
859 ) -> HashMap<JobId, (u64, HashSet<TableId>)> {
860 self.independent_checkpoint_job_controls
861 .iter()
862 .map(|(job_id, job)| (*job_id, job.pinned_upstream_log_epoch()))
863 .collect()
864 }
865
866 fn collect_no_shuffle_fragment_relations_for_reschedule_check(
867 &self,
868 ) -> Vec<(FragmentId, FragmentId)> {
869 let mut no_shuffle_relations = Vec::new();
870 for fragment in self.database_info.fragment_infos() {
871 let downstream_fragment_id = fragment.fragment_id;
872 visit_stream_node_cont(&fragment.nodes, |node| {
873 if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
874 && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
875 {
876 no_shuffle_relations.push((merge.upstream_fragment_id, downstream_fragment_id));
877 }
878 true
879 });
880 }
881
882 for job in self.independent_checkpoint_job_controls.values() {
883 if let Some(fragment_infos) = job.fragment_infos() {
884 for fragment_info in fragment_infos.values() {
885 let downstream_fragment_id = fragment_info.fragment_id;
886 visit_stream_node_cont(&fragment_info.nodes, |node| {
887 if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
888 && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
889 {
890 no_shuffle_relations
891 .push((merge.upstream_fragment_id, downstream_fragment_id));
892 }
893 true
894 });
895 }
896 }
897 }
898 no_shuffle_relations
899 }
900
901 fn collect_reschedule_blocked_jobs_for_independent_jobs_inflight(
902 &self,
903 ) -> MetaResult<HashSet<JobId>> {
904 let mut initial_blocked_fragment_ids = HashSet::new();
905 for job in self.independent_checkpoint_job_controls.values() {
906 if let Some(fragment_infos) = job.fragment_infos() {
907 for fragment_info in fragment_infos.values() {
908 if fragment_has_online_unreschedulable_scan(fragment_info) {
909 initial_blocked_fragment_ids.insert(fragment_info.fragment_id);
910 collect_fragment_upstream_fragment_ids(
911 fragment_info,
912 &mut initial_blocked_fragment_ids,
913 );
914 }
915 }
916 }
917 }
918
919 let mut blocked_fragment_ids = initial_blocked_fragment_ids.clone();
920 if !initial_blocked_fragment_ids.is_empty() {
921 let no_shuffle_relations =
922 self.collect_no_shuffle_fragment_relations_for_reschedule_check();
923 let (forward_edges, backward_edges) =
924 build_no_shuffle_fragment_graph_edges(no_shuffle_relations);
925 let initial_blocked_fragment_ids: Vec<_> =
926 initial_blocked_fragment_ids.iter().copied().collect();
927 for ensemble in find_no_shuffle_graphs(
928 &initial_blocked_fragment_ids,
929 &forward_edges,
930 &backward_edges,
931 )? {
932 blocked_fragment_ids.extend(ensemble.fragments());
933 }
934 }
935
936 let mut blocked_job_ids = HashSet::new();
937 blocked_job_ids.extend(
938 blocked_fragment_ids
939 .into_iter()
940 .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id)),
941 );
942 Ok(blocked_job_ids)
943 }
944
945 fn collect_reschedule_blocked_job_ids(
946 &self,
947 reschedules: &HashMap<FragmentId, Reschedule>,
948 fragment_actors: &HashMap<FragmentId, HashSet<ActorId>>,
949 blocked_job_ids: &HashSet<JobId>,
950 ) -> HashSet<JobId> {
951 let mut affected_fragment_ids: HashSet<FragmentId> = reschedules.keys().copied().collect();
952 affected_fragment_ids.extend(fragment_actors.keys().copied());
953 for reschedule in reschedules.values() {
954 affected_fragment_ids.extend(reschedule.downstream_fragment_ids.iter().copied());
955 affected_fragment_ids.extend(
956 reschedule
957 .upstream_fragment_dispatcher_ids
958 .iter()
959 .map(|(fragment_id, _)| *fragment_id),
960 );
961 }
962
963 affected_fragment_ids
964 .into_iter()
965 .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id))
966 .filter(|job_id| blocked_job_ids.contains(job_id))
967 .collect()
968 }
969
970 fn next_complete_barrier_task(
971 &mut self,
972 periodic_barriers: &mut PeriodicBarriers,
973 partial_graph_manager: &mut PartialGraphManager,
974 task: &mut Option<CompleteBarrierTask>,
975 hummock_version_stats: &HummockVersionStats,
976 ) {
977 let mut independent_jobs_task = vec![];
979 if let Some(committed_epoch) = self.committed_epoch {
980 let mut finished_jobs = Vec::new();
982 let min_upstream_inflight_barrier = partial_graph_manager
983 .first_inflight_barrier(self.partial_graph_id)
984 .map(|epoch| epoch.prev);
985 for (job_id, job) in &mut self.independent_checkpoint_job_controls {
986 match job {
987 IndependentCheckpointJobControl::CreatingStreamingJob(creating_job) => {
988 if let Some((epoch, resps, info, is_finish_epoch)) = creating_job
989 .start_completing(
990 partial_graph_manager,
991 min_upstream_inflight_barrier,
992 committed_epoch,
993 )
994 {
995 let resps = resps.into_values().collect_vec();
996 if is_finish_epoch {
997 assert!(info.notifier.is_none());
998 finished_jobs.push((*job_id, epoch, resps));
999 continue;
1000 };
1001 independent_jobs_task.push((*job_id, epoch, resps, info));
1002 }
1003 }
1004 IndependentCheckpointJobControl::BatchRefresh(batch_refresh_job) => {
1005 if let Some((epoch, resps, info, tracking_job)) =
1006 batch_refresh_job.start_completing(partial_graph_manager)
1007 {
1008 let resps = resps.into_values().collect_vec();
1009 if let Some(tracking_job) = tracking_job {
1010 let task = task.get_or_insert_default();
1011 task.finished_jobs.push(tracking_job);
1012 }
1013 independent_jobs_task.push((*job_id, epoch, resps, info));
1014 }
1015 }
1016 }
1017 }
1018 if !finished_jobs.is_empty() {
1019 partial_graph_manager.remove_partial_graphs(
1020 finished_jobs
1021 .iter()
1022 .map(|(job_id, ..)| to_partial_graph_id(self.database_id, Some(*job_id)))
1023 .collect(),
1024 );
1025 }
1026 for (job_id, epoch, resps) in finished_jobs {
1027 debug!(epoch, %job_id, "finish creating job");
1028 let Some(IndependentCheckpointJobControl::CreatingStreamingJob(
1031 creating_streaming_job,
1032 )) = self.independent_checkpoint_job_controls.remove(&job_id)
1033 else {
1034 panic!("finished job {job_id} should be a creating streaming job");
1035 };
1036 let tracking_job = creating_streaming_job.into_tracking_job();
1037 self.finishing_jobs_collector
1038 .collect(epoch, job_id, (resps, tracking_job));
1039 }
1040 }
1041 let mut observed_non_checkpoint = false;
1042 self.finishing_jobs_collector.advance_collected();
1043 let epoch_end_bound = self
1044 .finishing_jobs_collector
1045 .first_inflight_epoch()
1046 .map_or(Unbounded, |epoch| Excluded(epoch.prev));
1047 if let Some((epoch, resps, info)) = partial_graph_manager.start_completing(
1048 self.partial_graph_id,
1049 epoch_end_bound,
1050 |_, resps, post_collect_command| {
1051 observed_non_checkpoint = true;
1052 self.handle_refresh_table_info(task, &resps);
1053 self.database_info.apply_collected_command(
1054 &post_collect_command,
1055 &resps,
1056 hummock_version_stats,
1057 );
1058 },
1059 ) {
1060 self.handle_refresh_table_info(task, &resps);
1061 self.database_info.apply_collected_command(
1062 &info.post_collect_command,
1063 &resps,
1064 hummock_version_stats,
1065 );
1066 let mut resps_to_commit = resps.into_values().collect_vec();
1067 let mut staging_commit_info = self.database_info.take_staging_commit_info();
1068 if let Some((_, finished_jobs, _)) =
1069 self.finishing_jobs_collector
1070 .take_collected_if(|collected_epoch| {
1071 assert!(epoch <= collected_epoch.prev);
1072 epoch == collected_epoch.prev
1073 })
1074 {
1075 finished_jobs
1076 .into_iter()
1077 .for_each(|(_, (resps, tracking_job))| {
1078 resps_to_commit.extend(resps);
1079 staging_commit_info.finished_jobs.push(tracking_job);
1080 });
1081 }
1082 {
1083 let task = task.get_or_insert_default();
1084 Command::collect_commit_epoch_info(
1085 &self.database_info,
1086 &info,
1087 task,
1088 resps_to_commit,
1089 self.collect_backfill_pinned_upstream_log_epoch(),
1090 );
1091 self.completing_barrier = Some(info.barrier_info.epoch());
1092 task.finished_jobs.extend(staging_commit_info.finished_jobs);
1093 task.finished_cdc_table_backfill
1094 .extend(staging_commit_info.finished_cdc_table_backfill);
1095 task.epoch_infos
1096 .try_insert(self.partial_graph_id, info)
1097 .expect("non duplicate");
1098 task.commit_info
1099 .truncate_tables
1100 .extend(staging_commit_info.table_ids_to_truncate);
1101 }
1102 } else if observed_non_checkpoint
1103 && self.database_info.has_pending_finished_jobs()
1104 && !partial_graph_manager.has_pending_checkpoint_barrier(self.partial_graph_id)
1105 {
1106 periodic_barriers.force_checkpoint_in_next_barrier(self.database_id);
1107 }
1108 if !independent_jobs_task.is_empty() {
1109 let task = task.get_or_insert_default();
1110 for (job_id, epoch, resps, info) in independent_jobs_task {
1111 collect_independent_job_commit_epoch_info(task, epoch, resps, &info);
1112 task.epoch_infos
1113 .try_insert(to_partial_graph_id(self.database_id, Some(job_id)), info)
1114 .expect("non duplicate");
1115 }
1116 }
1117 }
1118
1119 fn ack_completed(
1120 &mut self,
1121 partial_graph_manager: &mut PartialGraphManager,
1122 command_prev_epoch: Option<u64>,
1123 independent_job_epochs: Vec<(JobId, u64)>,
1124 ) {
1125 {
1126 if let Some(epoch) = self.completing_barrier.take() {
1127 assert_eq!(command_prev_epoch, Some(epoch.prev));
1128 self.committed_epoch = Some(epoch.prev);
1129 partial_graph_manager.ack_completed(self.partial_graph_id, epoch.prev);
1130 self.last_committed_barrier_time
1131 .get_or_insert_with(|| {
1132 GLOBAL_META_METRICS
1133 .last_committed_barrier_time
1134 .with_guarded_label_values(&[&self.database_id.to_string()])
1135 })
1136 .set(Epoch(epoch.curr).as_unix_secs() as i64);
1137 } else {
1138 assert_eq!(command_prev_epoch, None);
1139 };
1140 for (job_id, epoch) in independent_job_epochs {
1141 if let Some(job) = self.independent_checkpoint_job_controls.get_mut(&job_id) {
1142 job.ack_completed(partial_graph_manager, epoch);
1143 }
1144 }
1147 }
1148 }
1149
1150 fn handle_refresh_table_info(
1151 &self,
1152 task: &mut Option<CompleteBarrierTask>,
1153 resps: &HashMap<WorkerId, BarrierCompleteResponse>,
1154 ) {
1155 let list_finished_info = resps
1156 .values()
1157 .flat_map(|resp| resp.list_finished_sources.clone())
1158 .collect::<Vec<_>>();
1159 if !list_finished_info.is_empty() {
1160 let task = task.get_or_insert_default();
1161 task.list_finished_source_ids.extend(list_finished_info);
1162 }
1163
1164 let load_finished_info = resps
1165 .values()
1166 .flat_map(|resp| resp.load_finished_sources.clone())
1167 .collect::<Vec<_>>();
1168 if !load_finished_info.is_empty() {
1169 let task = task.get_or_insert_default();
1170 task.load_finished_source_ids.extend(load_finished_info);
1171 }
1172
1173 let refresh_finished_table_ids: Vec<JobId> = resps
1174 .values()
1175 .flat_map(|resp| {
1176 resp.refresh_finished_tables
1177 .iter()
1178 .map(|table_id| table_id.as_job_id())
1179 })
1180 .collect::<Vec<_>>();
1181 if !refresh_finished_table_ids.is_empty() {
1182 let task = task.get_or_insert_default();
1183 task.refresh_finished_table_job_ids
1184 .extend(refresh_finished_table_ids);
1185 }
1186 }
1187}
1188
1189impl DatabaseCheckpointControl {
1190 fn handle_new_barrier(
1192 &mut self,
1193 command: Option<(Command, Notifier)>,
1194 checkpoint: bool,
1195 span: tracing::Span,
1196 partial_graph_manager: &mut PartialGraphManager,
1197 hummock_version_stats: &HummockVersionStats,
1198 worker_nodes: &HashMap<WorkerId, WorkerNode>,
1199 ) -> MetaResult<()> {
1200 let curr_epoch = self.state.in_flight_prev_epoch().next();
1201
1202 let (mut command, notifier) = if let Some((command, notifier)) = command {
1203 (Some(command), Some(notifier))
1204 } else {
1205 (None, None)
1206 };
1207
1208 debug_assert!(
1209 !matches!(
1210 &command,
1211 Some(Command::RescheduleIntent {
1212 reschedule_plan: None,
1213 ..
1214 })
1215 ),
1216 "reschedule intent should be resolved before injection"
1217 );
1218
1219 let mut notifier_start = notifier.map(Notifier::start);
1220 if let Some(Command::DropStreamingJobs {
1221 streaming_job_ids, ..
1222 }) = &mut command
1223 {
1224 streaming_job_ids.retain(|job_id| {
1225 let Some(job) = self.independent_checkpoint_job_controls.get_mut(job_id) else {
1226 return true;
1227 };
1228 !job.drop(notifier_start.as_mut(), partial_graph_manager)
1229 });
1230 if streaming_job_ids.is_empty() {
1231 if let Some(notifier) = notifier_start {
1232 notifier.started();
1233 }
1234 return Ok(());
1235 }
1236 }
1237
1238 if let Some(Command::RescheduleIntent {
1239 reschedule_plan: Some(reschedule_plan),
1240 ..
1241 }) = &command
1242 && !self.independent_checkpoint_job_controls.is_empty()
1243 {
1244 let blocked_job_ids =
1245 self.collect_reschedule_blocked_jobs_for_independent_jobs_inflight()?;
1246 let blocked_reschedule_job_ids = self.collect_reschedule_blocked_job_ids(
1247 &reschedule_plan.reschedules,
1248 &reschedule_plan.fragment_actors,
1249 &blocked_job_ids,
1250 );
1251 if !blocked_reschedule_job_ids.is_empty() {
1252 warn!(
1253 blocked_reschedule_job_ids = ?blocked_reschedule_job_ids,
1254 "reject reschedule fragments related to creating unreschedulable backfill jobs"
1255 );
1256 if let Some(notifier) = notifier_start {
1257 notifier.notify_start_failed(
1258 anyhow!(
1259 "cannot reschedule jobs {:?} when creating jobs with unreschedulable backfill fragments",
1260 blocked_reschedule_job_ids
1261 )
1262 .into(),
1263 );
1264 }
1265 return Ok(());
1266 }
1267 }
1268
1269 if !matches!(&command, Some(Command::CreateStreamingJob { .. }))
1270 && self.database_info.is_empty()
1271 {
1272 assert!(
1273 self.independent_checkpoint_job_controls.is_empty(),
1274 "should not have snapshot backfill job when there is no normal job in database"
1275 );
1276 self.last_committed_barrier_time = None;
1278 if let Some(notifier) = notifier_start {
1280 notifier.started();
1281 }
1282 return Ok(());
1283 };
1284
1285 if let Some(Command::CreateStreamingJob {
1286 job_type:
1287 CreateStreamingJobType::SnapshotBackfill { .. }
1288 | CreateStreamingJobType::BatchRefresh(_),
1289 ..
1290 }) = &command
1291 && self.state.is_paused()
1292 {
1293 warn!("cannot create streaming job with snapshot backfill when paused");
1294 if let Some(notifier) = notifier_start {
1295 notifier.notify_start_failed(
1296 anyhow!("cannot create streaming job with snapshot backfill when paused",)
1297 .into(),
1298 );
1299 }
1300 return Ok(());
1301 }
1302
1303 let barrier_info = self.state.next_barrier_info(checkpoint, curr_epoch);
1304 barrier_info.prev_epoch.span().in_scope(|| {
1306 tracing::info!(target: "rw_tracing", epoch = barrier_info.curr_epoch(), "new barrier enqueued");
1307 });
1308 span.record("epoch", barrier_info.curr_epoch());
1309
1310 let epoch = barrier_info.epoch();
1311 let ApplyCommandInfo { jobs_to_wait } = match self.apply_command(
1312 command,
1313 &mut notifier_start,
1314 barrier_info,
1315 partial_graph_manager,
1316 hummock_version_stats,
1317 worker_nodes,
1318 ) {
1319 Ok(info) => {
1320 assert!(notifier_start.is_none());
1321 info
1322 }
1323 Err(err) => {
1324 if let Some(notifier) = notifier_start {
1325 notifier.notify_start_failed(err.clone());
1326 }
1327 fail_point!("inject_barrier_err_success");
1328 return Err(err);
1329 }
1330 };
1331
1332 self.enqueue_command(epoch, jobs_to_wait);
1334
1335 Ok(())
1336 }
1337
1338 pub(crate) fn get_batch_refresh_trigger_info(&self, job_id: JobId) -> u64 {
1342 let job = self
1343 .independent_checkpoint_job_controls
1344 .get(&job_id)
1345 .expect("batch refresh job should exist");
1346 match job {
1347 IndependentCheckpointJobControl::BatchRefresh(br_job) => br_job
1348 .last_committed_epoch()
1349 .expect("idle job must have a last_committed_epoch"),
1350 _ => panic!("job {} should be a batch refresh job", job_id),
1351 }
1352 }
1353
1354 pub(crate) fn start_batch_refresh_run(
1358 &mut self,
1359 job_id: JobId,
1360 context: &BatchRefreshJobTriggerContext,
1361 worker_nodes: &HashMap<WorkerId, WorkerNode>,
1362 actor_id_counter: &AtomicU32,
1363 partial_graph_manager: &mut PartialGraphManager,
1364 ) -> MetaResult<bool> {
1365 let job = self
1366 .independent_checkpoint_job_controls
1367 .get_mut(&job_id)
1368 .expect("batch refresh job should exist");
1369 match job {
1370 IndependentCheckpointJobControl::BatchRefresh(br_job) => br_job.start_refresh_run(
1371 context,
1372 worker_nodes,
1373 actor_id_counter,
1374 partial_graph_manager,
1375 ),
1376 _ => panic!("job {} should be a batch refresh job", job_id),
1377 }
1378 }
1379}