1use std::collections::{BTreeMap, HashMap, HashSet};
16use std::sync::Arc;
17
18use itertools::Itertools;
19use risingwave_common::bail;
20use risingwave_common::catalog::TableId;
21use risingwave_common::config::meta::default::compaction_config;
22use risingwave_common::system_param::reader::SystemParamsRead;
23use risingwave_hummock_sdk::change_log::EpochNewChangeLog;
24use risingwave_hummock_sdk::compaction_group::group_split::split_sst_with_table_ids;
25use risingwave_hummock_sdk::sstable_info::SstableInfo;
26use risingwave_hummock_sdk::table_stats::{
27 PbTableStatsMap, add_prost_table_stats_map, purge_prost_table_stats, to_prost_table_stats_map,
28};
29use risingwave_hummock_sdk::table_watermark::TableWatermarks;
30use risingwave_hummock_sdk::vector_index::VectorIndexDelta;
31use risingwave_hummock_sdk::version::HummockVersionStateTableInfo;
32use risingwave_hummock_sdk::{
33 CompactionGroupId, HummockContextId, HummockSstableObjectId, LocalSstableInfo,
34};
35use risingwave_pb::hummock::{CompactionConfig, compact_task};
36use sea_orm::TransactionTrait;
37
38use crate::hummock::error::{Error, Result};
39use crate::hummock::manager::compaction_group_manager::CompactionGroupManager;
40use crate::hummock::manager::transaction::{
41 HummockVersionStatsTransaction, HummockVersionTransaction,
42};
43use crate::hummock::manager::versioning::Versioning;
44use crate::hummock::metrics_utils::{
45 get_or_create_local_table_stat, trigger_epoch_stat, trigger_local_table_stat, trigger_sst_stat,
46};
47use crate::hummock::model::CompactionGroup;
48use crate::hummock::sequence::{next_compaction_group_id, next_sstable_id};
49use crate::hummock::time_travel::should_mark_next_time_travel_version_snapshot;
50use crate::hummock::{HummockManager, commit_multi_var_with_provided_txn};
51
52pub struct NewTableFragmentInfo {
53 pub table_ids: HashSet<TableId>,
54}
55
56#[derive(Default)]
57pub struct CommitEpochInfo {
58 pub sstables: Vec<LocalSstableInfo>,
59 pub new_table_watermarks: HashMap<TableId, TableWatermarks>,
60 pub sst_to_context: HashMap<HummockSstableObjectId, HummockContextId>,
61 pub new_table_fragment_infos: Vec<NewTableFragmentInfo>,
62 pub change_log_delta: HashMap<TableId, EpochNewChangeLog>,
63 pub vector_index_delta: HashMap<TableId, VectorIndexDelta>,
64 pub tables_to_commit: HashMap<TableId, u64>,
66
67 pub truncate_tables: HashSet<TableId>,
68}
69
70impl HummockManager {
71 pub async fn commit_epoch(&self, commit_info: CommitEpochInfo) -> Result<()> {
74 let CommitEpochInfo {
75 mut sstables,
76 new_table_watermarks,
77 sst_to_context,
78 new_table_fragment_infos,
79 change_log_delta,
80 vector_index_delta,
81 tables_to_commit,
82 truncate_tables,
83 } = commit_info;
84 let mut versioning_guard = self
85 .versioning
86 .write_with_process_name("commit_epoch")
87 .await;
88 if versioning_guard.disable_commit_epochs {
90 return Ok(());
91 }
92
93 assert!(!tables_to_commit.is_empty());
94
95 let versioning: &mut Versioning = &mut versioning_guard;
96 self.commit_epoch_sanity_check(
97 &tables_to_commit,
98 &sstables,
99 &sst_to_context,
100 &versioning.current_version,
101 )
102 .await?;
103
104 let mut table_stats_change = PbTableStatsMap::default();
106 for s in &mut sstables {
107 add_prost_table_stats_map(
108 &mut table_stats_change,
109 &to_prost_table_stats_map(s.table_stats.clone()),
110 );
111 }
112
113 let table_change_log_object_ids_before_commit = versioning
114 .table_change_log
115 .values()
116 .flat_map(|l| l.get_object_ids())
117 .collect::<HashSet<_>>();
118
119 let mut version = HummockVersionTransaction::new(
120 &mut versioning.current_version,
121 &mut versioning.hummock_version_deltas,
122 &mut versioning.table_change_log,
123 self.env.notification_manager(),
124 Some(&self.table_committed_epoch_notifiers),
125 &self.metrics,
126 &self.env.opts,
127 &self.version_stat_tx,
128 );
129
130 let state_table_info = &version.latest_version().state_table_info;
131 let mut table_compaction_group_mapping = state_table_info.build_table_compaction_group_id();
132 let mut new_table_ids = HashMap::new();
133 let mut new_compaction_groups = Vec::new();
134 let mut compaction_group_manager_txn = None;
135 let mut compaction_group_config: Option<Arc<CompactionConfig>> = None;
136
137 for NewTableFragmentInfo { table_ids } in new_table_fragment_infos {
139 let (compaction_group_manager, compaction_group_config) =
140 if let Some(compaction_group_manager) = &mut compaction_group_manager_txn {
141 (
142 compaction_group_manager,
143 (*compaction_group_config
144 .as_ref()
145 .expect("must be set with compaction_group_manager_txn"))
146 .clone(),
147 )
148 } else {
149 let compaction_group_manager_guard = self
150 .compaction_group_manager
151 .write_with_process_name("commit_epoch")
152 .await;
153 let new_compaction_group_config =
154 compaction_group_manager_guard.default_compaction_config();
155 compaction_group_config = Some(new_compaction_group_config.clone());
156 (
157 compaction_group_manager_txn.insert(
158 CompactionGroupManager::start_owned_compaction_groups_txn(
159 compaction_group_manager_guard,
160 ),
161 ),
162 new_compaction_group_config,
163 )
164 };
165 let new_compaction_group_id = next_compaction_group_id(&self.env).await?;
166 let new_compaction_group = CompactionGroup {
167 group_id: new_compaction_group_id,
168 compaction_config: compaction_group_config.clone(),
169 };
170
171 new_compaction_groups.push(new_compaction_group.clone());
172 compaction_group_manager.insert(new_compaction_group_id, new_compaction_group);
173
174 on_handle_add_new_table(
175 state_table_info,
176 &table_ids,
177 new_compaction_group_id,
178 &mut table_compaction_group_mapping,
179 &mut new_table_ids,
180 )?;
181 }
182
183 let commit_sstables = self
184 .correct_commit_ssts(sstables, &table_compaction_group_mapping)
185 .await?;
186
187 let modified_compaction_groups = commit_sstables.keys().cloned().collect_vec();
188 let mut group_id_to_config = HashMap::new();
190 if let Some(compaction_group_manager) = compaction_group_manager_txn.as_ref() {
191 for cg_id in &modified_compaction_groups {
192 let compaction_group = compaction_group_manager
193 .get(cg_id)
194 .unwrap_or_else(|| panic!("compaction group {} should be created", cg_id))
195 .compaction_config();
196 group_id_to_config.insert(*cg_id, compaction_group);
197 }
198 } else {
199 let compaction_group_manager = self
200 .compaction_group_manager
201 .read_with_process_name("commit_epoch")
202 .await;
203 for cg_id in &modified_compaction_groups {
204 let compaction_group = compaction_group_manager
205 .try_get_compaction_group_config(*cg_id)
206 .unwrap_or_else(|| panic!("compaction group {} should be created", cg_id))
207 .compaction_config();
208 group_id_to_config.insert(*cg_id, compaction_group);
209 }
210 }
211
212 let group_id_to_sub_levels =
213 rewrite_commit_sstables_to_sub_level(commit_sstables, &group_id_to_config);
214
215 let mut group_id_to_truncate_tables: HashMap<CompactionGroupId, HashSet<TableId>> =
217 HashMap::new();
218 for table_id in &truncate_tables {
219 if let Some(compaction_group_id) = table_compaction_group_mapping.get(table_id) {
220 group_id_to_truncate_tables
221 .entry(*compaction_group_id)
222 .or_default()
223 .insert(*table_id);
224 } else {
225 bail!(
226 "table {} doesn't belong to any compaction group, skip truncating",
227 table_id
228 );
229 }
230 }
231
232 let time_travel_delta = version.pre_commit_epoch(
233 &tables_to_commit,
234 new_compaction_groups,
235 group_id_to_sub_levels,
236 &new_table_ids,
237 new_table_watermarks,
238 change_log_delta,
239 vector_index_delta,
240 group_id_to_truncate_tables,
241 );
242
243 if should_mark_next_time_travel_version_snapshot(&time_travel_delta) {
244 versioning.time_travel_snapshot_interval_counter = u64::MAX;
246 }
247
248 let mut version_stats = HummockVersionStatsTransaction::new(
250 &mut versioning.version_stats,
251 self.env.notification_manager(),
252 );
253 add_prost_table_stats_map(&mut version_stats.table_stats, &table_stats_change);
254 if purge_prost_table_stats(
255 &mut version_stats.table_stats,
256 version.latest_version(),
257 &truncate_tables,
258 ) {
259 self.metrics.version_stats.reset();
260 versioning.local_metrics.clear();
261 }
262
263 trigger_local_table_stat(
264 &self.metrics,
265 &mut versioning.local_metrics,
266 &version_stats,
267 &table_stats_change,
268 );
269 for (table_id, stats) in &table_stats_change {
270 if stats.total_key_size == 0
271 && stats.total_value_size == 0
272 && stats.total_key_count == 0
273 {
274 continue;
275 }
276 let stats_value = std::cmp::max(0, stats.total_key_size + stats.total_value_size);
277 let table_metrics = get_or_create_local_table_stat(
278 &self.metrics,
279 *table_id,
280 &mut versioning.local_metrics,
281 );
282 table_metrics.inc_write_throughput(stats_value as u64);
283 }
284 let mut time_travel_version = None;
285 if versioning.time_travel_snapshot_interval_counter
286 >= self.env.opts.hummock_time_travel_snapshot_interval
287 {
288 versioning.time_travel_snapshot_interval_counter = 0;
289 time_travel_version = Some(version.latest_version());
290 } else {
291 versioning.time_travel_snapshot_interval_counter = versioning
292 .time_travel_snapshot_interval_counter
293 .saturating_add(1);
294 }
295 let time_travel_tables_to_commit =
296 table_compaction_group_mapping
297 .iter()
298 .filter_map(|(table_id, cg_id)| {
299 tables_to_commit
300 .get(table_id)
301 .map(|committed_epoch| (table_id, cg_id, *committed_epoch))
302 });
303 let time_travel_table_ids: HashSet<_> = self
304 .metadata_manager
305 .catalog_controller
306 .list_time_travel_table_ids()
307 .await
308 .map_err(|e| Error::Internal(e.into()))?
309 .into_iter()
310 .collect();
311 let mut txn = self.env.meta_store_ref().conn.begin().await?;
312 let version_snapshot_sst_ids = self
313 .write_time_travel_metadata(
314 &txn,
315 time_travel_version,
316 time_travel_delta,
317 time_travel_table_ids,
318 &versioning.last_time_travel_snapshot_sst_ids,
319 time_travel_tables_to_commit,
320 )
321 .await?;
322 commit_multi_var_with_provided_txn!(
323 txn,
324 version,
325 version_stats,
326 compaction_group_manager_txn
327 )?;
328 if let Some(version_snapshot_sst_ids) = version_snapshot_sst_ids {
329 versioning.last_time_travel_snapshot_sst_ids = version_snapshot_sst_ids;
330 }
331
332 for compaction_group_id in &modified_compaction_groups {
333 trigger_sst_stat(
334 &self.metrics,
335 None,
336 &versioning.current_version,
337 *compaction_group_id,
338 );
339 }
340 trigger_epoch_stat(&self.metrics, &versioning.current_version);
341 let table_change_log_object_ids_after_commit = versioning
342 .table_change_log
343 .values()
344 .flat_map(|l| l.get_object_ids())
345 .collect::<HashSet<_>>();
346 drop(versioning_guard);
347 let may_delete_object_ids =
348 &table_change_log_object_ids_before_commit - &table_change_log_object_ids_after_commit;
349 self.gc_manager
350 .add_may_delete_object_ids(may_delete_object_ids.into_iter());
351
352 if !self.env.opts.compaction_deterministic_test {
354 for id in &modified_compaction_groups {
356 self.try_send_compaction_request(*id, compact_task::TaskType::Dynamic);
357 }
358 if !table_stats_change.is_empty() {
359 self.collect_table_write_throughput(table_stats_change)
360 .await;
361 }
362 }
363 if !modified_compaction_groups.is_empty() {
364 self.try_update_write_limits(&modified_compaction_groups)
365 .await;
366 }
367 #[cfg(test)]
368 {
369 self.check_state_consistency().await;
370 }
371 Ok(())
372 }
373
374 async fn collect_table_write_throughput(&self, table_stats: PbTableStatsMap) {
375 let params = self.env.system_params_reader().await;
376 let barrier_interval_ms = params.barrier_interval_ms() as u64;
377 let checkpoint_secs = std::cmp::max(
378 1,
379 params.checkpoint_frequency() * barrier_interval_ms / 1000,
380 );
381
382 let mut table_throughput_statistic_manager =
383 self.table_write_throughput_statistic_manager.write();
384 let timestamp = chrono::Utc::now().timestamp();
385
386 for (table_id, stat) in table_stats {
387 let throughput = ((stat.total_value_size + stat.total_key_size) as f64
388 / checkpoint_secs as f64) as u64;
389 table_throughput_statistic_manager
390 .add_table_throughput_with_ts(table_id, throughput, timestamp);
391 }
392 }
393
394 async fn correct_commit_ssts(
395 &self,
396 sstables: Vec<LocalSstableInfo>,
397 table_compaction_group_mapping: &HashMap<TableId, CompactionGroupId>,
398 ) -> Result<BTreeMap<CompactionGroupId, Vec<SstableInfo>>> {
399 let mut new_sst_id_number = 0;
400 let mut sst_to_cg_vec = Vec::with_capacity(sstables.len());
401 let commit_object_id_vec = sstables.iter().map(|s| s.sst_info.object_id).collect_vec();
402 for commit_sst in sstables {
403 let mut group_table_ids: BTreeMap<CompactionGroupId, Vec<TableId>> = BTreeMap::new();
404 for table_id in &commit_sst.sst_info.table_ids {
405 match table_compaction_group_mapping.get(table_id) {
406 Some(cg_id_from_meta) => {
407 group_table_ids
408 .entry(*cg_id_from_meta)
409 .or_default()
410 .push(*table_id);
411 }
412 None => {
413 tracing::warn!(
414 %table_id,
415 object_id = %commit_sst.sst_info.object_id,
416 "table doesn't belong to any compaction group",
417 );
418 }
419 }
420 }
421
422 new_sst_id_number += group_table_ids.len() * 2; sst_to_cg_vec.push((commit_sst, group_table_ids));
424 }
425
426 let mut new_sst_id = next_sstable_id(&self.env, new_sst_id_number).await?;
430 let mut commit_sstables: BTreeMap<CompactionGroupId, Vec<SstableInfo>> = BTreeMap::new();
431
432 for (mut sst, group_table_ids) in sst_to_cg_vec {
433 let len = group_table_ids.len();
434 for (index, (group_id, match_ids)) in group_table_ids.into_iter().enumerate() {
435 if sst.sst_info.table_ids == match_ids {
436 assert!(
438 index == len - 1,
439 "SST should be the last key in the group {} index {} len {}",
440 group_id,
441 index,
442 len
443 );
444 commit_sstables
445 .entry(group_id)
446 .or_default()
447 .push(sst.sst_info);
448 break;
449 }
450
451 let origin_sst_size = sst.sst_info.sst_size;
452 let new_sst_size = match_ids
453 .iter()
454 .map(|id| {
455 let stat = sst.table_stats.get(id).unwrap();
456 stat.total_compressed_size
457 })
458 .sum();
459
460 if new_sst_size == 0 {
461 tracing::warn!(
462 id = %sst.sst_info.sst_id,
463 object_id = %sst.sst_info.object_id,
464 match_ids = ?match_ids,
465 "Sstable doesn't contain any data for tables",
466 );
467 }
468
469 let old_sst_size = origin_sst_size.saturating_sub(new_sst_size);
470 if old_sst_size == 0 {
471 tracing::warn!(
472 id = %sst.sst_info.sst_id,
473 object_id = %sst.sst_info.object_id,
474 match_ids = ?match_ids,
475 origin_sst_size = origin_sst_size,
476 new_sst_size = new_sst_size,
477 "Sstable doesn't contain any data for tables",
478 );
479 }
480 let (modified_sst_info, branch_sst) = split_sst_with_table_ids(
481 &sst.sst_info,
482 &mut new_sst_id,
483 old_sst_size,
484 new_sst_size,
485 match_ids,
486 );
487 sst.sst_info = modified_sst_info;
488
489 commit_sstables
490 .entry(group_id)
491 .or_default()
492 .push(branch_sst);
493 }
494 }
495
496 for ssts in commit_sstables.values() {
498 let object_ids = ssts.iter().map(|s| s.object_id).collect_vec();
499 assert!(is_ordered_subset(&commit_object_id_vec, &object_ids));
500 }
501
502 Ok(commit_sstables)
503 }
504}
505
506fn on_handle_add_new_table(
507 state_table_info: &HummockVersionStateTableInfo,
508 table_ids: impl IntoIterator<Item = &TableId>,
509 compaction_group_id: CompactionGroupId,
510 table_compaction_group_mapping: &mut HashMap<TableId, CompactionGroupId>,
511 new_table_ids: &mut HashMap<TableId, CompactionGroupId>,
512) -> Result<()> {
513 for table_id in table_ids {
514 if let Some(info) = state_table_info.info().get(table_id) {
515 return Err(Error::CompactionGroup(format!(
516 "table {} already exist {:?}",
517 table_id, info,
518 )));
519 }
520 table_compaction_group_mapping.insert(*table_id, compaction_group_id);
521 new_table_ids.insert(*table_id, compaction_group_id);
522 }
523
524 Ok(())
525}
526
527fn rewrite_commit_sstables_to_sub_level(
530 commit_sstables: BTreeMap<CompactionGroupId, Vec<SstableInfo>>,
531 group_id_to_config: &HashMap<CompactionGroupId, Arc<CompactionConfig>>,
532) -> BTreeMap<CompactionGroupId, Vec<Vec<SstableInfo>>> {
533 let mut overlapping_sstables: BTreeMap<CompactionGroupId, Vec<Vec<SstableInfo>>> =
534 BTreeMap::new();
535 for (group_id, inserted_table_infos) in commit_sstables {
536 let config = group_id_to_config
537 .get(&group_id)
538 .expect("compaction group should exist");
539
540 let mut accumulated_size = 0;
541 let mut ssts = vec![];
542 let sub_level_size_limit = config
543 .max_overlapping_level_size
544 .unwrap_or(compaction_config::max_overlapping_level_size());
545
546 let level = overlapping_sstables.entry(group_id).or_default();
547
548 for sst in inserted_table_infos {
549 accumulated_size += sst.sst_size;
550 ssts.push(sst);
551 if accumulated_size > sub_level_size_limit {
552 level.push(ssts);
553
554 accumulated_size = 0;
556 ssts = vec![];
557 }
558 }
559
560 if !ssts.is_empty() {
561 level.push(ssts);
562 }
563
564 level.reverse();
566 }
567
568 overlapping_sstables
569}
570
571fn is_ordered_subset<T: PartialEq>(vec_1: &Vec<T>, vec_2: &Vec<T>) -> bool {
572 let mut vec_2_iter = vec_2.iter().peekable();
573 for item in vec_1 {
574 if vec_2_iter.peek() == Some(&item) {
575 vec_2_iter.next();
576 }
577 }
578
579 vec_2_iter.peek().is_none()
580}