1use std::collections::hash_map::Entry;
16use std::collections::{HashMap, HashSet};
17use std::mem::replace;
18use std::pin::pin;
19use std::sync::Arc;
20use std::time::Duration;
21
22use anyhow::anyhow;
23use arc_swap::ArcSwap;
24use futures::{TryFutureExt, pin_mut};
25use itertools::Itertools;
26use risingwave_common::catalog::DatabaseId;
27use risingwave_common::id::JobId;
28use risingwave_common::system_param::PAUSE_ON_NEXT_BOOTSTRAP_KEY;
29use risingwave_common::system_param::reader::SystemParamsRead;
30use risingwave_meta_model::WorkerId;
31use risingwave_pb::common::WorkerNode;
32use risingwave_pb::meta::Recovery;
33use risingwave_pb::meta::subscribe_response::{Info, Operation};
34use thiserror_ext::AsReport;
35use tokio::select;
36use tokio::sync::mpsc;
37use tokio::sync::oneshot::{Receiver, Sender};
38use tokio::task::JoinHandle;
39use tonic::Status;
40use tracing::{Instrument, debug, error, info, warn};
41
42use crate::barrier::checkpoint::{CheckpointControl, CheckpointControlEvent};
43use crate::barrier::complete_task::{BarrierCompleteOutput, CompletingTask};
44use crate::barrier::context::recovery::{RenderedDatabaseRuntimeInfo, render_runtime_info};
45use crate::barrier::context::{GlobalBarrierWorkerContext, GlobalBarrierWorkerContextImpl};
46use crate::barrier::info::InflightDatabaseInfo;
47use crate::barrier::rpc::{
48 DatabaseInitialBarrierCollector, database_partial_graphs, from_partial_graph_id,
49 merge_node_rpc_errors,
50};
51use crate::barrier::schedule::{MarkReadyOptions, PeriodicBarriers};
52use crate::barrier::{
53 BarrierManagerRequest, BarrierManagerStatus, BarrierWorkerRuntimeInfoSnapshot, Command,
54 CreateStreamingJobType, RecoveryReason, RescheduleContext, UpdateDatabaseBarrierRequest,
55 schedule,
56};
57use crate::controller::scale::{materialize_actor_assignments, preview_actor_assignments};
58use crate::error::MetaErrorInner;
59use crate::hummock::HummockManagerRef;
60use crate::manager::iceberg_compaction::IcebergCompactionManagerRef;
61use crate::manager::iceberg_pk_index_sink::IcebergPkIndexSinkManager;
62use crate::manager::sink_coordination::SinkCoordinatorManager;
63use crate::manager::{
64 ActiveStreamingWorkerChange, ActiveStreamingWorkerNodes, LocalNotification, MetaSrvEnv,
65 MetadataManager,
66};
67use crate::rpc::metrics::GLOBAL_META_METRICS;
68use crate::serving::ServingVnodeMappingRef;
69use crate::stream::{
70 GlobalRefreshManagerRef, ScaleControllerRef, SourceManagerRef, build_reschedule_commands,
71 rendered_layout_matches_current,
72};
73use crate::{MetaError, MetaResult};
74
75pub(super) struct GlobalBarrierWorker<C> {
85 enable_recovery: bool,
87
88 periodic_barriers: PeriodicBarriers,
90
91 system_enable_per_database_isolation: bool,
93
94 pub(super) context: Arc<C>,
95
96 env: MetaSrvEnv,
97
98 checkpoint_control: CheckpointControl,
99
100 completing_task: CompletingTask,
103
104 request_rx: mpsc::UnboundedReceiver<BarrierManagerRequest>,
105
106 active_streaming_nodes: ActiveStreamingWorkerNodes,
107
108 partial_graph_manager: PartialGraphManager,
109}
110
111#[cfg(test)]
112mod tests {
113 use std::collections::HashMap;
114
115 use super::*;
116 use crate::barrier::RescheduleContext;
117 use crate::barrier::notifier::Notifier;
118
119 #[tokio::test]
120 async fn test_reschedule_intent_without_workers_notifies_start_failed() {
121 let env = MetaSrvEnv::for_test().await;
122 let database_id = DatabaseId::new(1);
123 let database_info =
124 InflightDatabaseInfo::empty(database_id, env.shared_actor_infos().clone());
125 let (notifier, started_rx) = Notifier::new();
126
127 let new_barrier = schedule::NewBarrier {
128 database_id,
129 command: Some((
130 Command::RescheduleIntent {
131 context: RescheduleContext::empty(),
132 reschedule_plan: None,
133 },
134 notifier,
135 )),
136 span: tracing::Span::none(),
137 checkpoint: false,
138 };
139
140 let result =
141 resolve_reschedule_intent(env, HashMap::new(), Some(&database_info), new_barrier);
142
143 assert!(matches!(result, Ok(None)));
144 let started = started_rx.await.expect("started notifier dropped");
145 assert!(started.is_err());
146 }
147}
148
149impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
150 pub(super) async fn new_inner(
151 env: MetaSrvEnv,
152 request_rx: mpsc::UnboundedReceiver<BarrierManagerRequest>,
153 context: Arc<C>,
154 ) -> Self {
155 let enable_recovery = env.opts.enable_recovery;
156
157 let active_streaming_nodes = ActiveStreamingWorkerNodes::uninitialized();
158
159 let partial_graph_manager = PartialGraphManager::uninitialized(env.clone());
160
161 let reader = env.system_params_reader().await;
162 let system_enable_per_database_isolation = reader.per_database_isolation();
163 let periodic_barriers = PeriodicBarriers::default();
165
166 let checkpoint_control = CheckpointControl::new(env.clone());
167 Self {
168 enable_recovery,
169 periodic_barriers,
170 system_enable_per_database_isolation,
171 context,
172 env,
173 checkpoint_control,
174 completing_task: CompletingTask::None,
175 request_rx,
176 active_streaming_nodes,
177 partial_graph_manager,
178 }
179 }
180}
181
182fn resolve_reschedule_intent(
183 env: MetaSrvEnv,
184 worker_nodes: HashMap<WorkerId, WorkerNode>,
185 database_info: Option<&InflightDatabaseInfo>,
186 mut new_barrier: schedule::NewBarrier,
187) -> MetaResult<Option<schedule::NewBarrier>> {
188 let Some((command, notifier)) = new_barrier.command.take() else {
189 return Ok(Some(new_barrier));
190 };
191
192 match command {
193 Command::RescheduleIntent {
194 context,
195 reschedule_plan,
196 } => {
197 if let Some(reschedule_plan) = reschedule_plan {
198 new_barrier.command = Some((
199 Command::RescheduleIntent {
200 context,
201 reschedule_plan: Some(reschedule_plan),
202 },
203 notifier,
204 ));
205 return Ok(Some(new_barrier));
206 }
207 let span = tracing::info_span!(
208 "resolve_reschedule_intent",
209 database_id = %new_barrier.database_id
210 );
211 let reschedule_plan = {
212 let _guard = span.enter();
213 build_reschedule_from_context(
214 &env,
215 worker_nodes,
216 new_barrier.database_id,
217 context,
218 database_info.ok_or_else(|| {
219 anyhow!(
220 "database {} not found when resolving reschedule intent",
221 new_barrier.database_id
222 )
223 })?,
224 )
225 };
226 match reschedule_plan {
227 Ok(Some(reschedule_plan)) => {
228 new_barrier.command = Some((
229 Command::RescheduleIntent {
230 context: RescheduleContext::empty(),
231 reschedule_plan: Some(reschedule_plan),
232 },
233 notifier,
234 ));
235 Ok(Some(new_barrier))
236 }
237 Ok(None) => {
238 notifier.start().started();
240 Ok(None)
241 }
242 Err(err) => {
243 notifier.notify_start_failed(err);
244 Ok(None)
245 }
246 }
247 }
248 _ => {
249 new_barrier.command = Some((command, notifier));
250 Ok(Some(new_barrier))
251 }
252 }
253}
254
255fn build_reschedule_from_context(
256 env: &MetaSrvEnv,
257 worker_nodes: HashMap<WorkerId, WorkerNode>,
258 database_id: DatabaseId,
259 context: RescheduleContext,
260 database_info: &InflightDatabaseInfo,
261) -> MetaResult<Option<crate::barrier::ReschedulePlan>> {
262 if worker_nodes.is_empty() {
263 return Err(anyhow!("no active streaming workers for reschedule").into());
264 }
265
266 if context.is_empty() {
267 return Ok(None);
268 }
269
270 let all_prev_fragments = database_info
273 .fragment_infos()
274 .map(|fragment| (fragment.fragment_id, fragment))
275 .collect();
276
277 let previewed = preview_actor_assignments(&worker_nodes, &context.loaded)?;
278
279 if rendered_layout_matches_current(&previewed.fragments, &all_prev_fragments)? {
280 return Ok(None);
281 }
282
283 let actor_id_counter = env.actor_id_generator();
284 let rendered = materialize_actor_assignments(actor_id_counter, previewed);
287 let mut commands = build_reschedule_commands(rendered.fragments, context, all_prev_fragments)?;
288 Ok(commands.remove(&database_id))
289}
290
291impl GlobalBarrierWorker<GlobalBarrierWorkerContextImpl> {
292 #[expect(clippy::too_many_arguments)]
294 pub async fn new(
295 scheduled_barriers: schedule::ScheduledBarriers,
296 env: MetaSrvEnv,
297 metadata_manager: MetadataManager,
298 hummock_manager: HummockManagerRef,
299 serving_vnode_mapping: ServingVnodeMappingRef,
300 source_manager: SourceManagerRef,
301 sink_manager: SinkCoordinatorManager,
302 iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
303 iceberg_compaction_manager: IcebergCompactionManagerRef,
304 scale_controller: ScaleControllerRef,
305 request_rx: mpsc::UnboundedReceiver<BarrierManagerRequest>,
306 barrier_scheduler: schedule::BarrierScheduler,
307 refresh_manager: GlobalRefreshManagerRef,
308 ) -> Self {
309 let status = Arc::new(ArcSwap::new(Arc::new(BarrierManagerStatus::Starting)));
310
311 let context = Arc::new(GlobalBarrierWorkerContextImpl::new(
312 scheduled_barriers,
313 status,
314 metadata_manager,
315 hummock_manager,
316 serving_vnode_mapping,
317 source_manager,
318 scale_controller,
319 env.clone(),
320 barrier_scheduler,
321 refresh_manager,
322 sink_manager,
323 iceberg_pk_index_sink_manager,
324 iceberg_compaction_manager,
325 ));
326
327 Self::new_inner(env, request_rx, context).await
328 }
329
330 pub fn start(self) -> (JoinHandle<()>, Sender<()>) {
331 let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
332 let fut = (self.env.await_tree_reg())
333 .register_derived_root("Global Barrier Worker")
334 .instrument(self.run(shutdown_rx));
335 let join_handle = tokio::spawn(fut);
336
337 (join_handle, shutdown_tx)
338 }
339
340 async fn take_pause_on_bootstrap(&mut self) -> MetaResult<bool> {
342 let paused = self
343 .env
344 .system_params_reader()
345 .await
346 .pause_on_next_bootstrap()
347 || self.env.opts.pause_on_next_bootstrap_offline;
348
349 if paused {
350 warn!(
351 "The cluster will bootstrap with all data sources paused as specified by the system parameter `{}`. \
352 It will now be reset to `false`. \
353 To resume the data sources, either restart the cluster again or use `risectl meta resume`.",
354 PAUSE_ON_NEXT_BOOTSTRAP_KEY
355 );
356 self.env
357 .system_params_manager_impl_ref()
358 .set_param(PAUSE_ON_NEXT_BOOTSTRAP_KEY, Some("false".to_owned()))
359 .await?;
360 }
361 Ok(paused)
362 }
363
364 async fn run(mut self, shutdown_rx: Receiver<()>) {
366 tracing::info!(
367 "Starting barrier manager with: enable_recovery={}, in_flight_barrier_nums={}",
368 self.enable_recovery,
369 self.checkpoint_control.in_flight_barrier_nums,
370 );
371
372 if !self.enable_recovery {
373 let job_exist = self
374 .context
375 .metadata_manager
376 .catalog_controller
377 .has_any_streaming_jobs()
378 .await
379 .unwrap();
380 if job_exist {
381 panic!(
382 "Some streaming jobs already exist in meta, please start with recovery enabled \
383 or clean up the metadata using `./risedev clean-data`"
384 );
385 }
386 }
387
388 {
389 let span = tracing::info_span!("bootstrap_recovery");
394 crate::telemetry::report_event(
395 risingwave_pb::telemetry::TelemetryEventStage::Recovery,
396 "normal_recovery",
397 0,
398 None,
399 None,
400 None,
401 );
402
403 let paused = self.take_pause_on_bootstrap().await.unwrap_or(false);
404
405 Box::pin(
408 self.recovery(paused, RecoveryReason::Bootstrap)
409 .instrument(span),
410 )
411 .await;
412 }
413
414 Box::pin(self.run_inner(shutdown_rx)).await
415 }
416}
417
418impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
419 fn enable_per_database_isolation(&self) -> bool {
420 self.system_enable_per_database_isolation && {
421 if let Err(e) =
422 risingwave_common::license::Feature::DatabaseFailureIsolation.check_available()
423 {
424 warn!(error = %e.as_report(), "DatabaseFailureIsolation disabled by license");
425 false
426 } else {
427 true
428 }
429 }
430 }
431
432 async fn resolve_since_timestamp_snapshot_backfill(
433 &mut self,
434 new_barrier: &mut schedule::NewBarrier,
435 ) -> MetaResult<bool> {
436 let Some((
437 Command::CreateStreamingJob {
438 job_type:
439 CreateStreamingJobType::SnapshotBackfill {
440 snapshot_backfill_info,
441 since_epoch: Some(since_epoch),
442 },
443 ..
444 },
445 _,
446 )) = &mut new_barrier.command
447 else {
448 return Ok(true);
449 };
450 let since_timestamp_epoch = since_epoch.provided_since_epoch;
451
452 match self.completing_task.wait_completing_task().await {
455 Ok(Some(output)) => self
456 .checkpoint_control
457 .ack_completed(&mut self.partial_graph_manager, output),
458 Ok(None) => {}
459 Err(err) => {
460 error!(
461 err = %err.as_report(),
462 "failed to wait completing task before resolving since_timestamp"
463 );
464 return Err(err);
465 }
466 }
467
468 match self
469 .context
470 .resolve_log_store_epoch(
471 snapshot_backfill_info
472 .upstream_mv_table_id_to_backfill_epoch
473 .keys()
474 .copied(),
475 since_timestamp_epoch,
476 )
477 .await
478 {
479 Ok(upstream_log_epochs) => {
480 since_epoch.resolved = Some(upstream_log_epochs);
481 Ok(true)
482 }
483 Err(err) => {
484 error!(
485 err = %err.as_report(),
486 "failed to resolve log store epoch for since_timestamp"
487 );
488 let (_, notifier) = new_barrier
489 .command
490 .take()
491 .expect("matched command should still exist");
492 notifier.notify_start_failed(err);
493 Ok(false)
494 }
495 }
496 }
497
498 pub(super) async fn run_inner(mut self, mut shutdown_rx: Receiver<()>) {
499 let (local_notification_tx, mut local_notification_rx) =
500 tokio::sync::mpsc::unbounded_channel();
501 self.env
502 .notification_manager()
503 .insert_local_sender(local_notification_tx);
504
505 loop {
507 tokio::select! {
508 biased;
509
510 _ = &mut shutdown_rx => {
512 tracing::info!("Barrier manager is stopped");
513 break;
514 }
515
516 request = self.request_rx.recv() => {
517 if let Some(request) = request {
518 match request {
519 BarrierManagerRequest::GetBackfillProgress(result_tx) => {
520 let progress = self.checkpoint_control.gen_backfill_progress();
521 if result_tx.send(Ok(progress)).is_err() {
522 error!("failed to send get ddl progress");
523 }
524 }
525 BarrierManagerRequest::GetFragmentBackfillProgress(result_tx) => {
526 let progress =
527 self.checkpoint_control.gen_fragment_backfill_progress();
528 if result_tx.send(Ok(progress)).is_err() {
529 error!("failed to send get fragment backfill progress");
530 }
531 }
532 BarrierManagerRequest::GetCdcProgress(result_tx) => {
533 let progress = self.checkpoint_control.gen_cdc_progress();
534 if result_tx.send(Ok(progress)).is_err() {
535 error!("failed to send get ddl progress");
536 }
537 }
538 BarrierManagerRequest::AdhocRecovery(sender) => {
540 self.adhoc_recovery().await;
541 if sender.send(()).is_err() {
542 warn!("failed to notify finish of adhoc recovery");
543 }
544 }
545 BarrierManagerRequest::UpdateDatabaseBarrier( UpdateDatabaseBarrierRequest {
546 database_id,
547 barrier_interval_ms,
548 checkpoint_frequency,
549 sender,
550 }) => {
551 self.periodic_barriers
552 .update_database_barrier(
553 database_id,
554 barrier_interval_ms,
555 checkpoint_frequency,
556 );
557 if sender.send(()).is_err() {
558 warn!("failed to notify finish of update database barrier");
559 }
560 }
561 }
562 } else {
563 tracing::info!("end of request stream. meta node may be shutting down. Stop global barrier manager");
564 return;
565 }
566 }
567
568 changed_worker = self.active_streaming_nodes.changed() => {
569 #[cfg(debug_assertions)]
570 {
571 self.active_streaming_nodes.validate_change().await;
572 }
573
574 info!(?changed_worker, "worker changed");
575
576 match changed_worker {
577 ActiveStreamingWorkerChange::Add(node)
578 | ActiveStreamingWorkerChange::Update(node) => {
579 self.partial_graph_manager
580 .add_worker(node, self.context.clone())
581 .await;
582 }
583 ActiveStreamingWorkerChange::Remove(node) => {
584 self.partial_graph_manager.remove_worker(node);
585 }
586 }
587 }
588
589 notification = local_notification_rx.recv() => {
590 let notification = notification.unwrap();
591 if let LocalNotification::SystemParamsChange(p) = notification {
592 {
593 self.periodic_barriers.set_sys_barrier_interval(Duration::from_millis(p.barrier_interval_ms() as u64));
594 self.periodic_barriers
595 .set_sys_checkpoint_frequency(p.checkpoint_frequency());
596 self.system_enable_per_database_isolation = p.per_database_isolation();
597 }
598 }
599 }
600 complete_result = self
601 .completing_task
602 .next_completed_barrier(
603 &mut self.periodic_barriers,
604 &mut self.checkpoint_control,
605 &mut self.partial_graph_manager,
606 &self.context,
607 &self.env,
608 ) => {
609 match complete_result {
610 Ok(output) => {
611 self.checkpoint_control.ack_completed(&mut self.partial_graph_manager, output);
612 }
613 Err(e) => {
614 self.failure_recovery(e).await;
615 }
616 }
617 },
618 event = self.checkpoint_control.next_event() => {
619 let result: MetaResult<()> = try {
620 match event {
621 CheckpointControlEvent::EnteringInitializing(entering_initializing) => {
622 let database_id = entering_initializing.database_id();
623 let error = merge_node_rpc_errors(&format!("database {} reset", database_id), entering_initializing.action.0.iter().filter_map(|(worker_id, resp)| {
624 resp.root_err.as_ref().map(|root_err| {
625 (*worker_id, ScoredError {
626 error: Status::internal(&root_err.err_msg),
627 score: Score(root_err.score)
628 })
629 })
630 }));
631 Self::report_collect_failure(&self.env, &error);
632 self.context.notify_creating_job_failed(Some(database_id), format!("{}", error.as_report())).await;
633 let result: MetaResult<_> = try {
634 let runtime_info = self.context.reload_database_runtime_info(database_id).await.inspect_err(|err| {
635 warn!(%database_id, err = %err.as_report(), "reload runtime info failed");
636 })?;
637 let rendered_info = render_runtime_info(
638 self.env.actor_id_generator(),
639 &self.active_streaming_nodes,
640 &runtime_info.recovery_context,
641 database_id,
642 )
643 .inspect_err(|err: &MetaError| {
644 warn!(%database_id, err = %err.as_report(), "render runtime info failed");
645 })?;
646 if let Some(rendered_info) = rendered_info {
647 BarrierWorkerRuntimeInfoSnapshot::validate_database_info(
648 database_id,
649 &rendered_info.job_infos,
650 &self.active_streaming_nodes,
651 &rendered_info.stream_actors,
652 &runtime_info.state_table_committed_epochs,
653 )
654 .inspect_err(|err| {
655 warn!(%database_id, err = ?err.as_report(), "database runtime info failed validation");
656 })?;
657 Some((runtime_info, rendered_info))
658 } else {
659 None
660 }
661 };
662 match result {
663 Ok(Some((runtime_info, rendered_info))) => {
664 entering_initializing.enter(
665 runtime_info,
666 rendered_info,
667 &mut self.partial_graph_manager,
668 );
669 }
670 Ok(None) => {
671 info!(%database_id, "database removed after reloading empty runtime info");
672 entering_initializing.remove();
673 self.context
674 .refresh_table_refill_runtime_state_after_recovery()
675 .await?;
676 self.context.mark_ready(MarkReadyOptions::Database(database_id));
678 }
679 Err(e) => {
680 entering_initializing.fail_reload_runtime_info(e);
681 }
682 }
683 }
684 CheckpointControlEvent::EnteringRunning(entering_running) => {
685 let database_id = entering_running.database_id();
686 entering_running.enter();
687 self.context
688 .refresh_table_refill_runtime_state_after_recovery()
689 .await?;
690 self.context
691 .mark_ready(MarkReadyOptions::Database(database_id));
692 }
693 CheckpointControlEvent::BatchRefreshTrigger { database_id, job_id } => {
694 self.handle_batch_refresh_trigger(database_id, job_id).await?;
695 }
696 }
697 };
698 if let Err(e) = result {
699 self.failure_recovery(e).await;
700 }
701 }
702 event = self.partial_graph_manager.next_event(&self.context) => {
703 let result: MetaResult<()> = try {
704 match event {
705 PartialGraphManagerEvent::Worker(_worker_id, WorkerEvent::WorkerConnected) => {
706 }
708 PartialGraphManagerEvent::Worker(worker_id, WorkerEvent::WorkerError { err, affected_partial_graphs }) => {
709 let failed_databases = self
710 .checkpoint_control
711 .databases_failed_at_worker_err(worker_id)
712 .chain(
713 affected_partial_graphs
714 .into_iter()
715 .map(|partial_graph_id| {
716 let (database_id, _) = from_partial_graph_id(partial_graph_id);
717 database_id
718 })
719 )
720 .collect::<HashSet<_>>();
721 if !failed_databases.is_empty() {
722 if !self.enable_recovery {
723 panic!("control stream to worker {} failed but recovery not enabled: {}", worker_id, err.as_report());
724 }
725 if !self.enable_per_database_isolation() {
726 Err(err.clone())?;
727 }
728 Self::report_collect_failure(&self.env, &err);
729 for database_id in failed_databases {
730 if let Some(entering_recovery) = self.checkpoint_control.on_report_failure(database_id, &mut self.partial_graph_manager) {
731 warn!(%worker_id, %database_id, "database entering recovery on node failure");
732 self.context.abort_and_mark_blocked(Some(database_id), RecoveryReason::Failover(anyhow!("reset database: {}", database_id).into()));
733 self.context.notify_creating_job_failed(Some(database_id), format!("database {} reset due to node {} failure: {}", database_id, worker_id, err.as_report())).await;
734 let output = self.completing_task.wait_completing_task().await?;
736 entering_recovery.enter(output, &mut self.partial_graph_manager);
737 }
738 }
739 } else {
740 warn!(%worker_id, "no barrier to collect from worker, ignore err");
741 }
742 continue;
743 }
744 PartialGraphManagerEvent::PartialGraph(partial_graph_id, event) => {
745 let (database_id, _creating_job_id) = from_partial_graph_id(partial_graph_id);
746 match event {
747 PartialGraphEvent::BarrierCollected(collected_barrier) => {
748 self.checkpoint_control.barrier_collected(partial_graph_id, collected_barrier, &mut self.periodic_barriers)?;
749 }
750 PartialGraphEvent::Error(worker_id) => {
751 if !self.enable_recovery {
752 panic!("database {database_id} failure reported from {worker_id} but recovery not enabled")
753 }
754 if !self.enable_per_database_isolation() {
755 Err(MetaError::from(anyhow!("database {database_id} report failure from {worker_id}")))?;
756 }
757 if let Some(entering_recovery) = self.checkpoint_control.on_report_failure(database_id, &mut self.partial_graph_manager) {
758 warn!(%database_id, "database entering recovery");
759 self.context.abort_and_mark_blocked(Some(database_id), RecoveryReason::Failover(anyhow!("reset database: {}", database_id).into()));
760 let output = self.completing_task.wait_completing_task().await?;
762 entering_recovery.enter(output, &mut self.partial_graph_manager);
763 }
764 }
765 PartialGraphEvent::Reset(reset_resps) => {
766 self.checkpoint_control.on_partial_graph_reset(partial_graph_id, reset_resps);
767 }
768 PartialGraphEvent::Initialized => {
769 self.checkpoint_control.on_partial_graph_initialized(
770 partial_graph_id,
771 &mut self.partial_graph_manager,
772 )?;
773 }
774 }
775 }
776 };
777 };
778 if let Err(e) = result {
779 self.failure_recovery(e).await;
780 }
781 }
782 new_barrier = self.periodic_barriers.next_barrier(&*self.context) => {
783 let database_id = new_barrier.database_id;
784 let mut new_barrier = if matches!(
785 new_barrier.command,
786 Some((Command::RescheduleIntent { .. }, _))
787 ) {
788 let env = self.env.clone();
789 let worker_nodes = self
790 .active_streaming_nodes
791 .current()
792 .iter()
793 .map(|(worker_id, worker)| (*worker_id, worker.clone()))
794 .collect();
795 let database_info = self.checkpoint_control.database_info(database_id);
796 match resolve_reschedule_intent(
797 env,
798 worker_nodes,
799 database_info,
800 new_barrier,
801 ) {
802 Ok(Some(new_barrier)) => new_barrier,
803 Ok(None) => continue,
804 Err(err) => {
805 self.failure_recovery(err).await;
806 continue;
807 }
808 }
809 } else {
810 new_barrier
811 };
812 match self
813 .resolve_since_timestamp_snapshot_backfill(&mut new_barrier)
814 .await
815 {
816 Ok(true) => {}
817 Ok(false) => continue,
818 Err(err) => {
819 self.failure_recovery(err).await;
820 continue;
821 }
822 }
823 if let Err(e) = self.checkpoint_control.handle_new_barrier(
824 new_barrier,
825 &mut self.partial_graph_manager,
826 self.active_streaming_nodes.current()
827 ) {
828 if !self.enable_recovery {
829 panic!(
830 "failed to inject barrier to some databases but recovery not enabled: {:?}", (
831 database_id,
832 e.as_report()
833 )
834 );
835 }
836 let result: MetaResult<_> = try {
837 if !self.enable_per_database_isolation() {
838 let err = anyhow!("failed to inject barrier to databases: {:?}", (database_id, e.as_report()));
839 Err(MetaError::from(err))?;
840 } else if let Some(entering_recovery) = self.checkpoint_control.on_report_failure(database_id, &mut self.partial_graph_manager) {
841 warn!(%database_id, e = %e.as_report(),"database entering recovery on inject failure");
842 self.context.abort_and_mark_blocked(Some(database_id), RecoveryReason::Failover(anyhow!(e).context("inject barrier failure").into()));
843 let output = self.completing_task.wait_completing_task().await?;
845 entering_recovery.enter(output, &mut self.partial_graph_manager);
846 }
847 };
848 if let Err(e) = result {
849 self.failure_recovery(e).await;
850 }
851 }
852 }
853 }
854 }
855 }
856}
857
858impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
859 pub async fn clear_on_err(&mut self, err: &MetaError) {
861 match replace(&mut self.completing_task, CompletingTask::None) {
863 CompletingTask::None | CompletingTask::Err(_) => {}
864 CompletingTask::Completing {
865 epochs_to_ack,
866 join_handle,
867 ..
868 } => {
869 info!("waiting for completing command to finish in recovery");
870 match join_handle.await {
871 Err(e) => {
872 warn!(err = %e.as_report(), "failed to join completing task");
873 }
874 Ok(Err(e)) => {
875 warn!(
876 err = %e.as_report(),
877 "failed to complete barrier during clear"
878 );
879 }
880 Ok(Ok(hummock_version_stats)) => {
881 self.checkpoint_control.ack_completed(
882 &mut self.partial_graph_manager,
883 BarrierCompleteOutput {
884 epochs_to_ack,
885 hummock_version_stats,
886 },
887 );
888 }
889 }
890 }
891 };
892 self.partial_graph_manager.notify_all_err(err);
893 }
894}
895
896impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
897 async fn handle_batch_refresh_trigger(
900 &mut self,
901 database_id: DatabaseId,
902 job_id: JobId,
903 ) -> MetaResult<()> {
904 let last_committed_epoch = self
906 .checkpoint_control
907 .get_batch_refresh_trigger_info(database_id, job_id);
908
909 let context = self
911 .context
912 .load_batch_refresh_trigger_context(job_id, database_id, last_committed_epoch)
913 .await?;
914
915 let started = self.checkpoint_control.start_batch_refresh_run(
917 database_id,
918 job_id,
919 &context,
920 self.active_streaming_nodes.current(),
921 self.env.actor_id_generator(),
922 &mut self.partial_graph_manager,
923 )?;
924
925 if started {
927 self.checkpoint_control
928 .apply_batch_refresh_fragment_infos(database_id, job_id);
929 }
930
931 Ok(())
932 }
933}
934
935impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
936 async fn failure_recovery(&mut self, err: MetaError) {
938 self.clear_on_err(&err).await;
939
940 if self.enable_recovery {
941 let span = tracing::info_span!(
942 "failure_recovery",
943 error = %err.as_report(),
944 );
945
946 crate::telemetry::report_event(
947 risingwave_pb::telemetry::TelemetryEventStage::Recovery,
948 "failure_recovery",
949 0,
950 None,
951 None,
952 None,
953 );
954
955 let reason = RecoveryReason::Failover(err);
956
957 self.recovery(false, reason).instrument(span).await;
960 } else {
961 panic!(
962 "a streaming error occurred while recovery is disabled, aborting: {:?}",
963 err.as_report()
964 );
965 }
966 }
967
968 async fn adhoc_recovery(&mut self) {
969 let err = MetaErrorInner::AdhocRecovery.into();
970 self.clear_on_err(&err).await;
971
972 let span = tracing::info_span!(
973 "adhoc_recovery",
974 error = %err.as_report(),
975 );
976
977 crate::telemetry::report_event(
978 risingwave_pb::telemetry::TelemetryEventStage::Recovery,
979 "adhoc_recovery",
980 0,
981 None,
982 None,
983 None,
984 );
985
986 self.recovery(false, RecoveryReason::Adhoc)
989 .instrument(span)
990 .await;
991 }
992}
993
994impl<C> GlobalBarrierWorker<C> {
995 pub(super) fn report_collect_failure(env: &MetaSrvEnv, error: &MetaError) {
997 use risingwave_pb::meta::event_log;
999 let event = event_log::EventCollectBarrierFail {
1000 error: error.to_report_string(),
1001 };
1002 env.event_log_manager_ref()
1003 .add_event_logs(vec![event_log::Event::CollectBarrierFail(event)]);
1004 }
1005}
1006
1007mod retry_strategy {
1008 use std::time::Duration;
1009
1010 use risingwave_common::util::retry::exponential_backoff;
1011 use tokio_retry::strategy::jitter;
1012
1013 const RECOVERY_RETRY_BASE_INTERVAL: u64 = 20;
1015 const RECOVERY_RETRY_MAX_INTERVAL: Duration = Duration::from_secs(5);
1017
1018 pub(crate) type RetryBackoffFuture = std::pin::Pin<Box<tokio::time::Sleep>>;
1037
1038 pub(crate) fn get_retry_backoff_future(duration: Duration) -> RetryBackoffFuture {
1039 Box::pin(tokio::time::sleep(duration))
1040 }
1041
1042 pub(crate) type RetryBackoffStrategy =
1043 impl Iterator<Item = RetryBackoffFuture> + Send + 'static;
1044
1045 #[inline(always)]
1047 pub(crate) fn get_retry_strategy() -> impl Iterator<Item = Duration> + Send + 'static {
1048 exponential_backoff(
1049 Duration::from_millis(RECOVERY_RETRY_BASE_INTERVAL),
1050 RECOVERY_RETRY_BASE_INTERVAL,
1051 RECOVERY_RETRY_MAX_INTERVAL,
1052 )
1053 .map(jitter)
1054 }
1055
1056 #[define_opaque(RetryBackoffStrategy)]
1057 pub(crate) fn get_retry_backoff_strategy() -> RetryBackoffStrategy {
1058 get_retry_strategy().map(get_retry_backoff_future)
1059 }
1060}
1061
1062pub(crate) use retry_strategy::*;
1063use risingwave_common::error::tonic::extra::{Score, ScoredError};
1064use risingwave_pb::meta::event_log::{Event, EventRecovery};
1065
1066use crate::barrier::partial_graph::{
1067 PartialGraphEvent, PartialGraphManager, PartialGraphManagerEvent, WorkerEvent,
1068};
1069
1070impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
1071 pub async fn recovery(&mut self, is_paused: bool, recovery_reason: RecoveryReason) {
1079 self.partial_graph_manager.clear_worker();
1081
1082 let reason_str = match &recovery_reason {
1083 RecoveryReason::Bootstrap => "bootstrap".to_owned(),
1084 RecoveryReason::Failover(err) => {
1085 format!("failed over: {}", err.as_report())
1086 }
1087 RecoveryReason::Adhoc => "adhoc recovery".to_owned(),
1088 };
1089 self.context.abort_and_mark_blocked(None, recovery_reason);
1090
1091 self.recovery_inner(is_paused, reason_str).await;
1092 self.context.mark_ready(MarkReadyOptions::Global {
1093 blocked_databases: self.checkpoint_control.recovering_databases().collect(),
1094 });
1095 }
1096
1097 #[await_tree::instrument("recovery({recovery_reason})")]
1098 async fn recovery_inner(&mut self, is_paused: bool, recovery_reason: String) {
1099 let event_log_manager_ref = self.env.event_log_manager_ref();
1100
1101 tracing::info!("recovery start!");
1102 event_log_manager_ref.add_event_logs(vec![Event::Recovery(
1103 EventRecovery::global_recovery_start(recovery_reason.clone()),
1104 )]);
1105
1106 let retry_strategy = get_retry_strategy();
1107
1108 let recovery_timer = GLOBAL_META_METRICS
1111 .recovery_latency
1112 .with_label_values(&["global"])
1113 .start_timer();
1114
1115 let enable_per_database_isolation = self.enable_per_database_isolation();
1116
1117 let recovery_future = tokio_retry::Retry::spawn(retry_strategy, || async {
1118 self.env.stream_client_pool().invalidate_all();
1119 self.context
1122 .notify_creating_job_failed(None, recovery_reason.clone())
1123 .await;
1124
1125 let runtime_info_snapshot = self
1126 .context
1127 .reload_runtime_info()
1128 .await?;
1129 let BarrierWorkerRuntimeInfoSnapshot {
1130 active_streaming_nodes,
1131 recovery_context,
1132 mut state_table_committed_epochs,
1133 mut state_table_log_epochs,
1134 mut mv_depended_subscriptions,
1135 mut creating_jobs,
1136 hummock_version_stats,
1137 database_infos,
1138 mut cdc_table_snapshot_splits,
1139 } = runtime_info_snapshot;
1140
1141 let mut partial_graph_manager = PartialGraphManager::recover(
1142 self.env.clone(),
1143 active_streaming_nodes.current(),
1144 self.context.clone(),
1145 )
1146 .await;
1147 {
1148 let mut empty_databases = HashSet::new();
1149 let mut collected_databases = HashMap::new();
1150 let mut collecting_databases = HashMap::new();
1151 let mut failed_databases = HashMap::new();
1152 for &database_id in recovery_context.fragment_context.database_map.keys() {
1153 let mut recoverer = partial_graph_manager.start_recover();
1154 let result: MetaResult<_> = try {
1155 let Some(rendered_info) = render_runtime_info(
1156 self.env.actor_id_generator(),
1157 &active_streaming_nodes,
1158 &recovery_context,
1159 database_id,
1160 )
1161 .inspect_err(|err: &MetaError| {
1162 warn!(%database_id, err = %err.as_report(), "render runtime info failed");
1163 })? else {
1164 empty_databases.insert(database_id);
1165 continue;
1166 };
1167 BarrierWorkerRuntimeInfoSnapshot::validate_database_info(
1168 database_id,
1169 &rendered_info.job_infos,
1170 &active_streaming_nodes,
1171 &rendered_info.stream_actors,
1172 &state_table_committed_epochs,
1173 )
1174 .inspect_err(|err| {
1175 warn!(%database_id, err = %err.as_report(), "rendered runtime info failed validation");
1176 })?;
1177 let RenderedDatabaseRuntimeInfo {
1178 job_infos,
1179 stream_actors,
1180 mut source_splits,
1181 batch_refresh,
1182 } = rendered_info;
1183 recoverer.inject_database_initial_barrier(
1184 database_id,
1185 job_infos,
1186 &recovery_context.job_extra_info,
1187 &mut state_table_committed_epochs,
1188 &mut state_table_log_epochs,
1189 &recovery_context.fragment_relations,
1190 &stream_actors,
1191 &mut source_splits,
1192 &mut creating_jobs,
1193 &mut mv_depended_subscriptions,
1194 is_paused,
1195 &hummock_version_stats,
1196 &mut cdc_table_snapshot_splits,
1197 batch_refresh,
1198 )?
1199 };
1200 let collector = match result {
1201 Ok(database) => {
1202 DatabaseInitialBarrierCollector {
1203 database_id,
1204 initializing_partial_graphs: recoverer.all_initializing(),
1205 database,
1206 }
1207 }
1208 Err(e) => {
1209 warn!(%database_id, e = %e.as_report(), "failed to inject database initial barrier");
1210 assert!(failed_databases.insert(database_id, recoverer.failed()).is_none(), "non-duplicate");
1211 continue;
1212 }
1213 };
1214 if !collector.is_collected() {
1215 assert!(collecting_databases.insert(database_id, collector).is_none());
1216 } else {
1217 warn!(%database_id, "database has no node to inject initial barrier");
1218 assert!(collected_databases.insert(database_id, collector.finish()).is_none());
1219 }
1220 }
1221 if !empty_databases.is_empty() {
1222 info!(?empty_databases, "empty database in global recovery");
1223 }
1224 while !collecting_databases.is_empty() {
1225 match partial_graph_manager.next_event(&self.context).await {
1226 PartialGraphManagerEvent::Worker(_, WorkerEvent::WorkerConnected) => {
1227 }
1229 PartialGraphManagerEvent::Worker(worker_id, WorkerEvent::WorkerError { err, affected_partial_graphs }) => {
1230 let affected_databases: HashSet<_> = affected_partial_graphs.into_iter().map(|partial_graph_id| {
1231 let (database_id, _) = from_partial_graph_id(partial_graph_id);
1232 database_id
1233 }).collect();
1234 warn!(%worker_id, err = %err.as_report(), "worker node failure during recovery");
1235 for (failed_database_id, collector) in collecting_databases.extract_if(|database_id, collector| {
1236 !collector.is_valid_after_worker_err(worker_id) || affected_databases.contains(database_id)
1237 }) {
1238 warn!(%failed_database_id, %worker_id, "database failed to recovery in global recovery due to worker node err");
1239 let resetting_partial_graphs: HashSet<_> = collector.all_partial_graphs().collect();
1240 partial_graph_manager.reset_partial_graphs(resetting_partial_graphs.iter().copied());
1241 assert!(failed_databases.insert(failed_database_id, resetting_partial_graphs).is_none());
1242 }
1243 }
1244 PartialGraphManagerEvent::PartialGraph(partial_graph_id, event) => {
1245 match event {
1246 PartialGraphEvent::BarrierCollected(_) => {
1247 unreachable!("no barrier collected event on initializing")
1248 }
1249 PartialGraphEvent::Reset(_) => {
1250 let (database_id, _) =
1253 from_partial_graph_id(partial_graph_id);
1254 let resetting_partial_graphs = failed_databases
1255 .get_mut(&database_id)
1256 .expect("reset partial graph should belong to a failed database");
1257 assert!(
1258 resetting_partial_graphs.remove(&partial_graph_id),
1259 "partial graph {partial_graph_id} should be resetting"
1260 );
1261 }
1262 PartialGraphEvent::Error(worker_id) => {
1263 let (database_id, _) = from_partial_graph_id(partial_graph_id);
1264 if let Some(collector) = collecting_databases.remove(&database_id) {
1265 warn!(%database_id, %worker_id, "database reset during global recovery");
1266 let resetting_partial_graphs: HashSet<_> = collector.all_partial_graphs().collect();
1267 partial_graph_manager.reset_partial_graphs(resetting_partial_graphs.iter().copied());
1268 assert!(failed_databases.insert(database_id, resetting_partial_graphs).is_none());
1269 } else if let Some(database) = collected_databases.remove(&database_id) {
1270 warn!(%database_id, %worker_id, "database initialized but later reset during global recovery");
1271 let resetting_partial_graphs: HashSet<_> = database_partial_graphs(database_id, database.independent_checkpoint_job_controls.keys().copied()).collect();
1272 partial_graph_manager.reset_partial_graphs(resetting_partial_graphs.iter().copied());
1273 assert!(failed_databases.insert(database_id, resetting_partial_graphs).is_none());
1274 } else {
1275 assert!(failed_databases.contains_key(&database_id));
1276 }
1277 }
1278 PartialGraphEvent::Initialized => {
1279 let (database_id, _) = from_partial_graph_id(partial_graph_id);
1280 if failed_databases.contains_key(&database_id) {
1281 assert!(!collecting_databases.contains_key(&database_id));
1282 continue;
1284 }
1285 let Entry::Occupied(mut entry) = collecting_databases.entry(database_id) else {
1286 unreachable!("should exist")
1287 };
1288 let collector = entry.get_mut();
1289 collector.partial_graph_initialized(partial_graph_id);
1290 if collector.is_collected() {
1291 let collector = entry.remove();
1292 assert!(collected_databases.insert(database_id, collector.finish()).is_none());
1293 }
1294 }
1295 }
1296 }
1297 }
1298 }
1299 debug!("collected initial barrier");
1300 if !creating_jobs.is_empty() {
1301 warn!(job_ids = ?creating_jobs.iter().collect_vec(), "unused recovered creating jobs in recovery");
1302 }
1303 if !mv_depended_subscriptions.is_empty() {
1304 warn!(?mv_depended_subscriptions, "unused subscription infos in recovery");
1305 }
1306 if !state_table_committed_epochs.is_empty() {
1307 warn!(?state_table_committed_epochs, "unused state table committed epoch in recovery");
1308 }
1309 if !enable_per_database_isolation && !failed_databases.is_empty() {
1310 return Err(anyhow!(
1311 "global recovery failed due to failure of databases {:?}",
1312 failed_databases.keys().collect_vec()).into()
1313 );
1314 }
1315 let checkpoint_control = CheckpointControl::recover(
1316 collected_databases,
1317 failed_databases,
1318 hummock_version_stats,
1319 self.env.clone(),
1320 );
1321
1322 let reader = self.env.system_params_reader().await;
1323 let checkpoint_frequency = reader.checkpoint_frequency();
1324 let barrier_interval = Duration::from_millis(reader.barrier_interval_ms() as u64);
1325 let periodic_barriers = PeriodicBarriers::new(
1326 barrier_interval,
1327 checkpoint_frequency,
1328 database_infos,
1329 );
1330
1331 self.context
1332 .refresh_table_refill_runtime_state_after_recovery()
1333 .await?;
1334
1335 Ok((
1336 active_streaming_nodes,
1337 partial_graph_manager,
1338 checkpoint_control,
1339 periodic_barriers,
1340 ))
1341 }
1342 }.inspect_err(|err: &MetaError| {
1343 tracing::error!(error = %err.as_report(), "recovery failed");
1344 event_log_manager_ref.add_event_logs(vec![Event::Recovery(
1345 EventRecovery::global_recovery_failure(recovery_reason.clone(), err.to_report_string()),
1346 )]);
1347 GLOBAL_META_METRICS.recovery_failure_cnt.with_label_values(&["global"]).inc();
1348 }))
1349 .instrument(tracing::info_span!("recovery_attempt"));
1350
1351 let mut recover_txs = vec![];
1352 let mut update_barrier_requests = vec![];
1353 pin_mut!(recovery_future);
1354 let mut request_rx_closed = false;
1355 let new_state = loop {
1356 select! {
1357 biased;
1358 new_state = &mut recovery_future => {
1359 break new_state.expect("Retry until recovery success.");
1360 }
1361 request = pin!(self.request_rx.recv()), if !request_rx_closed => {
1362 let Some(request) = request else {
1363 warn!("request rx channel closed during recovery");
1364 request_rx_closed = true;
1365 continue;
1366 };
1367 match request {
1368 BarrierManagerRequest::GetBackfillProgress(tx) => {
1369 let _ = tx.send(Err(anyhow!("cluster under recovery[{}]", recovery_reason).into()));
1370 }
1371 BarrierManagerRequest::GetFragmentBackfillProgress(tx) => {
1372 let _ = tx.send(Err(anyhow!("cluster under recovery[{}]", recovery_reason).into()));
1373 }
1374 BarrierManagerRequest::GetCdcProgress(tx) => {
1375 let _ = tx.send(Err(anyhow!("cluster under recovery[{}]", recovery_reason).into()));
1376 }
1377 BarrierManagerRequest::AdhocRecovery(tx) => {
1378 recover_txs.push(tx);
1379 }
1380 BarrierManagerRequest::UpdateDatabaseBarrier(request) => {
1381 update_barrier_requests.push(request);
1382 }
1383 }
1384 }
1385 }
1386 };
1387
1388 let duration = recovery_timer.stop_and_record();
1389
1390 (
1391 self.active_streaming_nodes,
1392 self.partial_graph_manager,
1393 self.checkpoint_control,
1394 self.periodic_barriers,
1395 ) = new_state;
1396
1397 tracing::info!("recovery success");
1398
1399 for UpdateDatabaseBarrierRequest {
1400 database_id,
1401 barrier_interval_ms,
1402 checkpoint_frequency,
1403 sender,
1404 } in update_barrier_requests
1405 {
1406 self.periodic_barriers.update_database_barrier(
1407 database_id,
1408 barrier_interval_ms,
1409 checkpoint_frequency,
1410 );
1411 let _ = sender.send(());
1412 }
1413
1414 for tx in recover_txs {
1415 let _ = tx.send(());
1416 }
1417
1418 let recovering_databases = self
1419 .checkpoint_control
1420 .recovering_databases()
1421 .map(|database| database.as_raw_id())
1422 .collect_vec();
1423 let running_databases = self
1424 .checkpoint_control
1425 .running_databases()
1426 .map(|database| database.as_raw_id())
1427 .collect_vec();
1428
1429 event_log_manager_ref.add_event_logs(vec![Event::Recovery(
1430 EventRecovery::global_recovery_success(
1431 recovery_reason.clone(),
1432 duration as f32,
1433 running_databases,
1434 recovering_databases,
1435 ),
1436 )]);
1437
1438 self.env
1439 .notification_manager()
1440 .notify_frontend_without_version(Operation::Update, Info::Recovery(Recovery {}));
1441 self.env
1442 .notification_manager()
1443 .notify_compute_without_version(Operation::Update, Info::Recovery(Recovery {}));
1444 }
1445}