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