1use 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}
146pub struct HummockManager {
152 pub env: MetaSrvEnv,
153
154 metadata_manager: MetadataManager,
155 compaction: MonitoredRwLock<Compaction>,
159 versioning: MonitoredRwLock<Versioning>,
160 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 compactor_streams_change_tx:
183 UnboundedSender<(HummockContextId, Streaming<SubscribeCompactionEventRequest>)>,
184
185 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 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 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 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 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 instance.may_fill_backward_table_change_logs().await?;
405 instance.start_worker(rx);
406 instance.release_invalid_contexts().await?;
407 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 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 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 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 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 #[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 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}