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