Skip to main content

risingwave_meta/barrier/
manager.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    /// Serving `SHOW JOBS / SELECT * FROM rw_ddl_progress`
59    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        // If not in tracker, means the first barrier not collected yet.
68        // In that case just return progress 0.
69        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    /// Check the status of barrier manager, return error if it is not `Running`.
164    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}