Skip to main content

risingwave_meta/hummock/manager/
commit_epoch.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    /// `table_id` -> `committed_epoch`
65    pub tables_to_commit: HashMap<TableId, u64>,
66
67    pub truncate_tables: HashSet<TableId>,
68}
69
70impl HummockManager {
71    /// Caller should ensure `epoch` > `committed_epoch` of `tables_to_commit`
72    /// if tables are not newly added via `new_table_fragment_info`
73    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        // Prevent commit new epochs if this flag is set
89        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        // Consume and aggregate table stats.
105        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        // Add new table
138        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        // fill compaction_groups
189        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        // build group_id to truncate tables
216        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            // Unable to invoke mark_next_time_travel_version_snapshot because versioning is already mutable borrowed.
245            versioning.time_travel_snapshot_interval_counter = u64::MAX;
246        }
247
248        // Apply stats changes.
249        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        // Don't trigger compactions if we enable deterministic compaction
353        if !self.env.opts.compaction_deterministic_test {
354            // commit_epoch may contains SSTs from any compaction group
355            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; // `split_sst` will split the SST into two parts and consumer 2 SST IDs
423            sst_to_cg_vec.push((commit_sst, group_table_ids));
424        }
425
426        // Generate new SST IDs for each compaction group
427        // `next_sstable_id` will update the global SST ID and reserve the new SST IDs
428        // So we need to get the new SST ID first and then split the SSTs
429        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                    // The SST contains all the tables in the group should be last key
437                    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        // order check
497        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
527/// Rewrite the commit sstables to sub-levels based on the compaction group config.
528/// The type of `compaction_group_manager_txn` is too complex to be used in the function signature. So we use `HashMap` instead.
529fn 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                // reset the accumulated size and ssts
555                accumulated_size = 0;
556                ssts = vec![];
557            }
558        }
559
560        if !ssts.is_empty() {
561            level.push(ssts);
562        }
563
564        // The uploader organizes the ssts in decreasing epoch order, so the level needs to be reversed to ensure that the latest epoch is at the top.
565        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}