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