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