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