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_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    /// `table_id` -> `committed_epoch`
64    pub tables_to_commit: HashMap<TableId, u64>,
65
66    pub truncate_tables: HashSet<TableId>,
67}
68
69impl HummockManager {
70    /// Caller should ensure `epoch` > `committed_epoch` of `tables_to_commit`
71    /// if tables are not newly added via `new_table_fragment_info`
72    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        // Prevent commit new epochs if this flag is set
88        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        // Consume and aggregate table stats.
104        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        // Add new table
137        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        // fill compaction_groups
188        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        // build group_id to truncate tables
215        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            // Unable to invoke mark_next_time_travel_version_snapshot because versioning is already mutable borrowed.
244            versioning.time_travel_snapshot_interval_counter = u64::MAX;
245        }
246
247        // Apply stats changes.
248        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        // Publish candidates while the committed groups are protected from deletion.
346        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            // Hand off publication from versioning to statistics. A merge that sees the
353            // committed version must wait for its load observations, but the per-table
354            // sampling work itself runs after releasing the version lock.
355            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                // Commit membership is authoritative: raw SST statistics can still mention
360                // dropped tables. Empty commits for live tables are real zero observations.
361                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; // `split_sst` will split the SST into two parts and consumer 2 SST IDs
416            sst_to_cg_vec.push((commit_sst, group_table_ids));
417        }
418
419        // Generate new SST IDs for each compaction group
420        // `next_sstable_id` will update the global SST ID and reserve the new SST IDs
421        // So we need to get the new SST ID first and then split the SSTs
422        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                    // The SST contains all the tables in the group should be last key
430                    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        // order check
490        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
520/// Rewrite the commit sstables to sub-levels based on the compaction group config.
521/// The type of `compaction_group_manager_txn` is too complex to be used in the function signature. So we use `HashMap` instead.
522fn 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                // reset the accumulated size and ssts
548                accumulated_size = 0;
549                ssts = vec![];
550            }
551        }
552
553        if !ssts.is_empty() {
554            level.push(ssts);
555        }
556
557        // 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.
558        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}