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