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