Skip to main content

risingwave_meta/hummock/manager/
mod.rs

1// Copyright 2022 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::{BTreeMap, HashMap};
16use std::ops::{Deref, DerefMut};
17use std::sync::Arc;
18use std::sync::atomic::AtomicBool;
19
20use anyhow::anyhow;
21use bytes::Bytes;
22use futures::FutureExt;
23use itertools::Itertools;
24use parking_lot::lock_api::RwLock;
25use risingwave_common::catalog::{TableId, TableOption};
26use risingwave_common::monitor::MonitoredRwLock;
27use risingwave_common::system_param::reader::SystemParamsRead;
28use risingwave_hummock_sdk::change_log::TableChangeLog;
29use risingwave_hummock_sdk::version::{HummockVersion, HummockVersionDelta};
30use risingwave_hummock_sdk::{
31    CompactionGroupId, HummockCompactionTaskId, HummockContextId, HummockVersionId,
32    version_archive_dir, version_checkpoint_path,
33};
34use risingwave_meta_model::{
35    compaction_status, compaction_task, hummock_pinned_version, hummock_version_delta,
36    hummock_version_stats,
37};
38use risingwave_pb::hummock::compact_task::TaskStatus;
39use risingwave_pb::hummock::{
40    HummockVersionStats, PbCompactionGroupInfo, SubscribeCompactionEventRequest,
41};
42use table_write_throughput_statistic::TableWriteThroughputStatisticManager;
43use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel};
44use tokio::sync::{Mutex, Semaphore, oneshot};
45use tonic::Streaming;
46
47use crate::MetaResult;
48use crate::hummock::CompactorManagerRef;
49use crate::hummock::compaction::CompactStatus;
50use crate::hummock::error::Result;
51use crate::hummock::manager::checkpoint::HummockVersionCheckpoint;
52use crate::hummock::manager::context::ContextInfo;
53use crate::hummock::manager::gc::{FullGcState, GcManager};
54use crate::hummock::manager::sequence::PrefetchedSequence;
55use crate::hummock::model::ext::{compaction_task_model_to_assignment, to_table_change_log};
56use crate::manager::{MetaSrvEnv, MetadataManager};
57use crate::model::{ClusterId, MetadataModelError};
58use crate::rpc::metrics::MetaMetrics;
59
60mod context;
61mod gc;
62mod tests;
63mod versioning;
64pub use context::HummockVersionSafePoint;
65use versioning::*;
66pub(crate) mod checkpoint;
67mod commit_epoch;
68mod compaction;
69pub mod sequence;
70mod table_change_log;
71pub mod table_write_throughput_statistic;
72pub mod time_travel;
73mod timer_task;
74mod transaction;
75mod utils;
76mod worker;
77
78pub use commit_epoch::{CommitEpochInfo, NewTableFragmentInfo};
79pub use compaction::compaction_event_loop::*;
80use compaction::*;
81pub use compaction::{GroupState, GroupStateValidator, ManualCompactionTriggerResult};
82pub(crate) use utils::*;
83
84struct TableCommittedEpochNotifiers {
85    txs: HashMap<TableId, Vec<UnboundedSender<u64>>>,
86}
87
88impl TableCommittedEpochNotifiers {
89    fn notify_deltas(&mut self, deltas: &[HummockVersionDelta]) {
90        self.txs.retain(|table_id, txs| {
91            let mut is_dropped = false;
92            let mut committed_epoch = None;
93            for delta in deltas {
94                if delta.removed_table_ids.contains(table_id) {
95                    is_dropped = true;
96                    break;
97                }
98                if let Some(info) = delta.state_table_info_delta.get(table_id) {
99                    committed_epoch = Some(info.committed_epoch);
100                }
101            }
102            if is_dropped {
103                false
104            } else if let Some(committed_epoch) = committed_epoch {
105                txs.retain(|tx| tx.send(committed_epoch).is_ok());
106                !txs.is_empty()
107            } else {
108                true
109            }
110        })
111    }
112}
113
114#[derive(Clone, Debug)]
115struct CompactionTaskReportResult {
116    task_id: HummockCompactionTaskId,
117
118    task_status: TaskStatus,
119    reported: bool,
120}
121
122struct CompactionTaskReportNotifiers {
123    txs: HashMap<HummockCompactionTaskId, Vec<oneshot::Sender<CompactionTaskReportResult>>>,
124}
125
126impl CompactionTaskReportNotifiers {
127    fn register(
128        &mut self,
129        task_id: HummockCompactionTaskId,
130        tx: oneshot::Sender<CompactionTaskReportResult>,
131    ) {
132        self.txs.entry(task_id).or_default().push(tx);
133    }
134
135    fn remove(&mut self, task_id: HummockCompactionTaskId) {
136        self.txs.remove(&task_id);
137    }
138
139    fn notify(&mut self, result: CompactionTaskReportResult) {
140        if let Some(txs) = self.txs.remove(&result.task_id) {
141            for tx in txs {
142                let _ = tx.send(result.clone());
143            }
144        }
145    }
146}
147// Update to states are performed as follow:
148// - Initialize ValTransaction for the meta state to update
149// - Make changes on the ValTransaction.
150// - Call `commit_multi_var` to commit the changes via meta store transaction. If transaction
151//   succeeds, the in-mem state will be updated by the way.
152pub struct HummockManager {
153    pub env: MetaSrvEnv,
154
155    metadata_manager: MetadataManager,
156    /// Lock order: `compaction`, `versioning`, `compaction_group_manager`, `context_info`
157    /// - Lock `compaction` first, then `versioning`, then `compaction_group_manager` and finally `context_info`.
158    /// - This order should be strictly followed to prevent deadlock.
159    compaction: MonitoredRwLock<Compaction>,
160    versioning: MonitoredRwLock<Versioning>,
161    /// `CompactionGroupManager` manages compaction configs for compaction groups.
162    compaction_group_manager: MonitoredRwLock<CompactionGroupManager>,
163    context_info: MonitoredRwLock<ContextInfo>,
164
165    pub metrics: Arc<MetaMetrics>,
166
167    pub compactor_manager: CompactorManagerRef,
168    pub iceberg_compactor_manager: Arc<IcebergCompactorManager>,
169    event_sender: HummockManagerEventSender,
170    object_store: ObjectStoreRef,
171    version_checkpoint_path: String,
172    version_archive_dir: String,
173    pause_version_checkpoint: AtomicBool,
174    table_write_throughput_statistic_manager:
175        parking_lot::RwLock<TableWriteThroughputStatisticManager>,
176    table_committed_epoch_notifiers: parking_lot::Mutex<TableCommittedEpochNotifiers>,
177    compaction_task_report_notifiers: parking_lot::Mutex<CompactionTaskReportNotifiers>,
178    version_stat_tx: UnboundedSender<Arc<HummockVersion>>,
179
180    // for compactor
181    // `compactor_streams_change_tx` is used to pass the mapping from `context_id` to event_stream
182    // and is maintained in memory. All event_streams are consumed through a separate event loop
183    compactor_streams_change_tx:
184        UnboundedSender<(HummockContextId, Streaming<SubscribeCompactionEventRequest>)>,
185
186    // `compaction_state` will record the types of compact tasks that can be triggered in `hummock`
187    // and suggest types with a certain priority.
188    pub compaction_state: CompactionState,
189    full_gc_state: Arc<FullGcState>,
190    now: Mutex<u64>,
191    inflight_time_travel_query: Semaphore,
192    gc_manager: GcManager,
193
194    /// In-memory cache of prefetched compaction task ids to reduce per-task DB round-trips.
195    prefetched_compaction_task_ids: PrefetchedSequence,
196    table_id_to_table_option: parking_lot::RwLock<HashMap<TableId, TableOption>>,
197}
198
199pub type HummockManagerRef = Arc<HummockManager>;
200
201use risingwave_object_store::object::{ObjectError, ObjectStoreRef, build_remote_object_store};
202use risingwave_pb::catalog::Table;
203
204use super::IcebergCompactorManager;
205use crate::controller::SqlMetaStore;
206use crate::hummock::manager::compaction_group_manager::CompactionGroupManager;
207use crate::hummock::manager::worker::HummockManagerEventSender;
208
209impl HummockManager {
210    pub async fn new(
211        env: MetaSrvEnv,
212        metadata_manager: MetadataManager,
213        metrics: Arc<MetaMetrics>,
214        compactor_manager: CompactorManagerRef,
215        compactor_streams_change_tx: UnboundedSender<(
216            HummockContextId,
217            Streaming<SubscribeCompactionEventRequest>,
218        )>,
219    ) -> Result<HummockManagerRef> {
220        let compaction_group_manager = CompactionGroupManager::new(&env).await?;
221        Self::new_impl(
222            env,
223            metadata_manager,
224            metrics,
225            compactor_manager,
226            compaction_group_manager,
227            compactor_streams_change_tx,
228        )
229        .await
230    }
231
232    #[cfg(any(test, feature = "test"))]
233    pub(super) async fn with_config(
234        env: MetaSrvEnv,
235        cluster_controller: crate::controller::cluster::ClusterControllerRef,
236        catalog_controller: crate::controller::catalog::CatalogControllerRef,
237        metrics: Arc<MetaMetrics>,
238        compactor_manager: CompactorManagerRef,
239        config: risingwave_pb::hummock::CompactionConfig,
240        compactor_streams_change_tx: UnboundedSender<(
241            HummockContextId,
242            Streaming<SubscribeCompactionEventRequest>,
243        )>,
244    ) -> HummockManagerRef {
245        let compaction_group_manager = CompactionGroupManager::new_with_config(&env, config)
246            .await
247            .unwrap();
248        let metadata_manager = MetadataManager::new(cluster_controller, catalog_controller);
249        Self::new_impl(
250            env,
251            metadata_manager,
252            metrics,
253            compactor_manager,
254            compaction_group_manager,
255            compactor_streams_change_tx,
256        )
257        .await
258        .unwrap()
259    }
260
261    async fn new_impl(
262        env: MetaSrvEnv,
263        metadata_manager: MetadataManager,
264        metrics: Arc<MetaMetrics>,
265        compactor_manager: CompactorManagerRef,
266        compaction_group_manager: CompactionGroupManager,
267        compactor_streams_change_tx: UnboundedSender<(
268            HummockContextId,
269            Streaming<SubscribeCompactionEventRequest>,
270        )>,
271    ) -> Result<HummockManagerRef> {
272        let sys_params = env.system_params_reader().await;
273        let state_store_url = sys_params.state_store();
274        let state_store_url = state_store_url.expose();
275
276        let state_store_dir: &str = sys_params.data_directory();
277        let use_new_object_prefix_strategy: bool = sys_params.use_new_object_prefix_strategy();
278        let deterministic_mode = env.opts.compaction_deterministic_test;
279        let mut object_store_config = env.opts.object_store_config.clone();
280        // For fs and hdfs object store, operations are not always atomic.
281        // We should manually enable atomicity guarantee by setting the atomic_write_dir config when building services.
282        object_store_config.set_atomic_write_dir();
283        let object_store = Arc::new(
284            build_remote_object_store(
285                state_store_url.strip_prefix("hummock+").unwrap_or("memory"),
286                metrics.object_store_metric.clone(),
287                "Version Checkpoint",
288                Arc::new(object_store_config),
289            )
290            .await,
291        );
292        // Make sure data dir is not used by another cluster.
293        // Skip this check in e2e compaction test, which needs to start a secondary cluster with
294        // same bucket
295        if !deterministic_mode {
296            write_exclusive_cluster_id(
297                state_store_dir,
298                env.cluster_id().clone(),
299                object_store.clone(),
300            )
301            .await?;
302
303            // config bucket lifecycle for new cluster.
304            if let risingwave_object_store::object::ObjectStoreImpl::S3(s3) = object_store.as_ref()
305                && !env.opts.do_not_config_object_storage_lifecycle
306            {
307                let is_bucket_expiration_configured =
308                    s3.inner().configure_bucket_lifecycle(state_store_dir).await;
309                if is_bucket_expiration_configured {
310                    return Err(ObjectError::internal("Cluster cannot start with object expiration configured for bucket because RisingWave data will be lost when object expiration kicks in.
311                    Please disable object expiration and restart the cluster.")
312                    .into());
313                }
314            }
315        }
316        let version_checkpoint_path = version_checkpoint_path(state_store_dir);
317        let version_archive_dir = version_archive_dir(state_store_dir);
318        let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
319        let (version_stat_tx, mut version_stat_rx) = tokio::sync::mpsc::unbounded_channel();
320        let inflight_time_travel_query = env.opts.max_inflight_time_travel_query;
321        let gc_manager = GcManager::new(
322            object_store.clone(),
323            state_store_dir,
324            use_new_object_prefix_strategy,
325        );
326
327        let max_table_statistic_expired_time = std::cmp::max(
328            env.opts.table_stat_throuput_window_seconds_for_split,
329            env.opts.table_stat_throuput_window_seconds_for_merge,
330        ) as i64;
331
332        let iceberg_compactor_manager = Arc::new(IcebergCompactorManager::new());
333
334        let instance = HummockManager {
335            env,
336            versioning: MonitoredRwLock::new(
337                metrics.hummock_manager_lock_time.clone(),
338                metrics.hummock_manager_real_process_time.clone(),
339                Default::default(),
340                "hummock_manager::versioning",
341            ),
342            compaction: MonitoredRwLock::new(
343                metrics.hummock_manager_lock_time.clone(),
344                metrics.hummock_manager_real_process_time.clone(),
345                Default::default(),
346                "hummock_manager::compaction",
347            ),
348            compaction_group_manager: MonitoredRwLock::new(
349                metrics.hummock_manager_lock_time.clone(),
350                metrics.hummock_manager_real_process_time.clone(),
351                compaction_group_manager,
352                "hummock_manager::compaction_group_manager",
353            ),
354            context_info: MonitoredRwLock::new(
355                metrics.hummock_manager_lock_time.clone(),
356                metrics.hummock_manager_real_process_time.clone(),
357                Default::default(),
358                "hummock_manager::context_info",
359            ),
360            metrics,
361            metadata_manager,
362            compactor_manager,
363            iceberg_compactor_manager,
364            event_sender: tx,
365            object_store,
366            version_checkpoint_path,
367            version_archive_dir,
368            pause_version_checkpoint: AtomicBool::new(false),
369            table_write_throughput_statistic_manager: parking_lot::RwLock::new(
370                TableWriteThroughputStatisticManager::new(max_table_statistic_expired_time),
371            ),
372            table_committed_epoch_notifiers: parking_lot::Mutex::new(
373                TableCommittedEpochNotifiers {
374                    txs: Default::default(),
375                },
376            ),
377            compaction_task_report_notifiers: parking_lot::Mutex::new(
378                CompactionTaskReportNotifiers {
379                    txs: Default::default(),
380                },
381            ),
382            version_stat_tx,
383            compactor_streams_change_tx,
384            compaction_state: CompactionState::new(),
385            full_gc_state: FullGcState::new().into(),
386            now: Mutex::new(0),
387            inflight_time_travel_query: Semaphore::new(inflight_time_travel_query as usize),
388            gc_manager,
389            prefetched_compaction_task_ids: PrefetchedSequence::new(),
390            table_id_to_table_option: RwLock::new(HashMap::new()),
391        };
392        let instance = Arc::new(instance);
393        let version_stat_metrics = instance.metrics.clone();
394        tokio::spawn(async move {
395            while let Some(mut version) = version_stat_rx.recv().await {
396                while let Some(Some(next_version)) = version_stat_rx.recv().now_or_never() {
397                    version = next_version;
398                }
399                transaction::trigger_version_stat(&version_stat_metrics, version.as_ref());
400            }
401        });
402        instance.init_time_travel_state().await?;
403        instance.load_meta_store_state().await?;
404        instance.start_worker(rx);
405        instance.release_invalid_contexts().await?;
406        // Release snapshots pinned by meta on restarting.
407        instance.release_meta_context().await?;
408        Ok(instance)
409    }
410
411    fn meta_store_ref(&self) -> &SqlMetaStore {
412        self.env.meta_store_ref()
413    }
414
415    /// Load state from meta store.
416    async fn load_meta_store_state(&self) -> Result<()> {
417        let now = self.load_now().await?;
418        *self.now.lock().await = now.unwrap_or(0);
419
420        let mut compaction_guard = self
421            .compaction
422            .write_with_process_name("load_meta_store_state")
423            .await;
424        let mut versioning_guard = self
425            .versioning
426            .write_with_process_name("load_meta_store_state")
427            .await;
428        let mut context_info_guard = self
429            .context_info
430            .write_with_process_name("load_meta_store_state")
431            .await;
432        self.load_meta_store_state_impl(
433            &mut compaction_guard,
434            &mut versioning_guard,
435            &mut context_info_guard,
436        )
437        .await
438    }
439
440    /// Load state from meta store.
441    async fn load_meta_store_state_impl(
442        &self,
443        compaction_guard: &mut Compaction,
444        versioning_guard: &mut Versioning,
445        context_info: &mut ContextInfo,
446    ) -> Result<()> {
447        use sea_orm::EntityTrait;
448        let meta_store = self.meta_store_ref();
449        let compaction_statuses: BTreeMap<CompactionGroupId, CompactStatus> =
450            compaction_status::Entity::find()
451                .all(&meta_store.conn)
452                .await
453                .map_err(MetadataModelError::from)?
454                .into_iter()
455                .map(|m| (m.compaction_group_id as CompactionGroupId, m.into()))
456                .collect();
457        if !compaction_statuses.is_empty() {
458            compaction_guard.compaction_statuses = compaction_statuses;
459        }
460
461        compaction_guard.compact_task_assignment = compaction_task::Entity::find()
462            .all(&meta_store.conn)
463            .await
464            .map_err(MetadataModelError::from)?
465            .into_iter()
466            .map(|m| {
467                (
468                    m.id as HummockCompactionTaskId,
469                    compaction_task_model_to_assignment(m),
470                )
471            })
472            .collect();
473
474        let hummock_version_deltas: BTreeMap<HummockVersionId, HummockVersionDelta> =
475            hummock_version_delta::Entity::find()
476                .all(&meta_store.conn)
477                .await
478                .map_err(MetadataModelError::from)?
479                .into_iter()
480                .map(|m| {
481                    (
482                        m.id,
483                        HummockVersionDelta::from_persisted_protobuf_owned(m.into()),
484                    )
485                })
486                .collect();
487
488        let checkpoint = self.try_read_checkpoint().await?;
489        let mut redo_state = if let Some(c) = checkpoint {
490            versioning_guard.checkpoint = c;
491            versioning_guard.checkpoint.version.as_ref().clone()
492        } else {
493            let default_compaction_config = self
494                .compaction_group_manager
495                .read_with_process_name("load_meta_store_state")
496                .await
497                .default_compaction_config();
498            let checkpoint_version = HummockVersion::create_init_version(default_compaction_config);
499            tracing::info!("init hummock version checkpoint");
500            versioning_guard.checkpoint = HummockVersionCheckpoint {
501                version: Arc::new(checkpoint_version.clone()),
502                stale_objects: Default::default(),
503            };
504            self.write_checkpoint(&versioning_guard.checkpoint).await?;
505            checkpoint_version
506        };
507        let mut applied_delta_count = 0;
508        let total_to_apply = hummock_version_deltas.range(redo_state.id + 1..).count();
509        tracing::info!(
510            total_delta = hummock_version_deltas.len(),
511            total_to_apply,
512            "Start redo Hummock version."
513        );
514        for version_delta in hummock_version_deltas
515            .range(redo_state.id + 1..)
516            .map(|(_, v)| v)
517        {
518            assert_eq!(
519                version_delta.prev_id, redo_state.id,
520                "delta prev_id {}, redo state id {}",
521                version_delta.prev_id, redo_state.id
522            );
523            redo_state.apply_version_delta(version_delta);
524            applied_delta_count += 1;
525            if applied_delta_count % 1000 == 0 {
526                tracing::info!("Redo progress {applied_delta_count}/{total_to_apply}.");
527            }
528        }
529        tracing::info!("Finish redo Hummock version.");
530        let pruned_stale_table_id_count = redo_state.prune_stale_table_ids_from_ssts();
531        if pruned_stale_table_id_count > 0 {
532            tracing::warn!(
533                pruned_stale_table_id_count,
534                version_id = ?redo_state.id,
535                "Pruned stale table ids from recovered Hummock SST metadata."
536            );
537        }
538        versioning_guard.version_stats = hummock_version_stats::Entity::find()
539            .one(&meta_store.conn)
540            .await
541            .map_err(MetadataModelError::from)?
542            .map(HummockVersionStats::from)
543            .unwrap_or_else(|| HummockVersionStats {
544                // version_stats.hummock_version_id is always 0 in meta store.
545                hummock_version_id: 0.into(),
546                ..Default::default()
547            });
548
549        versioning_guard.current_version = Arc::new(redo_state);
550        versioning_guard.hummock_version_deltas = hummock_version_deltas;
551        versioning_guard.table_change_log =
552            risingwave_meta_model::hummock_table_change_log::Entity::find()
553                .all(&self.env.meta_store_ref().conn)
554                .await
555                .map_err(MetadataModelError::from)?
556                .into_iter()
557                .map(|m| (m.table_id, to_table_change_log(m)))
558                .into_group_map()
559                .into_iter()
560                .map(|(table_id, unordered_change_logs)| {
561                    (
562                        table_id,
563                        TableChangeLog::new(
564                            unordered_change_logs
565                                .into_iter()
566                                .sorted_by_key(|l| l.checkpoint_epoch),
567                        ),
568                    )
569                })
570                .collect();
571
572        context_info.pinned_versions = hummock_pinned_version::Entity::find()
573            .all(&meta_store.conn)
574            .await
575            .map_err(MetadataModelError::from)?
576            .into_iter()
577            .map(|m| (m.context_id as HummockContextId, m.into()))
578            .collect();
579
580        self.initial_compaction_group_config_after_load(
581            versioning_guard,
582            self.compaction_group_manager
583                .write_with_process_name("load_meta_store_state")
584                .await
585                .deref_mut(),
586        )
587        .await?;
588
589        Ok(())
590    }
591
592    pub fn init_metadata_for_version_replay(
593        &self,
594        _table_catalogs: Vec<Table>,
595        _compaction_groups: Vec<PbCompactionGroupInfo>,
596    ) -> Result<()> {
597        unimplemented!("kv meta store is deprecated");
598    }
599
600    /// Replay a version delta to current hummock version.
601    /// Returns the `version_id`, `max_committed_epoch` of the new version and the modified
602    /// compaction groups
603    #[cfg(any(test, feature = "test"))]
604    pub async fn replay_version_delta(
605        &self,
606        mut version_delta: HummockVersionDelta,
607    ) -> Result<(HummockVersion, Vec<CompactionGroupId>)> {
608        let mut versioning_guard = self
609            .versioning
610            .write_with_process_name("replay_version_delta")
611            .await;
612        // ensure the version id is ascending after replay
613        version_delta.id = versioning_guard.current_version.next_version_id();
614        version_delta.prev_id = versioning_guard.current_version.id;
615        let mut version_new = versioning_guard.current_version.as_ref().clone();
616        version_new.apply_version_delta(&version_delta);
617
618        let compaction_group_ids = version_delta.group_deltas.keys().cloned().collect();
619        versioning_guard.current_version = Arc::new(version_new.clone());
620        Ok((version_new, compaction_group_ids))
621    }
622
623    pub async fn disable_commit_epoch(&self) -> Arc<HummockVersion> {
624        let mut versioning_guard = self
625            .versioning
626            .write_with_process_name("disable_commit_epoch")
627            .await;
628        versioning_guard.disable_commit_epochs = true;
629        versioning_guard.current_version.clone()
630    }
631
632    pub fn metadata_manager(&self) -> &MetadataManager {
633        &self.metadata_manager
634    }
635
636    pub fn object_store_media_type(&self) -> &'static str {
637        self.object_store.media_type()
638    }
639
640    pub fn update_table_id_to_table_option(
641        &self,
642        new_table_id_to_table_option: HashMap<TableId, TableOption>,
643    ) {
644        *self.table_id_to_table_option.write() = new_table_id_to_table_option;
645    }
646
647    pub fn metadata_manager_ref(&self) -> &MetadataManager {
648        &self.metadata_manager
649    }
650
651    pub async fn subscribe_table_committed_epoch(
652        &self,
653        table_id: TableId,
654    ) -> MetaResult<(u64, UnboundedReceiver<u64>)> {
655        let version = self
656            .versioning
657            .read_with_process_name("subscribe_table_committed_epoch")
658            .await;
659        if let Some(epoch) = version.current_version.table_committed_epoch(table_id) {
660            let (tx, rx) = unbounded_channel();
661            self.table_committed_epoch_notifiers
662                .lock()
663                .txs
664                .entry(table_id)
665                .or_default()
666                .push(tx);
667            Ok((epoch, rx))
668        } else {
669            Err(anyhow!("table {} does not exist", table_id).into())
670        }
671    }
672}
673
674async fn write_exclusive_cluster_id(
675    state_store_dir: &str,
676    cluster_id: ClusterId,
677    object_store: ObjectStoreRef,
678) -> Result<()> {
679    const CLUSTER_ID_DIR: &str = "cluster_id";
680    const CLUSTER_ID_NAME: &str = "0";
681    let cluster_id_dir = format!("{}/{}/", state_store_dir, CLUSTER_ID_DIR);
682    let cluster_id_full_path = format!("{}{}", cluster_id_dir, CLUSTER_ID_NAME);
683    tracing::info!("try reading cluster_id");
684    match object_store.read(&cluster_id_full_path, ..).await {
685        Ok(stored_cluster_id) => {
686            let stored_cluster_id = String::from_utf8(stored_cluster_id.to_vec()).unwrap();
687            if cluster_id.deref() == stored_cluster_id {
688                return Ok(());
689            }
690
691            Err(ObjectError::internal(format!(
692                "Data directory is already used by another cluster with id {:?}, path {}.",
693                stored_cluster_id, cluster_id_full_path,
694            ))
695            .into())
696        }
697        Err(e) => {
698            if e.is_object_not_found_error() {
699                tracing::info!("cluster_id not found, writing cluster_id");
700                object_store
701                    .upload(&cluster_id_full_path, Bytes::from(String::from(cluster_id)))
702                    .await?;
703                return Ok(());
704            }
705            Err(e.into())
706        }
707    }
708}