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 does not exist while handling the command");
344 notifier.notify_start_failed(anyhow!("database {database_id} does not exist while 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_tables(&self) -> HashSet<TableId> {
857 self.independent_checkpoint_job_controls
858 .values()
859 .flat_map(|job| job.pinned_upstream_tables())
860 .collect()
861 }
862
863 fn collect_no_shuffle_fragment_relations_for_reschedule_check(
864 &self,
865 ) -> Vec<(FragmentId, FragmentId)> {
866 let mut no_shuffle_relations = Vec::new();
867 for fragment in self.database_info.fragment_infos() {
868 let downstream_fragment_id = fragment.fragment_id;
869 visit_stream_node_cont(&fragment.nodes, |node| {
870 if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
871 && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
872 {
873 no_shuffle_relations.push((merge.upstream_fragment_id, downstream_fragment_id));
874 }
875 true
876 });
877 }
878
879 for job in self.independent_checkpoint_job_controls.values() {
880 if let Some(fragment_infos) = job.fragment_infos() {
881 for fragment_info in fragment_infos.values() {
882 let downstream_fragment_id = fragment_info.fragment_id;
883 visit_stream_node_cont(&fragment_info.nodes, |node| {
884 if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
885 && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
886 {
887 no_shuffle_relations
888 .push((merge.upstream_fragment_id, downstream_fragment_id));
889 }
890 true
891 });
892 }
893 }
894 }
895 no_shuffle_relations
896 }
897
898 fn collect_reschedule_blocked_jobs_for_independent_jobs_inflight(
899 &self,
900 ) -> MetaResult<HashSet<JobId>> {
901 let mut initial_blocked_fragment_ids = HashSet::new();
902 for job in self.independent_checkpoint_job_controls.values() {
903 if let Some(fragment_infos) = job.fragment_infos() {
904 for fragment_info in fragment_infos.values() {
905 if fragment_has_online_unreschedulable_scan(fragment_info) {
906 initial_blocked_fragment_ids.insert(fragment_info.fragment_id);
907 collect_fragment_upstream_fragment_ids(
908 fragment_info,
909 &mut initial_blocked_fragment_ids,
910 );
911 }
912 }
913 }
914 }
915
916 let mut blocked_fragment_ids = initial_blocked_fragment_ids.clone();
917 if !initial_blocked_fragment_ids.is_empty() {
918 let no_shuffle_relations =
919 self.collect_no_shuffle_fragment_relations_for_reschedule_check();
920 let (forward_edges, backward_edges) =
921 build_no_shuffle_fragment_graph_edges(no_shuffle_relations);
922 let initial_blocked_fragment_ids: Vec<_> =
923 initial_blocked_fragment_ids.iter().copied().collect();
924 for ensemble in find_no_shuffle_graphs(
925 &initial_blocked_fragment_ids,
926 &forward_edges,
927 &backward_edges,
928 )? {
929 blocked_fragment_ids.extend(ensemble.fragments());
930 }
931 }
932
933 let mut blocked_job_ids = HashSet::new();
934 blocked_job_ids.extend(
935 blocked_fragment_ids
936 .into_iter()
937 .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id)),
938 );
939 Ok(blocked_job_ids)
940 }
941
942 fn collect_reschedule_blocked_job_ids(
943 &self,
944 reschedules: &HashMap<FragmentId, Reschedule>,
945 fragment_actors: &HashMap<FragmentId, HashSet<ActorId>>,
946 blocked_job_ids: &HashSet<JobId>,
947 ) -> HashSet<JobId> {
948 let mut affected_fragment_ids: HashSet<FragmentId> = reschedules.keys().copied().collect();
949 affected_fragment_ids.extend(fragment_actors.keys().copied());
950 for reschedule in reschedules.values() {
951 affected_fragment_ids.extend(reschedule.downstream_fragment_ids.iter().copied());
952 affected_fragment_ids.extend(
953 reschedule
954 .upstream_fragment_dispatcher_ids
955 .iter()
956 .map(|(fragment_id, _)| *fragment_id),
957 );
958 }
959
960 affected_fragment_ids
961 .into_iter()
962 .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id))
963 .filter(|job_id| blocked_job_ids.contains(job_id))
964 .collect()
965 }
966
967 fn next_complete_barrier_task(
968 &mut self,
969 periodic_barriers: &mut PeriodicBarriers,
970 partial_graph_manager: &mut PartialGraphManager,
971 task: &mut Option<CompleteBarrierTask>,
972 hummock_version_stats: &HummockVersionStats,
973 ) {
974 let mut independent_jobs_task = vec![];
976 if let Some(committed_epoch) = self.committed_epoch {
977 let mut finished_jobs = Vec::new();
979 let min_upstream_inflight_barrier = partial_graph_manager
980 .first_inflight_barrier(self.partial_graph_id)
981 .map(|epoch| epoch.prev);
982 for (job_id, job) in &mut self.independent_checkpoint_job_controls {
983 match job {
984 IndependentCheckpointJobControl::CreatingStreamingJob(creating_job) => {
985 if let Some((epoch, resps, info, is_finish_epoch)) = creating_job
986 .start_completing(
987 partial_graph_manager,
988 min_upstream_inflight_barrier,
989 committed_epoch,
990 )
991 {
992 let resps = resps.into_values().collect_vec();
993 if is_finish_epoch {
994 assert!(info.notifier.is_none());
995 finished_jobs.push((*job_id, epoch, resps));
996 continue;
997 };
998 independent_jobs_task.push((*job_id, epoch, resps, info));
999 }
1000 }
1001 IndependentCheckpointJobControl::BatchRefresh(batch_refresh_job) => {
1002 if let Some((epoch, resps, info, tracking_job)) = batch_refresh_job
1003 .start_completing(partial_graph_manager, committed_epoch)
1004 {
1005 let resps = resps.into_values().collect_vec();
1006 if let Some(tracking_job) = tracking_job {
1007 let task = task.get_or_insert_default();
1008 task.finished_jobs.push(tracking_job);
1009 }
1010 independent_jobs_task.push((*job_id, epoch, resps, info));
1011 }
1012 }
1013 }
1014 }
1015 if !finished_jobs.is_empty() {
1016 partial_graph_manager.remove_partial_graphs(
1017 finished_jobs
1018 .iter()
1019 .map(|(job_id, ..)| to_partial_graph_id(self.database_id, Some(*job_id)))
1020 .collect(),
1021 );
1022 }
1023 for (job_id, epoch, resps) in finished_jobs {
1024 debug!(epoch, %job_id, "finish creating job");
1025 let Some(IndependentCheckpointJobControl::CreatingStreamingJob(
1028 creating_streaming_job,
1029 )) = self.independent_checkpoint_job_controls.remove(&job_id)
1030 else {
1031 panic!("finished job {job_id} should be a creating streaming job");
1032 };
1033 let tracking_job = creating_streaming_job.into_tracking_job();
1034 self.finishing_jobs_collector
1035 .collect(epoch, job_id, (resps, tracking_job));
1036 }
1037 }
1038 let mut observed_non_checkpoint = false;
1039 self.finishing_jobs_collector.advance_collected();
1040 let epoch_end_bound = self
1041 .finishing_jobs_collector
1042 .first_inflight_epoch()
1043 .map_or(Unbounded, |epoch| Excluded(epoch.prev));
1044 if let Some((epoch, resps, info)) = partial_graph_manager.start_completing(
1045 self.partial_graph_id,
1046 epoch_end_bound,
1047 |_, resps, post_collect_command| {
1048 observed_non_checkpoint = true;
1049 self.handle_refresh_table_info(task, &resps);
1050 self.database_info.apply_collected_command(
1051 &post_collect_command,
1052 &resps,
1053 hummock_version_stats,
1054 );
1055 },
1056 ) {
1057 self.handle_refresh_table_info(task, &resps);
1058 self.database_info.apply_collected_command(
1059 &info.post_collect_command,
1060 &resps,
1061 hummock_version_stats,
1062 );
1063 let mut resps_to_commit = resps.into_values().collect_vec();
1064 let mut staging_commit_info = self.database_info.take_staging_commit_info();
1065 if let Some((_, finished_jobs, _)) =
1066 self.finishing_jobs_collector
1067 .take_collected_if(|collected_epoch| {
1068 assert!(epoch <= collected_epoch.prev);
1069 epoch == collected_epoch.prev
1070 })
1071 {
1072 finished_jobs
1073 .into_iter()
1074 .for_each(|(_, (resps, tracking_job))| {
1075 resps_to_commit.extend(resps);
1076 staging_commit_info.finished_jobs.push(tracking_job);
1077 });
1078 }
1079 {
1080 let task = task.get_or_insert_default();
1081 Command::collect_commit_epoch_info(
1082 &self.database_info,
1083 &info,
1084 task,
1085 resps_to_commit,
1086 self.collect_backfill_pinned_upstream_tables(),
1087 );
1088 self.completing_barrier = Some(info.barrier_info.epoch());
1089 task.finished_jobs.extend(staging_commit_info.finished_jobs);
1090 task.finished_cdc_table_backfill
1091 .extend(staging_commit_info.finished_cdc_table_backfill);
1092 task.epoch_infos
1093 .try_insert(self.partial_graph_id, info)
1094 .expect("non duplicate");
1095 task.commit_info
1096 .truncate_tables
1097 .extend(staging_commit_info.table_ids_to_truncate);
1098 }
1099 } else if observed_non_checkpoint
1100 && self.database_info.has_pending_finished_jobs()
1101 && !partial_graph_manager.has_pending_checkpoint_barrier(self.partial_graph_id)
1102 {
1103 periodic_barriers.force_checkpoint_in_next_barrier(self.database_id);
1104 }
1105 if !independent_jobs_task.is_empty() {
1106 let task = task.get_or_insert_default();
1107 for (job_id, epoch, resps, info) in independent_jobs_task {
1108 collect_independent_job_commit_epoch_info(task, epoch, resps, &info);
1109 task.epoch_infos
1110 .try_insert(to_partial_graph_id(self.database_id, Some(job_id)), info)
1111 .expect("non duplicate");
1112 }
1113 }
1114 }
1115
1116 fn ack_completed(
1117 &mut self,
1118 partial_graph_manager: &mut PartialGraphManager,
1119 command_prev_epoch: Option<u64>,
1120 independent_job_epochs: Vec<(JobId, u64)>,
1121 ) {
1122 {
1123 if let Some(epoch) = self.completing_barrier.take() {
1124 assert_eq!(command_prev_epoch, Some(epoch.prev));
1125 self.committed_epoch = Some(epoch.prev);
1126 partial_graph_manager.ack_completed(self.partial_graph_id, epoch.prev);
1127 self.last_committed_barrier_time
1128 .get_or_insert_with(|| {
1129 GLOBAL_META_METRICS
1130 .last_committed_barrier_time
1131 .with_guarded_label_values(&[&self.database_id.to_string()])
1132 })
1133 .set(Epoch(epoch.curr).as_unix_secs() as i64);
1134 } else {
1135 assert_eq!(command_prev_epoch, None);
1136 };
1137 for (job_id, epoch) in independent_job_epochs {
1138 if let Some(job) = self.independent_checkpoint_job_controls.get_mut(&job_id) {
1139 job.ack_completed(partial_graph_manager, epoch);
1140 }
1141 }
1144 }
1145 }
1146
1147 fn handle_refresh_table_info(
1148 &self,
1149 task: &mut Option<CompleteBarrierTask>,
1150 resps: &HashMap<WorkerId, BarrierCompleteResponse>,
1151 ) {
1152 let list_finished_info = resps
1153 .values()
1154 .flat_map(|resp| resp.list_finished_sources.clone())
1155 .collect::<Vec<_>>();
1156 if !list_finished_info.is_empty() {
1157 let task = task.get_or_insert_default();
1158 task.list_finished_source_ids.extend(list_finished_info);
1159 }
1160
1161 let load_finished_info = resps
1162 .values()
1163 .flat_map(|resp| resp.load_finished_sources.clone())
1164 .collect::<Vec<_>>();
1165 if !load_finished_info.is_empty() {
1166 let task = task.get_or_insert_default();
1167 task.load_finished_source_ids.extend(load_finished_info);
1168 }
1169
1170 let refresh_finished_table_ids: Vec<JobId> = resps
1171 .values()
1172 .flat_map(|resp| {
1173 resp.refresh_finished_tables
1174 .iter()
1175 .map(|table_id| table_id.as_job_id())
1176 })
1177 .collect::<Vec<_>>();
1178 if !refresh_finished_table_ids.is_empty() {
1179 let task = task.get_or_insert_default();
1180 task.refresh_finished_table_job_ids
1181 .extend(refresh_finished_table_ids);
1182 }
1183 }
1184}
1185
1186impl DatabaseCheckpointControl {
1187 fn handle_new_barrier(
1189 &mut self,
1190 command: Option<(Command, Notifier)>,
1191 checkpoint: bool,
1192 span: tracing::Span,
1193 partial_graph_manager: &mut PartialGraphManager,
1194 hummock_version_stats: &HummockVersionStats,
1195 worker_nodes: &HashMap<WorkerId, WorkerNode>,
1196 ) -> MetaResult<()> {
1197 let curr_epoch = self.state.in_flight_prev_epoch().next();
1198
1199 let (mut command, notifier) = if let Some((command, notifier)) = command {
1200 (Some(command), Some(notifier))
1201 } else {
1202 (None, None)
1203 };
1204
1205 debug_assert!(
1206 !matches!(
1207 &command,
1208 Some(Command::RescheduleIntent {
1209 reschedule_plan: None,
1210 ..
1211 })
1212 ),
1213 "reschedule intent should be resolved before injection"
1214 );
1215
1216 let mut notifier_start = notifier.map(Notifier::start);
1217 if let Some(Command::DropStreamingJobs {
1218 streaming_job_ids, ..
1219 }) = &mut command
1220 {
1221 streaming_job_ids.retain(|job_id| {
1222 let Some(job) = self.independent_checkpoint_job_controls.get_mut(job_id) else {
1223 return true;
1224 };
1225 !job.drop(notifier_start.as_mut(), partial_graph_manager)
1226 });
1227 if streaming_job_ids.is_empty() {
1228 if let Some(notifier) = notifier_start {
1229 notifier.started();
1230 }
1231 return Ok(());
1232 }
1233 }
1234
1235 if let Some(Command::RescheduleIntent {
1236 reschedule_plan: Some(reschedule_plan),
1237 ..
1238 }) = &command
1239 && !self.independent_checkpoint_job_controls.is_empty()
1240 {
1241 let blocked_job_ids =
1242 self.collect_reschedule_blocked_jobs_for_independent_jobs_inflight()?;
1243 let blocked_reschedule_job_ids = self.collect_reschedule_blocked_job_ids(
1244 &reschedule_plan.reschedules,
1245 &reschedule_plan.fragment_actors,
1246 &blocked_job_ids,
1247 );
1248 if !blocked_reschedule_job_ids.is_empty() {
1249 warn!(
1250 blocked_reschedule_job_ids = ?blocked_reschedule_job_ids,
1251 "reject reschedule fragments related to creating unreschedulable backfill jobs"
1252 );
1253 if let Some(notifier) = notifier_start {
1254 notifier.notify_start_failed(
1255 anyhow!(
1256 "cannot reschedule jobs {:?} when creating jobs with unreschedulable backfill fragments",
1257 blocked_reschedule_job_ids
1258 )
1259 .into(),
1260 );
1261 }
1262 return Ok(());
1263 }
1264 }
1265
1266 if !matches!(&command, Some(Command::CreateStreamingJob { .. }))
1267 && self.database_info.is_empty()
1268 {
1269 assert!(
1270 self.independent_checkpoint_job_controls.is_empty(),
1271 "should not have snapshot backfill job when there is no normal job in database"
1272 );
1273 self.last_committed_barrier_time = None;
1275 if let Some(notifier) = notifier_start {
1277 notifier.started();
1278 }
1279 return Ok(());
1280 };
1281
1282 if let Some(Command::CreateStreamingJob {
1283 job_type:
1284 CreateStreamingJobType::SnapshotBackfill { .. }
1285 | CreateStreamingJobType::BatchRefresh(_),
1286 ..
1287 }) = &command
1288 && self.state.is_paused()
1289 {
1290 warn!("cannot create streaming job with snapshot backfill when paused");
1291 if let Some(notifier) = notifier_start {
1292 notifier.notify_start_failed(
1293 anyhow!("cannot create streaming job with snapshot backfill when paused",)
1294 .into(),
1295 );
1296 }
1297 return Ok(());
1298 }
1299
1300 let barrier_info = self.state.next_barrier_info(checkpoint, curr_epoch);
1301 barrier_info.prev_epoch.span().in_scope(|| {
1303 tracing::info!(target: "rw_tracing", epoch = barrier_info.curr_epoch(), "new barrier enqueued");
1304 });
1305 span.record("epoch", barrier_info.curr_epoch());
1306
1307 let epoch = barrier_info.epoch();
1308 let ApplyCommandInfo { jobs_to_wait } = match self.apply_command(
1309 command,
1310 &mut notifier_start,
1311 barrier_info,
1312 partial_graph_manager,
1313 hummock_version_stats,
1314 worker_nodes,
1315 ) {
1316 Ok(info) => {
1317 assert!(notifier_start.is_none());
1318 info
1319 }
1320 Err(err) => {
1321 if let Some(notifier) = notifier_start {
1322 notifier.notify_start_failed(err.clone());
1323 }
1324 fail_point!("inject_barrier_err_success");
1325 return Err(err);
1326 }
1327 };
1328
1329 self.enqueue_command(epoch, jobs_to_wait);
1331
1332 Ok(())
1333 }
1334
1335 pub(crate) fn get_batch_refresh_trigger_info(&self, job_id: JobId) -> u64 {
1339 let job = self
1340 .independent_checkpoint_job_controls
1341 .get(&job_id)
1342 .expect("batch refresh job should exist");
1343 match job {
1344 IndependentCheckpointJobControl::BatchRefresh(br_job) => br_job
1345 .last_committed_epoch()
1346 .expect("idle job must have a last_committed_epoch"),
1347 _ => panic!("job {} should be a batch refresh job", job_id),
1348 }
1349 }
1350
1351 pub(crate) fn start_batch_refresh_run(
1355 &mut self,
1356 job_id: JobId,
1357 context: &BatchRefreshJobTriggerContext,
1358 worker_nodes: &HashMap<WorkerId, WorkerNode>,
1359 actor_id_counter: &AtomicU32,
1360 partial_graph_manager: &mut PartialGraphManager,
1361 ) -> MetaResult<bool> {
1362 let job = self
1363 .independent_checkpoint_job_controls
1364 .get_mut(&job_id)
1365 .expect("batch refresh job should exist");
1366 match job {
1367 IndependentCheckpointJobControl::BatchRefresh(br_job) => br_job.start_refresh_run(
1368 context,
1369 worker_nodes,
1370 actor_id_counter,
1371 partial_graph_manager,
1372 ),
1373 _ => panic!("job {} should be a batch refresh job", job_id),
1374 }
1375 }
1376}