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;
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}
149pub struct HummockManager {
155 pub env: MetaSrvEnv,
156
157 metadata_manager: MetadataManager,
158 compaction: MonitoredRwLock<Compaction>,
162 versioning: MonitoredRwLock<Versioning>,
163 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 compactor_streams_change_tx:
186 UnboundedSender<(HummockContextId, Streaming<SubscribeCompactionEventRequest>)>,
187
188 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 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 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 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 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 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 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 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 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 #[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 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}