1use std::collections::HashMap;
16use std::sync::Arc;
17
18use anyhow::Context;
19use arc_swap::ArcSwap;
20use risingwave_common::bail;
21use risingwave_common::cast::datetime_to_timestamp_millis;
22use risingwave_hummock_sdk::HummockVersionId;
23use risingwave_meta_model::DatabaseId;
24use risingwave_pb::ddl_service::{DdlProgress, PbBackfillType};
25use risingwave_pb::id::JobId;
26use risingwave_pb::meta::PbRecoveryStatus;
27use thiserror_ext::AsReport;
28use tokio::sync::mpsc::unbounded_channel;
29use tokio::sync::{mpsc, oneshot};
30use tokio::task::JoinHandle;
31use tracing::warn;
32
33use crate::MetaResult;
34use crate::barrier::cdc_progress::CdcProgress;
35use crate::barrier::worker::GlobalBarrierWorker;
36use crate::barrier::{
37 BackfillProgress, BarrierManagerRequest, BarrierManagerStatus, FragmentBackfillProgress,
38 RecoveryReason, UpdateDatabaseBarrierRequest, schedule,
39};
40use crate::hummock::HummockManagerRef;
41use crate::manager::iceberg_compaction::IcebergCompactionManagerRef;
42use crate::manager::iceberg_pk_index_sink::IcebergPkIndexSinkManager;
43use crate::manager::sink_coordination::SinkCoordinatorManager;
44use crate::manager::{MetaSrvEnv, MetadataManager};
45use crate::serving::ServingVnodeMappingRef;
46use crate::stream::{GlobalRefreshManagerRef, ScaleControllerRef, SourceManagerRef};
47
48pub struct GlobalBarrierManager {
49 status: Arc<ArcSwap<BarrierManagerStatus>>,
50 hummock_manager: HummockManagerRef,
51 request_tx: mpsc::UnboundedSender<BarrierManagerRequest>,
52 metadata_manager: MetadataManager,
53}
54
55pub type BarrierManagerRef = Arc<GlobalBarrierManager>;
56
57impl GlobalBarrierManager {
58 pub async fn get_ddl_progress(&self) -> MetaResult<Vec<DdlProgress>> {
60 let mut backfill_progress = {
61 let (tx, rx) = oneshot::channel();
62 self.request_tx
63 .send(BarrierManagerRequest::GetBackfillProgress(tx))
64 .context("failed to send get ddl progress request")?;
65 rx.await.context("failed to receive get ddl progress")?
66 };
67 let job_info = self
70 .metadata_manager
71 .catalog_controller
72 .list_creating_jobs(true, true, None)
73 .await?;
74 Ok(job_info
75 .into_iter()
76 .map(
77 |(job_id, definition, init_at, create_type, is_serverless_backfill)| {
78 let BackfillProgress {
79 progress,
80 backfill_type,
81 } = match &mut backfill_progress {
82 Ok(progress) => progress.remove(&job_id).unwrap_or_else(|| {
83 warn!(%job_id, "background job has no ddl progress");
84 BackfillProgress {
85 progress: "0.0%".into(),
86 backfill_type: PbBackfillType::NormalBackfill,
87 }
88 }),
89 Err(e) => BackfillProgress {
90 progress: format!("Err[{}]", e.as_report()),
91 backfill_type: PbBackfillType::NormalBackfill,
92 },
93 };
94 DdlProgress {
95 id: job_id.as_raw_id() as u64,
96 statement: definition,
97 create_type: create_type.as_str().into(),
98 initialized_at_time_millis: datetime_to_timestamp_millis(init_at),
99 progress,
100 is_serverless_backfill,
101 backfill_type: backfill_type as _,
102 }
103 },
104 )
105 .collect())
106 }
107
108 pub(crate) async fn get_fragment_backfill_progress(
109 &self,
110 ) -> MetaResult<Vec<FragmentBackfillProgress>> {
111 let (tx, rx) = oneshot::channel();
112 self.request_tx
113 .send(BarrierManagerRequest::GetFragmentBackfillProgress(tx))
114 .context("failed to send get fragment backfill progress request")?;
115 rx.await
116 .context("failed to receive get fragment backfill progress")?
117 }
118
119 pub async fn get_cdc_progress(&self) -> MetaResult<HashMap<JobId, CdcProgress>> {
120 let (tx, rx) = oneshot::channel();
121 self.request_tx
122 .send(BarrierManagerRequest::GetCdcProgress(tx))
123 .context("failed to send get ddl progress request")?;
124 rx.await.context("failed to receive get ddl progress")?
125 }
126
127 pub async fn adhoc_recovery(&self) -> MetaResult<()> {
128 let (tx, rx) = oneshot::channel();
129 self.request_tx
130 .send(BarrierManagerRequest::AdhocRecovery(tx))
131 .context("failed to send adhoc recovery request")?;
132 rx.await.context("failed to wait adhoc recovery")?;
133 Ok(())
134 }
135
136 pub async fn update_database_barrier(
137 &self,
138 database_id: DatabaseId,
139 barrier_interval_ms: Option<u32>,
140 checkpoint_frequency: Option<u64>,
141 ) -> MetaResult<()> {
142 let (tx, rx) = oneshot::channel();
143 self.request_tx
144 .send(BarrierManagerRequest::UpdateDatabaseBarrier(
145 UpdateDatabaseBarrierRequest {
146 database_id,
147 barrier_interval_ms,
148 checkpoint_frequency,
149 sender: tx,
150 },
151 ))
152 .context("failed to send update database barrier request")?;
153 rx.await.context("failed to wait update database barrier")?;
154 Ok(())
155 }
156
157 pub async fn get_hummock_version_id(&self) -> HummockVersionId {
158 self.hummock_manager.get_version_id().await
159 }
160}
161
162impl GlobalBarrierManager {
163 pub fn check_status_running(&self) -> MetaResult<()> {
165 let status = self.status.load();
166 match &**status {
167 BarrierManagerStatus::Starting
168 | BarrierManagerStatus::Recovering(RecoveryReason::Bootstrap) => {
169 bail!("The cluster is bootstrapping")
170 }
171 BarrierManagerStatus::Recovering(RecoveryReason::Failover(e)) => {
172 Err(anyhow::anyhow!(e.clone()).context("The cluster is recovering"))?
173 }
174 BarrierManagerStatus::Recovering(RecoveryReason::Adhoc) => {
175 bail!("The cluster is recovering-adhoc")
176 }
177 BarrierManagerStatus::Running => Ok(()),
178 }
179 }
180
181 pub fn get_recovery_status(&self) -> PbRecoveryStatus {
182 (&**self.status.load()).into()
183 }
184}
185
186impl GlobalBarrierManager {
187 #[expect(clippy::too_many_arguments)]
188 pub async fn start(
189 scheduled_barriers: schedule::ScheduledBarriers,
190 env: MetaSrvEnv,
191 metadata_manager: MetadataManager,
192 hummock_manager: HummockManagerRef,
193 serving_vnode_mapping: ServingVnodeMappingRef,
194 source_manager: SourceManagerRef,
195 sink_manager: SinkCoordinatorManager,
196 iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
197 iceberg_compaction_manager: IcebergCompactionManagerRef,
198 scale_controller: ScaleControllerRef,
199 barrier_scheduler: schedule::BarrierScheduler,
200 refresh_manager: GlobalRefreshManagerRef,
201 ) -> (Arc<Self>, JoinHandle<()>, oneshot::Sender<()>) {
202 let (request_tx, request_rx) = unbounded_channel();
203 let hummock_manager_clone = hummock_manager.clone();
204 let metadata_manager_clone = metadata_manager.clone();
205 let barrier_worker = GlobalBarrierWorker::new(
206 scheduled_barriers,
207 env,
208 metadata_manager,
209 hummock_manager,
210 serving_vnode_mapping,
211 source_manager,
212 sink_manager,
213 iceberg_pk_index_sink_manager,
214 iceberg_compaction_manager,
215 scale_controller,
216 request_rx,
217 barrier_scheduler,
218 refresh_manager,
219 )
220 .await;
221 let manager = Self {
222 status: barrier_worker.context.status(),
223 hummock_manager: hummock_manager_clone,
224 request_tx,
225 metadata_manager: metadata_manager_clone,
226 };
227 let (join_handle, shutdown_tx) = barrier_worker.start();
228 (Arc::new(manager), join_handle, shutdown_tx)
229 }
230}