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::stream_node::NodeBody;
35use risingwave_pb::stream_plan::{DispatcherType as PbDispatcherType, PbSubscriptionUpstreamInfo};
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, IndependentCheckpointJob, 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::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;
63use crate::notification::Notifier;
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::FinishRefresh { .. }
337 | Command::ResetSource { .. }
338 | Command::ResumeBackfill { .. }
339 | Command::InjectSourceOffsets { .. } => {
340 if cfg!(debug_assertions) {
341 panic!(
342 "new database graph info can only be created for normal creating streaming job, but get command: {} {:?}",
343 database_id, command
344 )
345 } else {
346 warn!(%database_id, ?command, "database does not exist while handling the command");
347 notifier.notify_start_failed(anyhow!("database {database_id} does not exist while handling command {command:?}").into());
348 return Ok(());
349 }
350 }
351 },
352 };
353
354 database.handle_new_barrier(
355 Some((command, notifier)),
356 checkpoint,
357 span,
358 partial_graph_manager,
359 &self.hummock_version_stats,
360 worker_nodes,
361 )
362 } else {
363 let database = match self.databases.entry(database_id) {
364 Entry::Occupied(entry) => entry.into_mut(),
365 Entry::Vacant(_) => {
366 return Ok(());
369 }
370 };
371 let Some(database) = database.running_state_mut() else {
372 return Ok(());
374 };
375 if partial_graph_manager.pending_barrier_num(database.partial_graph_id)
376 >= self.in_flight_barrier_nums
377 {
378 return Ok(());
380 }
381 database.handle_new_barrier(
382 None,
383 checkpoint,
384 span,
385 partial_graph_manager,
386 &self.hummock_version_stats,
387 worker_nodes,
388 )
389 }
390 }
391
392 pub(crate) fn gen_backfill_progress(&self) -> HashMap<JobId, BackfillProgress> {
393 let mut progress = HashMap::new();
394 for status in self.databases.values() {
395 let Some(database_checkpoint_control) = status.running_state() else {
396 continue;
397 };
398 progress.extend(
400 database_checkpoint_control
401 .database_info
402 .gen_backfill_progress(),
403 );
404 for (job_id, job) in &database_checkpoint_control.independent_checkpoint_job_controls {
406 if let Some(p) = job.gen_backfill_progress() {
407 progress.insert(*job_id, p);
408 }
409 }
410 }
411 progress
412 }
413
414 pub(crate) fn gen_fragment_backfill_progress(&self) -> Vec<FragmentBackfillProgress> {
415 let mut progress = Vec::new();
416 for status in self.databases.values() {
417 let Some(database_checkpoint_control) = status.running_state() else {
418 continue;
419 };
420 progress.extend(
421 database_checkpoint_control
422 .database_info
423 .gen_fragment_backfill_progress(),
424 );
425 for job in database_checkpoint_control
426 .independent_checkpoint_job_controls
427 .values()
428 {
429 progress.extend(job.gen_fragment_backfill_progress());
430 }
431 }
432 progress
433 }
434
435 pub(crate) fn gen_cdc_progress(&self) -> HashMap<JobId, CdcProgress> {
436 let mut progress = HashMap::new();
437 for status in self.databases.values() {
438 let Some(database_checkpoint_control) = status.running_state() else {
439 continue;
440 };
441 progress.extend(database_checkpoint_control.database_info.gen_cdc_progress());
443 }
444 progress
445 }
446
447 pub(crate) fn databases_failed_at_worker_err(
448 &mut self,
449 worker_id: WorkerId,
450 ) -> impl Iterator<Item = DatabaseId> + '_ {
451 self.databases
452 .iter_mut()
453 .filter_map(
454 move |(database_id, database_status)| match database_status {
455 DatabaseCheckpointControlStatus::Running(control) => {
456 if !control.is_valid_after_worker_err(worker_id) {
457 Some(*database_id)
458 } else {
459 None
460 }
461 }
462 DatabaseCheckpointControlStatus::Recovering(state) => {
463 if !state.is_valid_after_worker_err(worker_id) {
464 Some(*database_id)
465 } else {
466 None
467 }
468 }
469 },
470 )
471 }
472
473 pub(crate) fn get_batch_refresh_trigger_info(
476 &self,
477 database_id: DatabaseId,
478 job_id: JobId,
479 ) -> u64 {
480 let database = self
481 .databases
482 .get(&database_id)
483 .and_then(|s| s.running_state())
484 .expect("database should be running for batch refresh trigger");
485 database.get_batch_refresh_trigger_info(job_id)
486 }
487
488 pub(crate) fn start_batch_refresh_run(
489 &mut self,
490 database_id: DatabaseId,
491 job_id: JobId,
492 context: &BatchRefreshJobTriggerContext,
493 worker_nodes: &HashMap<WorkerId, WorkerNode>,
494 actor_id_counter: &AtomicU32,
495 partial_graph_manager: &mut PartialGraphManager,
496 ) -> MetaResult<bool> {
497 let database = self
498 .databases
499 .get_mut(&database_id)
500 .and_then(|s| s.running_state_mut())
501 .expect("database should be running");
502 database.start_batch_refresh_run(
503 job_id,
504 context,
505 worker_nodes,
506 actor_id_counter,
507 partial_graph_manager,
508 )
509 }
510
511 pub(crate) fn apply_batch_refresh_fragment_infos(
512 &mut self,
513 database_id: DatabaseId,
514 job_id: JobId,
515 ) {
516 let database = self
519 .databases
520 .get_mut(&database_id)
521 .and_then(|s| s.running_state_mut())
522 .expect("database should be running");
523 let br_job = match database
524 .independent_checkpoint_job_controls
525 .get(&job_id)
526 .expect("job should exist")
527 .running()
528 .expect("job should be running")
529 {
530 IndependentCheckpointJob::BatchRefresh(job) => job,
531 _ => panic!("expected batch refresh job"),
532 };
533 if let Some(fragment_infos) = br_job.fragment_infos() {
534 database
535 .database_info
536 .shared_actor_infos
537 .upsert(database_id, fragment_infos.values().map(|f| (f, job_id)));
538 }
539 }
540}
541
542pub(crate) enum CheckpointControlEvent<'a> {
543 EnteringInitializing(DatabaseStatusAction<'a, EnterInitializing>),
544 EnteringRunning(DatabaseStatusAction<'a, EnterRunning>),
545 BatchRefreshTrigger {
548 database_id: DatabaseId,
549 job_id: JobId,
550 },
551}
552
553impl CheckpointControl {
554 pub(crate) fn on_partial_graph_reset(
555 &mut self,
556 partial_graph_id: PartialGraphId,
557 reset_resps: HashMap<WorkerId, ResetPartialGraphResponse>,
558 ) {
559 let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
560 match self.databases.get_mut(&database_id).expect("should exist") {
561 DatabaseCheckpointControlStatus::Running(database) => {
562 if let Some(independent_job_id) = independent_job_id {
563 match database
564 .independent_checkpoint_job_controls
565 .remove(&independent_job_id)
566 {
567 Some(independent_job) => {
568 database
569 .pending_independent_job_subscriptions_to_drop
570 .extend(independent_job.on_partial_graph_reset());
571 }
572 None => {
573 if cfg!(debug_assertions) {
574 panic!(
575 "receive reset partial graph resp on non-existing independent job {independent_job_id} in database {database_id}"
576 )
577 }
578 warn!(
579 %database_id,
580 %independent_job_id,
581 "ignore reset partial graph resp on non-existing independent job on running database"
582 );
583 }
584 }
585 } else {
586 unreachable!("should not receive reset database resp when database running")
587 }
588 }
589 DatabaseCheckpointControlStatus::Recovering(state) => {
590 state.on_partial_graph_reset(partial_graph_id, reset_resps);
591 }
592 }
593 }
594
595 pub(crate) fn on_partial_graph_initialized(
596 &mut self,
597 partial_graph_id: PartialGraphId,
598 partial_graph_manager: &mut PartialGraphManager,
599 ) -> MetaResult<()> {
600 let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
601 match self.databases.get_mut(&database_id).expect("should exist") {
602 DatabaseCheckpointControlStatus::Running(database) => {
603 let Some(independent_job_id) = independent_job_id else {
604 unreachable!("database partial graph should not initialize when running")
605 };
606 let job = database
607 .independent_checkpoint_job_controls
608 .get_mut(&independent_job_id)
609 .expect("independent job should exist");
610 match job.running_mut().expect("job should be running") {
611 IndependentCheckpointJob::BatchRefresh(job) => {
612 job.on_log_store_initialized(partial_graph_manager)
613 }
614 IndependentCheckpointJob::CreatingStreamingJob(_) => {
615 unreachable!("creating streaming job should not initialize when running")
616 }
617 }
618 }
619 DatabaseCheckpointControlStatus::Recovering(state) => {
620 state.partial_graph_initialized(partial_graph_id);
621 Ok(())
622 }
623 }
624 }
625
626 pub(crate) fn next_event(
627 &mut self,
628 ) -> impl Future<Output = CheckpointControlEvent<'_>> + Send + '_ {
629 let mut this = Some(self);
630 poll_fn(move |cx| {
631 let Some(this_mut) = this.as_mut() else {
632 unreachable!("should not be polled after poll ready")
633 };
634 for (&database_id, database_status) in &mut this_mut.databases {
635 match database_status {
636 DatabaseCheckpointControlStatus::Running(database) => {
637 if let Some(committed_epoch) = database.committed_epoch {
639 for (job_id, job) in &database.independent_checkpoint_job_controls {
640 if let Some(IndependentCheckpointJob::BatchRefresh(br_job)) =
641 job.running()
642 && br_job.should_start_refresh(committed_epoch)
643 {
644 let job_id = *job_id;
645 let _ = this.take().expect("checked Some");
646 return Poll::Ready(
647 CheckpointControlEvent::BatchRefreshTrigger {
648 database_id,
649 job_id,
650 },
651 );
652 }
653 }
654 }
655 }
656 DatabaseCheckpointControlStatus::Recovering(state) => {
657 let poll_result = state.poll_next_event(cx);
658 if let Poll::Ready(action) = poll_result {
659 let this = this.take().expect("checked Some");
660 return Poll::Ready(match action {
661 RecoveringStateAction::EnterInitializing(reset_workers) => {
662 CheckpointControlEvent::EnteringInitializing(
663 this.new_database_status_action(
664 database_id,
665 EnterInitializing(reset_workers),
666 ),
667 )
668 }
669 RecoveringStateAction::EnterRunning => {
670 CheckpointControlEvent::EnteringRunning(
671 this.new_database_status_action(database_id, EnterRunning),
672 )
673 }
674 });
675 }
676 }
677 }
678 }
679 Poll::Pending
680 })
681 }
682}
683
684pub(crate) enum DatabaseCheckpointControlStatus {
685 Running(DatabaseCheckpointControl),
686 Recovering(DatabaseRecoveringState),
687}
688
689impl DatabaseCheckpointControlStatus {
690 fn running_state(&self) -> Option<&DatabaseCheckpointControl> {
691 match self {
692 DatabaseCheckpointControlStatus::Running(state) => Some(state),
693 DatabaseCheckpointControlStatus::Recovering(_) => None,
694 }
695 }
696
697 fn running_state_mut(&mut self) -> Option<&mut DatabaseCheckpointControl> {
698 match self {
699 DatabaseCheckpointControlStatus::Running(state) => Some(state),
700 DatabaseCheckpointControlStatus::Recovering(_) => None,
701 }
702 }
703
704 fn expect_running(&mut self, reason: &'static str) -> &mut DatabaseCheckpointControl {
705 match self {
706 DatabaseCheckpointControlStatus::Running(state) => state,
707 DatabaseCheckpointControlStatus::Recovering(_) => {
708 panic!("should be at running: {}", reason)
709 }
710 }
711 }
712}
713
714pub(in crate::barrier) struct DatabaseCheckpointControlMetrics {
715 barrier_latency: LabelGuardedHistogram,
716 in_flight_barrier_nums: LabelGuardedIntGauge,
717 all_barrier_nums: LabelGuardedIntGauge,
718}
719
720impl DatabaseCheckpointControlMetrics {
721 pub(in crate::barrier) fn new(database_id: DatabaseId) -> Self {
722 let database_id_str = database_id.to_string();
723 let barrier_latency = GLOBAL_META_METRICS
724 .barrier_latency
725 .with_guarded_label_values(&[&database_id_str]);
726 let in_flight_barrier_nums = GLOBAL_META_METRICS
727 .in_flight_barrier_nums
728 .with_guarded_label_values(&[&database_id_str]);
729 let all_barrier_nums = GLOBAL_META_METRICS
730 .all_barrier_nums
731 .with_guarded_label_values(&[&database_id_str]);
732 Self {
733 barrier_latency,
734 in_flight_barrier_nums,
735 all_barrier_nums,
736 }
737 }
738}
739
740impl PartialGraphStat for DatabaseCheckpointControlMetrics {
741 fn observe_barrier_latency(&self, _epoch: EpochPair, barrier_latency_secs: f64) {
742 self.barrier_latency.observe(barrier_latency_secs);
743 }
744
745 fn observe_barrier_num(&self, inflight_barrier_num: usize, collected_barrier_num: usize) {
746 self.in_flight_barrier_nums.set(inflight_barrier_num as _);
747 self.all_barrier_nums
748 .set((inflight_barrier_num + collected_barrier_num) as _);
749 }
750}
751
752pub(in crate::barrier) struct DatabaseCheckpointControl {
754 pub(super) database_id: DatabaseId,
755 pub(super) term_id: String,
757 partial_graph_id: PartialGraphId,
758 pub(super) state: BarrierWorkerState,
759
760 finishing_jobs_collector:
761 BarrierItemCollector<JobId, (Vec<BarrierCompleteResponse>, TrackingJob), ()>,
762 completing_barrier: Option<EpochPair>,
764
765 committed_epoch: Option<u64>,
766
767 last_committed_barrier_time: Option<LabelGuardedIntGauge>,
770
771 pub(super) database_info: InflightDatabaseInfo,
772 pub independent_checkpoint_job_controls: HashMap<JobId, IndependentCheckpointJobControl>,
773 pub(super) pending_independent_job_subscriptions_to_drop: Vec<PbSubscriptionUpstreamInfo>,
774}
775
776impl DatabaseCheckpointControl {
777 fn new(database_id: DatabaseId, shared_actor_infos: SharedActorInfos) -> Self {
778 Self {
779 database_id,
780 term_id: Uuid::new_v4().to_string(),
781 partial_graph_id: to_partial_graph_id(database_id, None),
782 state: BarrierWorkerState::new(),
783 finishing_jobs_collector: BarrierItemCollector::new(false),
784 completing_barrier: None,
785 committed_epoch: None,
786 last_committed_barrier_time: None,
787 database_info: InflightDatabaseInfo::empty(database_id, shared_actor_infos),
788 independent_checkpoint_job_controls: Default::default(),
789 pending_independent_job_subscriptions_to_drop: Default::default(),
790 }
791 }
792
793 pub(crate) fn recovery(
794 database_id: DatabaseId,
795 term_id: String,
796 state: BarrierWorkerState,
797 committed_epoch: u64,
798 database_info: InflightDatabaseInfo,
799 independent_checkpoint_job_controls: HashMap<JobId, IndependentCheckpointJobControl>,
800 ) -> Self {
801 Self {
802 database_id,
803 term_id,
804 partial_graph_id: to_partial_graph_id(database_id, None),
805 state,
806 finishing_jobs_collector: BarrierItemCollector::new(false),
807 completing_barrier: None,
808 committed_epoch: Some(committed_epoch),
809 last_committed_barrier_time: None,
810 database_info,
811 independent_checkpoint_job_controls,
812 pending_independent_job_subscriptions_to_drop: Default::default(),
813 }
814 }
815
816 pub(super) fn term_id(&self) -> &str {
817 &self.term_id
818 }
819
820 pub(crate) fn is_valid_after_worker_err(&self, worker_id: WorkerId) -> bool {
821 !self.database_info.contains_worker(worker_id as _)
822 && self
823 .independent_checkpoint_job_controls
824 .values()
825 .all(|job| {
826 job.fragment_infos()
827 .map(|fragment_infos| {
828 !InflightFragmentInfo::contains_worker(
829 fragment_infos.values(),
830 worker_id,
831 )
832 })
833 .unwrap_or(true)
834 })
835 }
836
837 fn enqueue_command(&mut self, epoch: EpochPair, independent_jobs_to_wait: HashSet<JobId>) {
839 let prev_epoch = epoch.prev;
840 tracing::trace!(prev_epoch, ?independent_jobs_to_wait, "enqueue command");
841 if !independent_jobs_to_wait.is_empty() {
842 self.finishing_jobs_collector
843 .enqueue(epoch, independent_jobs_to_wait, ());
844 }
845 }
846
847 fn barrier_collected(
850 &mut self,
851 partial_graph_id: PartialGraphId,
852 collected_barrier: CollectedBarrier<'_>,
853 periodic_barriers: &mut PeriodicBarriers,
854 ) -> MetaResult<()> {
855 let prev_epoch = collected_barrier.epoch.prev;
856 tracing::trace!(
857 prev_epoch,
858 partial_graph_id = %partial_graph_id,
859 "barrier collected"
860 );
861 let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
862 assert_eq!(self.database_id, database_id);
863 if let Some(independent_job_id) = independent_job_id {
864 let job = self
865 .independent_checkpoint_job_controls
866 .get_mut(&independent_job_id)
867 .expect("should exist");
868 let should_force_checkpoint = job.collect(collected_barrier);
869 if should_force_checkpoint {
870 periodic_barriers.force_checkpoint_in_next_barrier(self.database_id);
871 }
872 }
873 Ok(())
874 }
875}
876
877impl DatabaseCheckpointControl {
878 fn collect_backfill_pinned_upstream_tables(&self) -> HashSet<TableId> {
879 self.independent_checkpoint_job_controls
880 .values()
881 .flat_map(|job| job.pinned_upstream_tables().iter().copied())
882 .collect()
883 }
884
885 fn collect_no_shuffle_fragment_relations_for_reschedule_check(
886 &self,
887 ) -> Vec<(FragmentId, FragmentId)> {
888 let mut no_shuffle_relations = Vec::new();
889 for fragment in self.database_info.fragment_infos() {
890 let downstream_fragment_id = fragment.fragment_id;
891 visit_stream_node_cont(&fragment.nodes, |node| {
892 if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
893 && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
894 {
895 no_shuffle_relations.push((merge.upstream_fragment_id, downstream_fragment_id));
896 }
897 true
898 });
899 }
900
901 for job in self.independent_checkpoint_job_controls.values() {
902 if let Some(fragment_infos) = job.fragment_infos() {
903 for fragment_info in fragment_infos.values() {
904 let downstream_fragment_id = fragment_info.fragment_id;
905 visit_stream_node_cont(&fragment_info.nodes, |node| {
906 if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
907 && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
908 {
909 no_shuffle_relations
910 .push((merge.upstream_fragment_id, downstream_fragment_id));
911 }
912 true
913 });
914 }
915 }
916 }
917 no_shuffle_relations
918 }
919
920 fn collect_reschedule_blocked_jobs_for_independent_jobs_inflight(
921 &self,
922 ) -> MetaResult<HashSet<JobId>> {
923 let mut initial_blocked_fragment_ids = HashSet::new();
924 for job in self.independent_checkpoint_job_controls.values() {
925 if let Some(fragment_infos) = job.fragment_infos() {
926 for fragment_info in fragment_infos.values() {
927 if fragment_has_online_unreschedulable_scan(fragment_info) {
928 initial_blocked_fragment_ids.insert(fragment_info.fragment_id);
929 collect_fragment_upstream_fragment_ids(
930 fragment_info,
931 &mut initial_blocked_fragment_ids,
932 );
933 }
934 }
935 }
936 }
937
938 let mut blocked_fragment_ids = initial_blocked_fragment_ids.clone();
939 if !initial_blocked_fragment_ids.is_empty() {
940 let no_shuffle_relations =
941 self.collect_no_shuffle_fragment_relations_for_reschedule_check();
942 let (forward_edges, backward_edges) =
943 build_no_shuffle_fragment_graph_edges(no_shuffle_relations);
944 let initial_blocked_fragment_ids: Vec<_> =
945 initial_blocked_fragment_ids.iter().copied().collect();
946 for ensemble in find_no_shuffle_graphs(
947 &initial_blocked_fragment_ids,
948 &forward_edges,
949 &backward_edges,
950 )? {
951 blocked_fragment_ids.extend(ensemble.fragments());
952 }
953 }
954
955 let mut blocked_job_ids = HashSet::new();
956 blocked_job_ids.extend(
957 blocked_fragment_ids
958 .into_iter()
959 .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id)),
960 );
961 Ok(blocked_job_ids)
962 }
963
964 fn collect_reschedule_blocked_job_ids(
965 &self,
966 reschedules: &HashMap<FragmentId, Reschedule>,
967 fragment_actors: &HashMap<FragmentId, HashSet<ActorId>>,
968 blocked_job_ids: &HashSet<JobId>,
969 ) -> HashSet<JobId> {
970 let mut affected_fragment_ids: HashSet<FragmentId> = reschedules.keys().copied().collect();
971 affected_fragment_ids.extend(fragment_actors.keys().copied());
972 for reschedule in reschedules.values() {
973 affected_fragment_ids.extend(reschedule.downstream_fragment_ids.iter().copied());
974 affected_fragment_ids.extend(
975 reschedule
976 .upstream_fragment_dispatcher_ids
977 .iter()
978 .map(|(fragment_id, _)| *fragment_id),
979 );
980 }
981
982 affected_fragment_ids
983 .into_iter()
984 .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id))
985 .filter(|job_id| blocked_job_ids.contains(job_id))
986 .collect()
987 }
988
989 fn next_complete_barrier_task(
990 &mut self,
991 periodic_barriers: &mut PeriodicBarriers,
992 partial_graph_manager: &mut PartialGraphManager,
993 task: &mut Option<CompleteBarrierTask>,
994 hummock_version_stats: &HummockVersionStats,
995 ) {
996 let mut independent_jobs_task = vec![];
998 let mut finished_jobs = Vec::new();
999 let min_upstream_inflight_barrier = partial_graph_manager
1000 .first_inflight_barrier(self.partial_graph_id)
1001 .map(|epoch| epoch.prev);
1002 for (job_id, job) in &mut self.independent_checkpoint_job_controls {
1003 let Some(job) = job.ready_mut() else {
1004 continue;
1005 };
1006 match job {
1007 IndependentCheckpointJob::CreatingStreamingJob(creating_job) => {
1008 if let Some((epoch, resps, info, is_finish_epoch)) = creating_job
1009 .start_completing(partial_graph_manager, min_upstream_inflight_barrier)
1010 {
1011 let resps = resps.into_values().collect_vec();
1012 if is_finish_epoch {
1013 assert!(info.notifier.is_none());
1014 finished_jobs.push((*job_id, epoch, resps));
1015 continue;
1016 };
1017 independent_jobs_task.push((*job_id, epoch, resps, info));
1018 }
1019 }
1020 IndependentCheckpointJob::BatchRefresh(batch_refresh_job) => {
1021 if let Some((epoch, resps, info, tracking_job)) =
1022 batch_refresh_job.start_completing(partial_graph_manager)
1023 {
1024 let resps = resps.into_values().collect_vec();
1025 if let Some(tracking_job) = tracking_job {
1026 let task = task.get_or_insert_default();
1027 task.finished_jobs.push(tracking_job);
1028 }
1029 independent_jobs_task.push((*job_id, epoch, resps, info));
1030 }
1031 }
1032 }
1033 }
1034 if !finished_jobs.is_empty() {
1035 partial_graph_manager.remove_partial_graphs(
1036 finished_jobs
1037 .iter()
1038 .map(|(job_id, ..)| to_partial_graph_id(self.database_id, Some(*job_id)))
1039 .collect(),
1040 );
1041 }
1042 for (job_id, epoch, resps) in finished_jobs {
1043 debug!(epoch, %job_id, "finish creating job");
1044 let Some(IndependentCheckpointJobControl::Running {
1049 job: IndependentCheckpointJob::CreatingStreamingJob(creating_streaming_job),
1050 ..
1051 }) = self.independent_checkpoint_job_controls.remove(&job_id)
1052 else {
1053 panic!("finished job {job_id} should be a creating streaming job");
1054 };
1055 let tracking_job = creating_streaming_job.into_tracking_job();
1056 self.finishing_jobs_collector
1057 .collect(epoch, job_id, (resps, tracking_job));
1058 }
1059 let mut observed_non_checkpoint = false;
1060 self.finishing_jobs_collector.advance_collected();
1061 let epoch_end_bound = self
1062 .finishing_jobs_collector
1063 .first_inflight_epoch()
1064 .map_or(Unbounded, |epoch| Excluded(epoch.prev));
1065 if let Some((epoch, resps, info)) = partial_graph_manager.start_completing(
1066 self.partial_graph_id,
1067 epoch_end_bound,
1068 |_, resps, post_collect_command| {
1069 observed_non_checkpoint = true;
1070 self.handle_refresh_table_info(task, &resps);
1071 self.database_info.apply_collected_command(
1072 &post_collect_command,
1073 &resps,
1074 hummock_version_stats,
1075 );
1076 },
1077 ) {
1078 self.handle_refresh_table_info(task, &resps);
1079 self.database_info.apply_collected_command(
1080 &info.post_collect_command,
1081 &resps,
1082 hummock_version_stats,
1083 );
1084 let mut resps_to_commit = resps.into_values().collect_vec();
1085 let mut staging_commit_info = self.database_info.take_staging_commit_info();
1086 if let Some((_, finished_jobs, _)) =
1087 self.finishing_jobs_collector
1088 .take_collected_if(|collected_epoch| {
1089 assert!(epoch <= collected_epoch.prev);
1090 epoch == collected_epoch.prev
1091 })
1092 {
1093 finished_jobs
1094 .into_iter()
1095 .for_each(|(_, (resps, tracking_job))| {
1096 resps_to_commit.extend(resps);
1097 staging_commit_info.finished_jobs.push(tracking_job);
1098 });
1099 }
1100 {
1101 let task = task.get_or_insert_default();
1102 Command::collect_commit_epoch_info(
1103 &self.database_info,
1104 &info,
1105 task,
1106 resps_to_commit,
1107 self.collect_backfill_pinned_upstream_tables(),
1108 );
1109 self.completing_barrier = Some(info.barrier_info.epoch());
1110 task.finished_jobs.extend(staging_commit_info.finished_jobs);
1111 task.finished_cdc_table_backfill
1112 .extend(staging_commit_info.finished_cdc_table_backfill);
1113 task.epoch_infos
1114 .try_insert(self.partial_graph_id, info)
1115 .expect("non duplicate");
1116 task.commit_info
1117 .truncate_tables
1118 .extend(staging_commit_info.table_ids_to_truncate);
1119 }
1120 } else if observed_non_checkpoint
1121 && self.database_info.has_pending_finished_jobs()
1122 && !partial_graph_manager.has_pending_checkpoint_barrier(self.partial_graph_id)
1123 {
1124 periodic_barriers.force_checkpoint_in_next_barrier(self.database_id);
1125 }
1126 if !independent_jobs_task.is_empty() {
1127 let task = task.get_or_insert_default();
1128 for (job_id, epoch, resps, info) in independent_jobs_task {
1129 collect_independent_job_commit_epoch_info(task, epoch, resps, &info);
1130 task.epoch_infos
1131 .try_insert(to_partial_graph_id(self.database_id, Some(job_id)), info)
1132 .expect("non duplicate");
1133 }
1134 }
1135 }
1136
1137 fn ack_completed(
1138 &mut self,
1139 partial_graph_manager: &mut PartialGraphManager,
1140 command_prev_epoch: Option<u64>,
1141 independent_job_epochs: Vec<(JobId, u64)>,
1142 ) {
1143 {
1144 if let Some(epoch) = self.completing_barrier.take() {
1145 assert_eq!(command_prev_epoch, Some(epoch.prev));
1146 self.committed_epoch = Some(epoch.prev);
1147 partial_graph_manager.ack_completed(self.partial_graph_id, epoch.prev);
1148 for job in self.independent_checkpoint_job_controls.values_mut() {
1149 job.on_upstream_database_ack_completed(epoch.prev);
1150 }
1151 self.last_committed_barrier_time
1152 .get_or_insert_with(|| {
1153 GLOBAL_META_METRICS
1154 .last_committed_barrier_time
1155 .with_guarded_label_values(&[&self.database_id.to_string()])
1156 })
1157 .set(Epoch(epoch.curr).as_unix_secs() as i64);
1158 } else {
1159 assert_eq!(command_prev_epoch, None);
1160 };
1161 for (job_id, epoch) in independent_job_epochs {
1162 if let Some(job) = self.independent_checkpoint_job_controls.get_mut(&job_id) {
1163 job.ack_completed(partial_graph_manager, epoch);
1164 }
1165 }
1168 }
1169 }
1170
1171 fn handle_refresh_table_info(
1172 &self,
1173 task: &mut Option<CompleteBarrierTask>,
1174 resps: &HashMap<WorkerId, BarrierCompleteResponse>,
1175 ) {
1176 let list_finished_info = resps
1177 .values()
1178 .flat_map(|resp| resp.list_finished_sources.clone())
1179 .collect::<Vec<_>>();
1180 if !list_finished_info.is_empty() {
1181 let task = task.get_or_insert_default();
1182 task.list_finished_source_ids.extend(list_finished_info);
1183 }
1184
1185 let load_finished_info = resps
1186 .values()
1187 .flat_map(|resp| resp.load_finished_sources.clone())
1188 .collect::<Vec<_>>();
1189 if !load_finished_info.is_empty() {
1190 let task = task.get_or_insert_default();
1191 task.load_finished_source_ids.extend(load_finished_info);
1192 }
1193
1194 let refresh_finished_actors = resps
1195 .values()
1196 .flat_map(|resp| resp.refresh_finished_actors.clone())
1197 .collect::<Vec<_>>();
1198 if !refresh_finished_actors.is_empty() {
1199 let task = task.get_or_insert_default();
1200 task.refresh_finished_actors.extend(refresh_finished_actors);
1201 }
1202 }
1203}
1204
1205impl DatabaseCheckpointControl {
1206 fn handle_new_barrier(
1208 &mut self,
1209 command: Option<(Command, Notifier)>,
1210 checkpoint: bool,
1211 span: tracing::Span,
1212 partial_graph_manager: &mut PartialGraphManager,
1213 hummock_version_stats: &HummockVersionStats,
1214 worker_nodes: &HashMap<WorkerId, WorkerNode>,
1215 ) -> MetaResult<()> {
1216 let curr_epoch = self.state.in_flight_prev_epoch().next();
1217
1218 let (mut command, notifier) = if let Some((command, notifier)) = command {
1219 (Some(command), Some(notifier))
1220 } else {
1221 (None, None)
1222 };
1223
1224 debug_assert!(
1225 !matches!(
1226 &command,
1227 Some(Command::RescheduleIntent {
1228 reschedule_plan: None,
1229 ..
1230 })
1231 ),
1232 "reschedule intent should be resolved before injection"
1233 );
1234
1235 let mut notifier_start = notifier.map(Notifier::start);
1236 if let Some(Command::DropStreamingJobs {
1237 streaming_job_ids, ..
1238 }) = &mut command
1239 {
1240 streaming_job_ids.retain(|job_id| {
1241 let Some(job) = self.independent_checkpoint_job_controls.get_mut(job_id) else {
1242 return true;
1243 };
1244 !job.drop(notifier_start.as_mut(), partial_graph_manager)
1245 });
1246 if streaming_job_ids.is_empty() {
1247 if let Some(notifier) = notifier_start {
1248 notifier.started();
1249 }
1250 return Ok(());
1251 }
1252 }
1253
1254 if let Some(Command::RescheduleIntent {
1255 reschedule_plan: Some(reschedule_plan),
1256 ..
1257 }) = &command
1258 && !self.independent_checkpoint_job_controls.is_empty()
1259 {
1260 let blocked_job_ids =
1261 self.collect_reschedule_blocked_jobs_for_independent_jobs_inflight()?;
1262 let blocked_reschedule_job_ids = self.collect_reschedule_blocked_job_ids(
1263 &reschedule_plan.reschedules,
1264 &reschedule_plan.fragment_actors,
1265 &blocked_job_ids,
1266 );
1267 if !blocked_reschedule_job_ids.is_empty() {
1268 warn!(
1269 blocked_reschedule_job_ids = ?blocked_reschedule_job_ids,
1270 "reject reschedule fragments related to creating unreschedulable backfill jobs"
1271 );
1272 if let Some(notifier) = notifier_start {
1273 notifier.notify_start_failed(
1274 anyhow!(
1275 "cannot reschedule jobs {:?} when creating jobs with unreschedulable backfill fragments",
1276 blocked_reschedule_job_ids
1277 )
1278 .into(),
1279 );
1280 }
1281 return Ok(());
1282 }
1283 }
1284
1285 if !matches!(&command, Some(Command::CreateStreamingJob { .. }))
1286 && self.database_info.is_empty()
1287 {
1288 assert!(
1289 self.independent_checkpoint_job_controls.is_empty(),
1290 "should not have snapshot backfill job when there is no normal job in database"
1291 );
1292 self.last_committed_barrier_time = None;
1294 if let Some(notifier) = notifier_start {
1296 notifier.started();
1297 }
1298 return Ok(());
1299 };
1300
1301 if let Some(Command::CreateStreamingJob {
1302 job_type: CreateStreamingJobType::Independent { .. },
1303 ..
1304 }) = &command
1305 && self.state.is_paused()
1306 {
1307 warn!("cannot create streaming job with snapshot backfill when paused");
1308 if let Some(notifier) = notifier_start {
1309 notifier.notify_start_failed(
1310 anyhow!("cannot create streaming job with snapshot backfill when paused",)
1311 .into(),
1312 );
1313 }
1314 return Ok(());
1315 }
1316
1317 let barrier_info = self.state.next_barrier_info(checkpoint, curr_epoch);
1318 barrier_info.prev_epoch.span().in_scope(|| {
1320 tracing::info!(target: "rw_tracing", epoch = barrier_info.curr_epoch(), "new barrier enqueued");
1321 });
1322 span.record("epoch", barrier_info.curr_epoch());
1323
1324 let epoch = barrier_info.epoch();
1325 let ApplyCommandInfo { jobs_to_wait } = match self.apply_command(
1326 command,
1327 &mut notifier_start,
1328 barrier_info,
1329 partial_graph_manager,
1330 hummock_version_stats,
1331 worker_nodes,
1332 ) {
1333 Ok(info) => {
1334 assert!(notifier_start.is_none());
1335 info
1336 }
1337 Err(err) => {
1338 if let Some(notifier) = notifier_start {
1339 notifier.notify_start_failed(err.clone());
1340 }
1341 fail_point!("inject_barrier_err_success");
1342 return Err(err);
1343 }
1344 };
1345
1346 self.enqueue_command(epoch, jobs_to_wait);
1348
1349 Ok(())
1350 }
1351
1352 pub(crate) fn get_batch_refresh_trigger_info(&self, job_id: JobId) -> u64 {
1356 let job = self
1359 .independent_checkpoint_job_controls
1360 .get(&job_id)
1361 .expect("batch refresh job should exist")
1362 .running()
1363 .expect("batch refresh job should be running");
1364 match job {
1365 IndependentCheckpointJob::BatchRefresh(br_job) => br_job
1366 .last_committed_epoch()
1367 .expect("idle job must have a last_committed_epoch"),
1368 _ => panic!("job {} should be a batch refresh job", job_id),
1369 }
1370 }
1371
1372 pub(crate) fn start_batch_refresh_run(
1376 &mut self,
1377 job_id: JobId,
1378 context: &BatchRefreshJobTriggerContext,
1379 worker_nodes: &HashMap<WorkerId, WorkerNode>,
1380 actor_id_counter: &AtomicU32,
1381 partial_graph_manager: &mut PartialGraphManager,
1382 ) -> MetaResult<bool> {
1383 let term_id = self.term_id.as_str();
1386 let job = self
1387 .independent_checkpoint_job_controls
1388 .get_mut(&job_id)
1389 .expect("batch refresh job should exist")
1390 .running_mut()
1391 .expect("batch refresh job should be running");
1392 match job {
1393 IndependentCheckpointJob::BatchRefresh(br_job) => br_job.start_refresh_run(
1394 context,
1395 worker_nodes,
1396 actor_id_counter,
1397 term_id,
1398 partial_graph_manager,
1399 ),
1400 _ => panic!("job {} should be a batch refresh job", job_id),
1401 }
1402 }
1403}