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.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.compaction.write().await;
420 let mut versioning_guard = self.versioning.write().await;
421 let mut context_info_guard = self.context_info.write().await;
422 self.load_meta_store_state_impl(
423 &mut compaction_guard,
424 &mut versioning_guard,
425 &mut context_info_guard,
426 )
427 .await
428 }
429
430 async fn load_meta_store_state_impl(
432 &self,
433 compaction_guard: &mut Compaction,
434 versioning_guard: &mut Versioning,
435 context_info: &mut ContextInfo,
436 ) -> Result<()> {
437 use sea_orm::EntityTrait;
438 let meta_store = self.meta_store_ref();
439 let compaction_statuses: BTreeMap<CompactionGroupId, CompactStatus> =
440 compaction_status::Entity::find()
441 .all(&meta_store.conn)
442 .await
443 .map_err(MetadataModelError::from)?
444 .into_iter()
445 .map(|m| (m.compaction_group_id as CompactionGroupId, m.into()))
446 .collect();
447 if !compaction_statuses.is_empty() {
448 compaction_guard.compaction_statuses = compaction_statuses;
449 }
450
451 compaction_guard.compact_task_assignment = compaction_task::Entity::find()
452 .all(&meta_store.conn)
453 .await
454 .map_err(MetadataModelError::from)?
455 .into_iter()
456 .map(|m| {
457 (
458 m.id as HummockCompactionTaskId,
459 compaction_task_model_to_assignment(m),
460 )
461 })
462 .collect();
463
464 let hummock_version_deltas: BTreeMap<HummockVersionId, HummockVersionDelta> =
465 hummock_version_delta::Entity::find()
466 .all(&meta_store.conn)
467 .await
468 .map_err(MetadataModelError::from)?
469 .into_iter()
470 .map(|m| {
471 (
472 m.id,
473 HummockVersionDelta::from_persisted_protobuf_owned(m.into()),
474 )
475 })
476 .collect();
477
478 let checkpoint = self.try_read_checkpoint().await?;
479 let mut redo_state = if let Some(c) = checkpoint {
480 versioning_guard.checkpoint = c;
481 versioning_guard.checkpoint.version.as_ref().clone()
482 } else {
483 let default_compaction_config = self
484 .compaction_group_manager
485 .read()
486 .await
487 .default_compaction_config();
488 let checkpoint_version = HummockVersion::create_init_version(default_compaction_config);
489 tracing::info!("init hummock version checkpoint");
490 versioning_guard.checkpoint = HummockVersionCheckpoint {
491 version: Arc::new(checkpoint_version.clone()),
492 stale_objects: Default::default(),
493 };
494 self.write_checkpoint(&versioning_guard.checkpoint).await?;
495 checkpoint_version
496 };
497 let mut applied_delta_count = 0;
498 let total_to_apply = hummock_version_deltas.range(redo_state.id + 1..).count();
499 tracing::info!(
500 total_delta = hummock_version_deltas.len(),
501 total_to_apply,
502 "Start redo Hummock version."
503 );
504 for version_delta in hummock_version_deltas
505 .range(redo_state.id + 1..)
506 .map(|(_, v)| v)
507 {
508 assert_eq!(
509 version_delta.prev_id, redo_state.id,
510 "delta prev_id {}, redo state id {}",
511 version_delta.prev_id, redo_state.id
512 );
513 redo_state.apply_version_delta(version_delta);
514 applied_delta_count += 1;
515 if applied_delta_count % 1000 == 0 {
516 tracing::info!("Redo progress {applied_delta_count}/{total_to_apply}.");
517 }
518 }
519 tracing::info!("Finish redo Hummock version.");
520 let pruned_stale_table_id_count = redo_state.prune_stale_table_ids_from_ssts();
521 if pruned_stale_table_id_count > 0 {
522 tracing::warn!(
523 pruned_stale_table_id_count,
524 version_id = ?redo_state.id,
525 "Pruned stale table ids from recovered Hummock SST metadata."
526 );
527 }
528 versioning_guard.version_stats = hummock_version_stats::Entity::find()
529 .one(&meta_store.conn)
530 .await
531 .map_err(MetadataModelError::from)?
532 .map(HummockVersionStats::from)
533 .unwrap_or_else(|| HummockVersionStats {
534 hummock_version_id: 0.into(),
536 ..Default::default()
537 });
538
539 versioning_guard.current_version = Arc::new(redo_state);
540 versioning_guard.hummock_version_deltas = hummock_version_deltas;
541 versioning_guard.table_change_log =
542 risingwave_meta_model::hummock_table_change_log::Entity::find()
543 .all(&self.env.meta_store_ref().conn)
544 .await
545 .map_err(MetadataModelError::from)?
546 .into_iter()
547 .map(|m| (m.table_id, to_table_change_log(m)))
548 .into_group_map()
549 .into_iter()
550 .map(|(table_id, unordered_change_logs)| {
551 (
552 table_id,
553 TableChangeLog::new(
554 unordered_change_logs
555 .into_iter()
556 .sorted_by_key(|l| l.checkpoint_epoch),
557 ),
558 )
559 })
560 .collect();
561
562 context_info.pinned_versions = hummock_pinned_version::Entity::find()
563 .all(&meta_store.conn)
564 .await
565 .map_err(MetadataModelError::from)?
566 .into_iter()
567 .map(|m| (m.context_id as HummockContextId, m.into()))
568 .collect();
569
570 self.initial_compaction_group_config_after_load(
571 versioning_guard,
572 self.compaction_group_manager.write().await.deref_mut(),
573 )
574 .await?;
575
576 Ok(())
577 }
578
579 pub fn init_metadata_for_version_replay(
580 &self,
581 _table_catalogs: Vec<Table>,
582 _compaction_groups: Vec<PbCompactionGroupInfo>,
583 ) -> Result<()> {
584 unimplemented!("kv meta store is deprecated");
585 }
586
587 #[cfg(any(test, feature = "test"))]
591 pub async fn replay_version_delta(
592 &self,
593 mut version_delta: HummockVersionDelta,
594 ) -> Result<(HummockVersion, Vec<CompactionGroupId>)> {
595 let mut versioning_guard = self.versioning.write().await;
596 version_delta.id = versioning_guard.current_version.next_version_id();
598 version_delta.prev_id = versioning_guard.current_version.id;
599 let mut version_new = versioning_guard.current_version.as_ref().clone();
600 version_new.apply_version_delta(&version_delta);
601
602 let compaction_group_ids = version_delta.group_deltas.keys().cloned().collect();
603 versioning_guard.current_version = Arc::new(version_new.clone());
604 Ok((version_new, compaction_group_ids))
605 }
606
607 pub async fn disable_commit_epoch(&self) -> Arc<HummockVersion> {
608 let mut versioning_guard = self.versioning.write().await;
609 versioning_guard.disable_commit_epochs = true;
610 versioning_guard.current_version.clone()
611 }
612
613 pub fn metadata_manager(&self) -> &MetadataManager {
614 &self.metadata_manager
615 }
616
617 pub fn object_store_media_type(&self) -> &'static str {
618 self.object_store.media_type()
619 }
620
621 pub fn update_table_id_to_table_option(
622 &self,
623 new_table_id_to_table_option: HashMap<TableId, TableOption>,
624 ) {
625 *self.table_id_to_table_option.write() = new_table_id_to_table_option;
626 }
627
628 pub fn metadata_manager_ref(&self) -> &MetadataManager {
629 &self.metadata_manager
630 }
631
632 pub async fn subscribe_table_committed_epoch(
633 &self,
634 table_id: TableId,
635 ) -> MetaResult<(u64, UnboundedReceiver<u64>)> {
636 let version = self.versioning.read().await;
637 if let Some(epoch) = version.current_version.table_committed_epoch(table_id) {
638 let (tx, rx) = unbounded_channel();
639 self.table_committed_epoch_notifiers
640 .lock()
641 .txs
642 .entry(table_id)
643 .or_default()
644 .push(tx);
645 Ok((epoch, rx))
646 } else {
647 Err(anyhow!("table {} not exist", table_id).into())
648 }
649 }
650}
651
652async fn write_exclusive_cluster_id(
653 state_store_dir: &str,
654 cluster_id: ClusterId,
655 object_store: ObjectStoreRef,
656) -> Result<()> {
657 const CLUSTER_ID_DIR: &str = "cluster_id";
658 const CLUSTER_ID_NAME: &str = "0";
659 let cluster_id_dir = format!("{}/{}/", state_store_dir, CLUSTER_ID_DIR);
660 let cluster_id_full_path = format!("{}{}", cluster_id_dir, CLUSTER_ID_NAME);
661 tracing::info!("try reading cluster_id");
662 match object_store.read(&cluster_id_full_path, ..).await {
663 Ok(stored_cluster_id) => {
664 let stored_cluster_id = String::from_utf8(stored_cluster_id.to_vec()).unwrap();
665 if cluster_id.deref() == stored_cluster_id {
666 return Ok(());
667 }
668
669 Err(ObjectError::internal(format!(
670 "Data directory is already used by another cluster with id {:?}, path {}.",
671 stored_cluster_id, cluster_id_full_path,
672 ))
673 .into())
674 }
675 Err(e) => {
676 if e.is_object_not_found_error() {
677 tracing::info!("cluster_id not found, writing cluster_id");
678 object_store
679 .upload(&cluster_id_full_path, Bytes::from(String::from(cluster_id)))
680 .await?;
681 return Ok(());
682 }
683 Err(e.into())
684 }
685 }
686}