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