From b2cd68b811351af7ef2cb157181c2fbeb85d3fb1 Mon Sep 17 00:00:00 2001 From: meiyi Date: Tue, 1 Sep 2026 10:08:09 +0800 Subject: [PATCH 1/3] [feature](cloud) support single rowset grouped compaction (#65907) Cloud mode does not support seg compaction during load data. If one rowset has too many segments, the compaction may consume much memory. This pr support compacting one rowset which has more than `cloud_single_rowset_compaction_min_segments` segments by compacting every `cloud_single_rowset_compaction_segment_group_size` segments. For example, compact seg0 - seg63 first, then seg64 - seg127 and so on. If the segment group is more than 1, record the partial segments nonoverlap relation in `segment_group_sizes` in `RowsetMetaPB`(Each value is the number of consecutive output segments in one non-overlapping group.) and `NONOVERLAPPING_WITHIN_GROUP` in `SegmentsOverlapPB`. And in `VerticalBlockReader`, only read the first segment in the each group to reduce memory. --- be/src/cloud/cloud_cumulative_compaction.cpp | 205 ++++++- be/src/cloud/cloud_cumulative_compaction.h | 34 ++ be/src/cloud/cloud_schema_change_job.cpp | 8 +- be/src/cloud/cloud_snapshot_mgr.cpp | 2 + be/src/cloud/config.cpp | 3 + be/src/cloud/config.h | 3 + be/src/cloud/pb_convert.cpp | 4 + be/src/storage/compaction/compaction.cpp | 72 ++- be/src/storage/compaction/compaction.h | 32 +- .../iterator/vertical_block_reader.cpp | 70 ++- .../storage/iterator/vertical_block_reader.h | 9 + be/src/storage/merger.cpp | 67 ++- be/src/storage/merger.h | 17 +- be/src/storage/rowid_conversion.h | 17 +- be/src/storage/rowset/rowset_meta.h | 26 +- .../rowset/vertical_beta_rowset_writer.cpp | 6 +- be/test/cloud/cloud_compaction_test.cpp | 369 +++++++++++- ...loud_cumulative_compaction_policy_test.cpp | 19 + be/test/cloud/cloud_snapshot_mgr_test.cpp | 7 + ...cloud_file_cache_write_index_only_test.cpp | 5 + .../iterator/vertical_block_reader_test.cpp | 136 +++++ be/test/storage/pb_convert_test.cpp | 24 + be/test/storage/rowid_conversion_test.cpp | 384 +++++++++++- be/test/storage/rowset/rowset_meta_test.cpp | 21 + gensrc/proto/olap_file.proto | 11 + ...ud_single_rowset_grouped_compaction.groovy | 563 ++++++++++++++++++ 26 files changed, 2016 insertions(+), 98 deletions(-) create mode 100644 be/test/storage/iterator/vertical_block_reader_test.cpp create mode 100644 regression-test/suites/cloud_p0/compaction/test_cloud_single_rowset_grouped_compaction.groovy diff --git a/be/src/cloud/cloud_cumulative_compaction.cpp b/be/src/cloud/cloud_cumulative_compaction.cpp index cd61ea9c4261d2..c1b0d467bd21d4 100644 --- a/be/src/cloud/cloud_cumulative_compaction.cpp +++ b/be/src/cloud/cloud_cumulative_compaction.cpp @@ -17,6 +17,8 @@ #include "cloud/cloud_cumulative_compaction.h" +#include +#include #include #include @@ -26,12 +28,17 @@ #include "cloud/config.h" #include "common/config.h" #include "common/logging.h" +#include "common/metrics/doris_metrics.h" #include "common/status.h" #include "cpp/sync_point.h" #include "service/backend_options.h" #include "storage/compaction/compaction.h" #include "storage/compaction/cumulative_compaction_policy.h" #include "storage/compaction/cumulative_compaction_time_series_policy.h" +#include "storage/merger.h" +#include "storage/rowset/rowset_reader.h" +#include "storage/rowset/rowset_writer.h" +#include "storage/tablet/tablet_schema.h" #include "util/debug_points.h" #include "util/trace.h" #include "util/uuid_generator.h" @@ -40,6 +47,74 @@ namespace doris { #include "common/compile_check_begin.h" using namespace ErrorCode; +namespace cloud { + +bool is_single_rowset_compaction_candidate(const RowsetSharedPtr& rowset) { + const auto& rowset_meta = rowset->rowset_meta(); + const int64_t overlap_unit_count = + rowset_meta->segments_overlap() == NONOVERLAPPING_WITHIN_GROUP + ? static_cast(rowset_meta->segment_group_sizes().size()) + : rowset->num_segments(); + return !rowset_meta->has_delete_predicate() && rowset_meta->is_segments_overlapping() && + overlap_unit_count >= config::cloud_single_rowset_compaction_min_segments; +} + +bool should_use_single_rowset_grouped_compaction(const std::vector& input_rowsets, + const TabletSchema& tablet_schema, + std::string_view compaction_policy) { + return compaction_policy == CUMULATIVE_SIZE_BASED_POLICY && + tablet_schema.num_key_columns() > 0 && tablet_schema.cluster_key_uids().empty() && + config::enable_cloud_single_rowset_compaction && input_rowsets.size() == 1 && + is_single_rowset_compaction_candidate(input_rowsets.front()); +} + +std::vector build_segment_group_merge_ranges(const RowsetMeta& rowset_meta, + int64_t segment_group_size) { + DORIS_CHECK_GT(segment_group_size, 1); + DORIS_CHECK_GT(rowset_meta.num_segments(), 0); + + std::vector ranges; + if (rowset_meta.segments_overlap() == NONOVERLAPPING_WITHIN_GROUP) { + const auto& input_segment_group_sizes = rowset_meta.segment_group_sizes(); + const int64_t input_group_count = cast_set(input_segment_group_sizes.size()); + DORIS_CHECK_GT(input_group_count, 0); + ranges.reserve(cast_set((input_group_count + segment_group_size - 1) / + segment_group_size)); + + int64_t segment_end = 0; + for (int64_t group_start = 0; group_start < input_group_count; + group_start += segment_group_size) { + const int64_t group_end = std::min(group_start + segment_group_size, input_group_count); + const int64_t segment_start = segment_end; + for (int64_t group_index = group_start; group_index < group_end; ++group_index) { + const int32_t input_group_size = + input_segment_group_sizes.Get(cast_set(group_index)); + DORIS_CHECK_GT(input_group_size, 0); + segment_end += input_group_size; + } + + ranges.push_back({.segment_start = segment_start, + .segment_end = segment_end, + .merge_way_num = group_end - group_start}); + } + DORIS_CHECK_EQ(segment_end, rowset_meta.num_segments()); + } else { + ranges.reserve(cast_set((rowset_meta.num_segments() + segment_group_size - 1) / + segment_group_size)); + for (int64_t segment_start = 0; segment_start < rowset_meta.num_segments(); + segment_start += segment_group_size) { + const int64_t segment_end = + std::min(segment_start + segment_group_size, rowset_meta.num_segments()); + ranges.push_back({.segment_start = segment_start, + .segment_end = segment_end, + .merge_way_num = segment_end - segment_start}); + } + } + return ranges; +} + +} // namespace cloud + bvar::Adder cumu_output_size("cumu_compaction", "output_size"); bvar::LatencyRecorder g_cu_compaction_hold_delete_bitmap_lock_time_ms( "cu_compaction_hold_delete_bitmap_lock_time_ms"); @@ -250,6 +325,18 @@ Status CloudCumulativeCompaction::execute_compact() { return st; } +bool CloudCumulativeCompaction::should_calculate_new_cumulative_point( + int64_t input_cumulative_point) const { + if (!_single_rowset_compaction_segment_group_size.has_value()) { + return true; + } + + DORIS_CHECK_EQ(_input_rowsets.size(), 1); + DORIS_CHECK(_output_rowset != nullptr); + return _input_rowsets.front()->start_version() == input_cumulative_point && + _output_rowset->rowset_meta()->segments_overlap() == NONOVERLAPPING; +} + Status CloudCumulativeCompaction::modify_rowsets() { // calculate new cumulative point int64_t input_cumulative_point; @@ -267,17 +354,19 @@ Status CloudCumulativeCompaction::modify_rowsets() { } auto compaction_policy = cloud_tablet()->tablet_meta()->compaction_policy(); int64_t new_cumulative_point = input_cumulative_point; - if (!_enable_parallel_cumu_compaction && input_tablet_state == TABLET_NOTREADY && - _output_rowset->start_version() > input_cumulative_point) { - // Historical rowsets are absent from a schema-change target until conversion finishes. - DORIS_CHECK_LE(input_cumulative_point, input_alter_version); - DORIS_CHECK_GT(_output_rowset->start_version(), input_alter_version); - } else if (!_enable_parallel_cumu_compaction || - _output_rowset->start_version() == input_cumulative_point) { - new_cumulative_point = - _engine.cumu_compaction_policy(compaction_policy) - ->new_cumulative_point(cloud_tablet(), _output_rowset, _last_delete_version, - input_cumulative_point); + if (should_calculate_new_cumulative_point(input_cumulative_point)) { + if (!_enable_parallel_cumu_compaction && input_tablet_state == TABLET_NOTREADY && + _output_rowset->start_version() > input_cumulative_point) { + // Historical rowsets are absent from a schema-change target until conversion finishes. + DORIS_CHECK_LE(input_cumulative_point, input_alter_version); + DORIS_CHECK_GT(_output_rowset->start_version(), input_alter_version); + } else if (!_enable_parallel_cumu_compaction || + _output_rowset->start_version() == input_cumulative_point) { + new_cumulative_point = + _engine.cumu_compaction_policy(compaction_policy) + ->new_cumulative_point(cloud_tablet(), _output_rowset, + _last_delete_version, input_cumulative_point); + } } // commit compaction job cloud::TabletJobInfoPB job; @@ -608,6 +697,7 @@ Status CloudCumulativeCompaction::advance_cumulative_point_before_pick( Status CloudCumulativeCompaction::pick_rowsets_to_compact() { _input_rowsets.clear(); + _single_rowset_compaction_segment_group_size.reset(); int64_t min_conflict_version = _min_conflict_version; int64_t max_conflict_version = _max_conflict_version; @@ -646,6 +736,20 @@ Status CloudCumulativeCompaction::pick_rowsets_to_compact() { config::cumulative_compaction_min_deltas, &_input_rowsets, &_last_delete_version, &compaction_score); + const int64_t segment_group_size = + config::cloud_single_rowset_compaction_segment_group_size; + if (config::enable_cloud_single_rowset_compaction && segment_group_size > 1) { + for (const auto& rowset : _input_rowsets) { + if (cloud::should_use_single_rowset_grouped_compaction( + {rowset}, *cloud_tablet()->tablet_schema(), compaction_policy)) { + auto grouped_input_rowset = rowset; + _input_rowsets = {std::move(grouped_input_rowset)}; + _single_rowset_compaction_segment_group_size = segment_group_size; + return Status::OK(); + } + } + } + if (_input_rowsets.empty()) { return Status::Error( "no suitable versions: input rowsets empty"); @@ -709,6 +813,85 @@ Status CloudCumulativeCompaction::pick_rowsets_to_compact() { return Status::OK(); } +Status CloudCumulativeCompaction::prepare_merge_input_rowsets(MergeInputRowsetsResult* result) { + if (!_single_rowset_compaction_segment_group_size.has_value()) { + return Status::OK(); + } + + const int64_t segment_group_size = *_single_rowset_compaction_segment_group_size; + DORIS_CHECK_GT(segment_group_size, 1); + result->is_segment_grouped = true; + result->segment_group_size = segment_group_size; + return Status::OK(); +} + +Status CloudCumulativeCompaction::do_merge_input_rowsets( + const std::vector& input_rs_readers, + MergeInputRowsetsResult* result) { + if (!result->is_segment_grouped) { + return Compaction::do_merge_input_rowsets(input_rs_readers, result); + } + + const int64_t segment_group_size = result->segment_group_size; + const auto& input_rowset = _input_rowsets.front(); + const auto segment_ranges = cloud::build_segment_group_merge_ranges( + *input_rowset->rowset_meta(), segment_group_size); + for (size_t range_index = 0; range_index < segment_ranges.size(); ++range_index) { + const auto& range = segment_ranges[range_index]; + const int32_t output_segment_start = _output_rs_writer->get_allocated_segment_id(); + + RowsetReaderSharedPtr rs_reader; + RETURN_IF_ERROR(input_rowset->create_reader(&rs_reader)); + std::vector group_readers; + group_readers.push_back(std::move(rs_reader)); + + Merger::Statistics group_stats; + group_stats.rowid_conversion = _stats.rowid_conversion; + RETURN_IF_ERROR(execute_merge(group_readers, range.merge_way_num, &group_stats, + std::make_pair(range.segment_start, range.segment_end), + {.total_ranges = cast_set(segment_ranges.size()), + .range_index = cast_set(range_index)})); + + _stats.output_rows += group_stats.output_rows; + _stats.merged_rows += group_stats.merged_rows; + _stats.filtered_rows += group_stats.filtered_rows; + _stats.bytes_read_from_local += group_stats.bytes_read_from_local; + _stats.bytes_read_from_remote += group_stats.bytes_read_from_remote; + _stats.cached_bytes_total += group_stats.cached_bytes_total; + _stats.cloud_local_read_time += group_stats.cloud_local_read_time; + _stats.cloud_remote_read_time += group_stats.cloud_remote_read_time; + + const int32_t output_segment_end = _output_rs_writer->get_allocated_segment_id(); + const int32_t output_group_size = output_segment_end - output_segment_start; + if (output_group_size > 0) { + result->output_segment_group_sizes.push_back(output_group_size); + } + } + return Status::OK(); +} + +void CloudCumulativeCompaction::update_output_rowset_after_build( + const MergeInputRowsetsResult& result) { + if (!result.is_segment_grouped) { + return; + } + if (result.output_segment_group_sizes.size() > 1) { + _output_rowset->rowset_meta()->set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + _output_rowset->rowset_meta()->set_segment_group_sizes(result.output_segment_group_sizes); + } + + const auto& input_rowset = _input_rowsets.front(); + LOG_INFO("finish single rowset grouped compaction, tablet_id={}, version=[{}-{}]", + _tablet->tablet_id(), input_rowset->start_version(), input_rowset->end_version()) + .tag("job_id", _uuid) + .tag("input_segments", input_rowset->num_segments()) + .tag("segment_group_size", result.segment_group_size) + .tag("output_segments", _output_rowset->num_segments()) + .tag("output_groups", result.output_segment_group_sizes.size()) + .tag("output_segment_group_sizes", + fmt::format("[{}]", fmt::join(result.output_segment_group_sizes, ", "))); +} + void CloudCumulativeCompaction::update_cumulative_point(int64_t input_cumulative_point, int64_t output_cumulative_point) { DORIS_CHECK_LT(input_cumulative_point, output_cumulative_point); diff --git a/be/src/cloud/cloud_cumulative_compaction.h b/be/src/cloud/cloud_cumulative_compaction.h index 28f88c6432f55a..004bd7a5db0b65 100644 --- a/be/src/cloud/cloud_cumulative_compaction.h +++ b/be/src/cloud/cloud_cumulative_compaction.h @@ -20,15 +20,39 @@ #include #include #include +#include +#include #include "cloud/cloud_storage_engine.h" #include "cloud/cloud_tablet.h" #include "storage/compaction/compaction.h" #include "storage/compaction_task_tracker.h" +#include "storage/tablet/tablet_fwd.h" namespace doris { #include "common/compile_check_begin.h" +class RowsetMeta; + +namespace cloud { + +struct SegmentGroupMergeRange { + int64_t segment_start; + int64_t segment_end; + int64_t merge_way_num; +}; + +bool is_single_rowset_compaction_candidate(const RowsetSharedPtr& rowset); + +bool should_use_single_rowset_grouped_compaction(const std::vector& input_rowsets, + const TabletSchema& tablet_schema, + std::string_view compaction_policy); + +std::vector build_segment_group_merge_ranges(const RowsetMeta& rowset_meta, + int64_t segment_group_size); + +} // namespace cloud + class CloudCumulativeCompaction : public CloudCompactionMixin { public: CloudCumulativeCompaction(CloudStorageEngine& engine, CloudTabletSPtr tablet); @@ -54,6 +78,15 @@ class CloudCumulativeCompaction : public CloudCompactionMixin { Status pick_rowsets_to_compact(); + Status prepare_merge_input_rowsets(MergeInputRowsetsResult* result) override; + + Status do_merge_input_rowsets(const std::vector& input_rs_readers, + MergeInputRowsetsResult* result) override; + + void update_output_rowset_after_build(const MergeInputRowsetsResult& result) override; + + bool should_calculate_new_cumulative_point(int64_t input_cumulative_point) const; + std::string_view compaction_name() const override { return "CloudCumulativeCompaction"; } protected: @@ -76,6 +109,7 @@ class CloudCumulativeCompaction : public CloudCompactionMixin { int64_t _cumulative_compaction_cnt = 0; int64_t _picked_cumulative_point = 0; Version _last_delete_version {-1, -1}; + std::optional _single_rowset_compaction_segment_group_size; }; #include "common/compile_check_end.h" diff --git a/be/src/cloud/cloud_schema_change_job.cpp b/be/src/cloud/cloud_schema_change_job.cpp index 32b4386e0595e8..b63f6ee53bb29a 100644 --- a/be/src/cloud/cloud_schema_change_job.cpp +++ b/be/src/cloud/cloud_schema_change_job.cpp @@ -381,7 +381,13 @@ Status CloudSchemaChangeJob::_convert_historical_rowsets(const SchemaChangeParam context.txn_expiration = _expiration; context.version = rs_reader->version(); context.rowset_state = VISIBLE; - context.segments_overlap = rs_reader->rowset()->rowset_meta()->segments_overlap(); + const auto input_segments_overlap = rs_reader->rowset()->rowset_meta()->segments_overlap(); + // Cloud schema change rewrites remote rowsets, so the input group layout is no longer + // applicable. Fall back to the conservative overlap state without assuming that the + // rewritten segments are globally ordered. + context.segments_overlap = input_segments_overlap == NONOVERLAPPING_WITHIN_GROUP + ? OVERLAPPING + : input_segments_overlap; context.tablet_schema = _new_tablet->tablet_schema(); context.newest_write_timestamp = rs_reader->newest_write_timestamp(); context.storage_resource = _cloud_storage_engine.get_storage_resource(sc_params.vault_id); diff --git a/be/src/cloud/cloud_snapshot_mgr.cpp b/be/src/cloud/cloud_snapshot_mgr.cpp index bd4e3fee76757f..1f096c7b87a483 100644 --- a/be/src/cloud/cloud_snapshot_mgr.cpp +++ b/be/src/cloud/cloud_snapshot_mgr.cpp @@ -271,6 +271,8 @@ Status CloudSnapshotMgr::_create_rowset_meta( new_rowset_meta_pb->set_creation_time(time(nullptr)); new_rowset_meta_pb->set_num_segments(source_meta_pb.num_segments()); new_rowset_meta_pb->set_rowset_state(source_meta_pb.rowset_state()); + new_rowset_meta_pb->mutable_segment_group_sizes()->CopyFrom( + source_meta_pb.segment_group_sizes()); new_rowset_meta_pb->clear_segments_key_bounds(); for (const auto& key_bound : source_meta_pb.segments_key_bounds()) { diff --git a/be/src/cloud/config.cpp b/be/src/cloud/config.cpp index cd180c2e5b3b87..1de5a5cc60b257 100644 --- a/be/src/cloud/config.cpp +++ b/be/src/cloud/config.cpp @@ -57,6 +57,9 @@ DEFINE_mInt32(max_base_compaction_task_num_per_disk, "2"); DEFINE_mBool(prioritize_query_perf_in_compaction, "false"); DEFINE_mInt32(compaction_max_rowset_count, "10000"); DEFINE_mInt64(compaction_txn_max_size_bytes, "7340032"); // 7MB +DEFINE_mBool(enable_cloud_single_rowset_compaction, "false"); +DEFINE_mInt32(cloud_single_rowset_compaction_min_segments, "512"); +DEFINE_mInt32(cloud_single_rowset_compaction_segment_group_size, "64"); DEFINE_mInt32(refresh_s3_info_interval_s, "60"); DEFINE_mInt32(vacuum_stale_rowsets_interval_s, "300"); diff --git a/be/src/cloud/config.h b/be/src/cloud/config.h index 5ef03b9df446d5..270ca667688b38 100644 --- a/be/src/cloud/config.h +++ b/be/src/cloud/config.h @@ -97,6 +97,9 @@ DECLARE_mInt32(max_base_compaction_task_num_per_disk); DECLARE_mBool(prioritize_query_perf_in_compaction); DECLARE_mInt32(compaction_max_rowset_count); DECLARE_mInt64(compaction_txn_max_size_bytes); +DECLARE_mBool(enable_cloud_single_rowset_compaction); +DECLARE_mInt32(cloud_single_rowset_compaction_min_segments); +DECLARE_mInt32(cloud_single_rowset_compaction_segment_group_size); // CloudStorageEngine config DECLARE_mInt32(refresh_s3_info_interval_s); diff --git a/be/src/cloud/pb_convert.cpp b/be/src/cloud/pb_convert.cpp index a90f93feeda77f..b8d386b201ba39 100644 --- a/be/src/cloud/pb_convert.cpp +++ b/be/src/cloud/pb_convert.cpp @@ -80,6 +80,7 @@ void doris_rowset_meta_to_cloud(RowsetMetaCloudPB* out, const RowsetMetaPB& in) } out->set_txn_expiration(in.txn_expiration()); out->set_segments_overlap_pb(in.segments_overlap_pb()); + out->mutable_segment_group_sizes()->CopyFrom(in.segment_group_sizes()); if (in.has_segments_key_bounds_truncated()) { out->set_segments_key_bounds_truncated(in.segments_key_bounds_truncated()); } @@ -163,6 +164,7 @@ void doris_rowset_meta_to_cloud(RowsetMetaCloudPB* out, RowsetMetaPB&& in) { } out->set_txn_expiration(in.txn_expiration()); out->set_segments_overlap_pb(in.segments_overlap_pb()); + out->mutable_segment_group_sizes()->Swap(in.mutable_segment_group_sizes()); if (in.has_segments_key_bounds_truncated()) { out->set_segments_key_bounds_truncated(in.segments_key_bounds_truncated()); } @@ -258,6 +260,7 @@ void cloud_rowset_meta_to_doris(RowsetMetaPB* out, const RowsetMetaCloudPB& in) } out->set_txn_expiration(in.txn_expiration()); out->set_segments_overlap_pb(in.segments_overlap_pb()); + out->mutable_segment_group_sizes()->CopyFrom(in.segment_group_sizes()); if (in.has_segments_key_bounds_truncated()) { out->set_segments_key_bounds_truncated(in.segments_key_bounds_truncated()); } @@ -341,6 +344,7 @@ void cloud_rowset_meta_to_doris(RowsetMetaPB* out, RowsetMetaCloudPB&& in) { } out->set_txn_expiration(in.txn_expiration()); out->set_segments_overlap_pb(in.segments_overlap_pb()); + out->mutable_segment_group_sizes()->Swap(in.mutable_segment_group_sizes()); if (in.has_segments_key_bounds_truncated()) { out->set_segments_key_bounds_truncated(in.segments_key_bounds_truncated()); } diff --git a/be/src/storage/compaction/compaction.cpp b/be/src/storage/compaction/compaction.cpp index b2aa69ffe6dee6..aa61df40dcc998 100644 --- a/be/src/storage/compaction/compaction.cpp +++ b/be/src/storage/compaction/compaction.cpp @@ -256,6 +256,9 @@ int64_t Compaction::merge_way_num() { } Status Compaction::merge_input_rowsets() { + MergeInputRowsetsResult result; + RETURN_IF_ERROR(prepare_merge_input_rowsets(&result)); + std::vector input_rs_readers; input_rs_readers.reserve(_input_rowsets.size()); for (auto& rowset : _input_rowsets) { @@ -280,37 +283,10 @@ Status Compaction::merge_input_rowsets() { _stats.rowid_conversion = _rowid_conversion.get(); } - int64_t way_num = merge_way_num(); - - Status res; { SCOPED_TIMER(_merge_rowsets_latency_timer); // 1. Merge segment files and write bkd inverted index - if (_is_vertical) { - if (!_tablet->tablet_schema()->cluster_key_uids().empty()) { - RETURN_IF_ERROR(update_delete_bitmap()); - } - auto progress_cb = [compaction_id = this->_compaction_id](int64_t total, - int64_t completed) { - CompactionTaskTracker::instance()->update_progress(compaction_id, total, completed); - }; - res = Merger::vertical_merge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema, - input_rs_readers, _output_rs_writer.get(), - cast_set(get_avg_segment_rows()), - way_num, &_stats, progress_cb); - } else { - if (!_tablet->tablet_schema()->cluster_key_uids().empty()) { - return Status::InternalError( - "mow table with cluster keys does not support non vertical compaction"); - } - res = Merger::vmerge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema, - input_rs_readers, _output_rs_writer.get(), &_stats); - } - - _tablet->last_compaction_status = res; - if (!res.ok()) { - return res; - } + RETURN_IF_ERROR(do_merge_input_rowsets(input_rs_readers, &result)); // 2. Merge the remaining inverted index files of the string type RETURN_IF_ERROR(do_inverted_index_compaction()); } @@ -335,6 +311,7 @@ Status Compaction::merge_input_rowsets() { //RETURN_IF_ERROR(_engine.meta_mgr().commit_rowset(*_output_rowset->rowset_meta().get())); set_delete_predicate_for_output_rowset(); + update_output_rowset_after_build(result); _local_read_bytes_total = _stats.bytes_read_from_local; _remote_read_bytes_total = _stats.bytes_read_from_remote; @@ -351,6 +328,45 @@ Status Compaction::merge_input_rowsets() { return check_correctness(); } +Status Compaction::do_merge_input_rowsets( + const std::vector& input_rs_readers, + MergeInputRowsetsResult* /*result*/) { + return execute_merge(input_rs_readers, merge_way_num(), &_stats); +} + +Status Compaction::execute_merge(const std::vector& input_rs_readers, + int64_t merge_way_num, Merger::Statistics* stats, + std::optional> segment_range, + VerticalMergeProgressContext progress) { + Status status; + if (_is_vertical) { + if (!_tablet->tablet_schema()->cluster_key_uids().empty() && !segment_range.has_value()) { + RETURN_IF_ERROR(update_delete_bitmap()); + } + auto progress_cb = [compaction_id = this->_compaction_id, progress](int64_t total, + int64_t completed) { + CompactionTaskTracker::instance()->update_progress( + compaction_id, total * progress.total_ranges, + total * progress.range_index + completed); + }; + status = Merger::vertical_merge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema, + input_rs_readers, _output_rs_writer.get(), + cast_set(get_avg_segment_rows()), + merge_way_num, stats, progress_cb, segment_range); + } else { + if (!_tablet->tablet_schema()->cluster_key_uids().empty()) { + return Status::InternalError( + "mow table with cluster keys does not support non vertical compaction"); + } + status = Merger::vmerge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema, + input_rs_readers, _output_rs_writer.get(), stats, + segment_range); + } + + _tablet->last_compaction_status = status; + return status; +} + void Compaction::set_delete_predicate_for_output_rowset() { // Now we support delete in cumu compaction, to make all data in rowsets whose version // is below output_version to be delete in the future base compaction, we should carry diff --git a/be/src/storage/compaction/compaction.h b/be/src/storage/compaction/compaction.h index 4bc0b4e3ad20c1..c24cde1aa61676 100644 --- a/be/src/storage/compaction/compaction.h +++ b/be/src/storage/compaction/compaction.h @@ -30,6 +30,7 @@ #endif #include #include +#include #include #include "cloud/cloud_tablet.h" @@ -102,8 +103,35 @@ class Compaction { void set_delete_predicate_for_output_rowset(); protected: + struct MergeInputRowsetsResult { + bool is_segment_grouped = false; + int64_t segment_group_size = 0; + std::vector output_segment_group_sizes; + }; + Status merge_input_rowsets(); + virtual Status prepare_merge_input_rowsets(MergeInputRowsetsResult* /*result*/) { + return Status::OK(); + } + + virtual Status do_merge_input_rowsets( + const std::vector& input_rs_readers, + MergeInputRowsetsResult* result); + + virtual void update_output_rowset_after_build(const MergeInputRowsetsResult& /*result*/) {} + + // Maps one segment-range merge's column-group progress into the whole compaction task. + struct VerticalMergeProgressContext { + int64_t total_ranges; + int64_t range_index; + }; + + Status execute_merge(const std::vector& input_rs_readers, + int64_t merge_way_num, Merger::Statistics* stats, + std::optional> segment_range = std::nullopt, + VerticalMergeProgressContext progress = {1, 0}); + // merge inverted index files Status do_inverted_index_compaction(); @@ -250,6 +278,8 @@ class CloudCompactionMixin : public Compaction { virtual Status garbage_collection(); + Status construct_output_rowset_writer(RowsetWriterContext& ctx) override; + // Helper function to apply truncation and log the result // Returns the number of rowsets that were truncated size_t apply_txn_size_truncation_and_log(const std::string& compaction_name); @@ -266,8 +296,6 @@ class CloudCompactionMixin : public Compaction { virtual Status rebuild_tablet_schema() { return Status::OK(); } private: - Status construct_output_rowset_writer(RowsetWriterContext& ctx) override; - Status set_storage_resource_from_input_rowsets(RowsetWriterContext& ctx); Status execute_compact_impl(int64_t permits); diff --git a/be/src/storage/iterator/vertical_block_reader.cpp b/be/src/storage/iterator/vertical_block_reader.cpp index 1fe1274dac55f1..697d838b14466b 100644 --- a/be/src/storage/iterator/vertical_block_reader.cpp +++ b/be/src/storage/iterator/vertical_block_reader.cpp @@ -56,6 +56,40 @@ VerticalBlockReader::~VerticalBlockReader() { } } +void VerticalBlockReader::_append_grouped_iterator_init_flags( + const RowsetMeta& rowset_meta, std::pair segment_offsets, + size_t added_iterators, std::vector* iterator_init_flag) { + auto [segment_start, segment_end] = segment_offsets; + if (segment_start == segment_end) { + segment_start = 0; + segment_end = rowset_meta.num_segments(); + } + + DORIS_CHECK_GE(segment_start, 0); + DORIS_CHECK_LE(segment_start, segment_end); + DORIS_CHECK_LE(segment_end, rowset_meta.num_segments()); + DORIS_CHECK_EQ(added_iterators, static_cast(segment_end - segment_start)); + DORIS_CHECK_GT(rowset_meta.segment_group_sizes().size(), 1); + + const size_t flags_start = iterator_init_flag->size(); + int64_t group_start = 0; + for (const auto group_size : rowset_meta.segment_group_sizes()) { + DORIS_CHECK_GT(group_size, 0); + const int64_t group_end = group_start + group_size; + const int64_t selected_start = std::max(group_start, segment_start); + const int64_t selected_end = std::min(group_end, segment_end); + if (selected_start < selected_end) { + iterator_init_flag->push_back(true); + iterator_init_flag->insert(iterator_init_flag->end(), + static_cast(selected_end - selected_start - 1), + false); + } + group_start = group_end; + } + DORIS_CHECK_EQ(group_start, rowset_meta.num_segments()); + DORIS_CHECK_EQ(iterator_init_flag->size() - flags_start, added_iterators); +} + Status VerticalBlockReader::next_block_with_aggregation(Block* block, bool* eof) { auto res = (this->*_next_block_func)(block, eof); if (!config::is_cloud_mode()) { @@ -80,23 +114,35 @@ Status VerticalBlockReader::_get_segment_iterators(const ReaderParams& read_para return res; } for (const auto& rs_split : read_params.rs_splits) { + RETURN_IF_ERROR(rs_split.rs_reader->init(&_reader_context, rs_split)); + const auto rowset = rs_split.rs_reader->rowset(); // segment iterator will be inited here // In vertical compaction, every group will load segment so we should cache // segment to avoid tot many s3 head request - bool use_cache = !rs_split.rs_reader->rowset()->is_local(); + bool use_cache = !rowset->is_local(); + size_t iterators_start = segment_iters->size(); RETURN_IF_ERROR(rs_split.rs_reader->get_segment_iterators(&_reader_context, segment_iters, use_cache)); - // if segments overlapping, all segment iterator should be inited in - // heap merge iterator. If segments are none overlapping, only first segment of this - // rowset will be inited and push to heap, other segment will be inited later when current - // segment reached it's end. - // Use this iterator_init_flag so we can load few segments in HeapMergeIterator to save memory - if (rs_split.rs_reader->rowset()->is_segments_overlapping()) { - for (int i = 0; i < rs_split.rs_reader->rowset()->num_segments(); ++i) { - iterator_init_flag->push_back(true); + // If segments overlap, all segment iterators should be initialized in the heap merge + // iterator. For grouped segments, initialize the first iterator of each group because + // segments are ordered within a group but groups may overlap. If segments do not overlap, + // only the first segment of the rowset is initialized; later segments are initialized when + // the current segment reaches its end. + // Use this iterator_init_flag so we can load few segments in HeapMergeIterator to save + // memory. + size_t added_iterators = segment_iters->size() - iterators_start; + if (rowset->is_segments_overlapping()) { + const auto& rowset_meta = *rowset->rowset_meta(); + if (rowset_meta.segments_overlap() == NONOVERLAPPING_WITHIN_GROUP) { + _append_grouped_iterator_init_flags(rowset_meta, rs_split.segment_offsets, + added_iterators, iterator_init_flag); + } else { + for (size_t i = 0; i < added_iterators; ++i) { + iterator_init_flag->push_back(true); + } } } else { - for (int i = 0; i < rs_split.rs_reader->rowset()->num_segments(); ++i) { + for (size_t i = 0; i < added_iterators; ++i) { if (i == 0) { iterator_init_flag->push_back(true); continue; @@ -104,8 +150,8 @@ Status VerticalBlockReader::_get_segment_iterators(const ReaderParams& read_para iterator_init_flag->push_back(false); } } - for (int i = 0; i < rs_split.rs_reader->rowset()->num_segments(); ++i) { - rowset_ids->push_back(rs_split.rs_reader->rowset()->rowset_id()); + for (size_t i = 0; i < added_iterators; ++i) { + rowset_ids->push_back(rowset->rowset_id()); } rs_split.rs_reader->reset_read_options(); } diff --git a/be/src/storage/iterator/vertical_block_reader.h b/be/src/storage/iterator/vertical_block_reader.h index baba2a215f8d19..2907ef52b2058e 100644 --- a/be/src/storage/iterator/vertical_block_reader.h +++ b/be/src/storage/iterator/vertical_block_reader.h @@ -43,8 +43,10 @@ namespace doris { struct RowsetId; class RowSourcesBuffer; +class RowsetMeta; struct RowBatch; struct VerticalCompactionContextStats; +class VerticalBlockReaderTestAccessor; class VerticalBlockReader final : public TabletReader { public: @@ -72,6 +74,13 @@ class VerticalBlockReader final : public TabletReader { static uint64_t nextId; private: + friend class VerticalBlockReaderTestAccessor; + + static void _append_grouped_iterator_init_flags(const RowsetMeta& rowset_meta, + std::pair segment_offsets, + size_t added_iterators, + std::vector* iterator_init_flag); + // Directly read row from rowset and pass to upper caller. No need to do aggregation. // This is usually used for DUPLICATE KEY tables Status _direct_next_block(Block* block, bool* eof); diff --git a/be/src/storage/merger.cpp b/be/src/storage/merger.cpp index f5364d487877aa..2d7712afe86819 100644 --- a/be/src/storage/merger.cpp +++ b/be/src/storage/merger.cpp @@ -68,7 +68,8 @@ namespace doris { Status Merger::vmerge_rowsets(BaseTabletSPtr tablet, ReaderType reader_type, const TabletSchema& cur_tablet_schema, const std::vector& src_rowset_readers, - RowsetWriter* dst_rowset_writer, Statistics* stats_output) { + RowsetWriter* dst_rowset_writer, Statistics* stats_output, + std::optional> segment_range) { if (!cur_tablet_schema.cluster_key_uids().empty()) { return Status::InternalError( "mow table with cluster keys does not support non vertical compaction"); @@ -81,7 +82,11 @@ Status Merger::vmerge_rowsets(BaseTabletSPtr tablet, ReaderType reader_type, TabletReadSource read_source; read_source.rs_splits.reserve(src_rowset_readers.size()); for (const RowsetReaderSharedPtr& rs_reader : src_rowset_readers) { - read_source.rs_splits.emplace_back(rs_reader); + auto& rs_split = read_source.rs_splits.emplace_back(rs_reader); + if (segment_range.has_value()) { + DCHECK_EQ(src_rowset_readers.size(), 1); + rs_split.segment_offsets = segment_range.value(); + } } read_source.fill_delete_predicates(); reader_params.set_read_source(std::move(read_source)); @@ -252,7 +257,7 @@ Status Merger::vertical_compact_one_group( RowsetWriter* dst_rowset_writer, uint32_t max_rows_per_segment, Statistics* stats_output, std::vector key_group_cluster_key_idxes, int64_t batch_size, CompactionSampleInfo* sample_info, VerticalCompactionContextStats* context_stats, - bool enable_sparse_optimization) { + bool enable_sparse_optimization, std::optional> segment_range) { // build tablet reader VLOG_NOTICE << "vertical compact one group, max_rows_per_segment=" << max_rows_per_segment; VerticalBlockReader reader(row_source_buf, context_stats); @@ -266,7 +271,11 @@ Status Merger::vertical_compact_one_group( TabletReadSource read_source; read_source.rs_splits.reserve(src_rowset_readers.size()); for (const RowsetReaderSharedPtr& rs_reader : src_rowset_readers) { - read_source.rs_splits.emplace_back(rs_reader); + auto& rs_split = read_source.rs_splits.emplace_back(rs_reader); + if (segment_range.has_value()) { + DCHECK_EQ(src_rowset_readers.size(), 1); + rs_split.segment_offsets = segment_range.value(); + } } read_source.fill_delete_predicates(); reader_params.set_read_source(std::move(read_source)); @@ -494,7 +503,8 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t RowsetWriter* dst_rowset_writer, uint32_t max_rows_per_segment, int64_t merge_way_num, Statistics* stats_output, - VerticalCompactionProgressCallback progress_cb) { + VerticalCompactionProgressCallback progress_cb, + std::optional> segment_range) { LOG(INFO) << "Start to do vertical compaction, tablet_id: " << tablet->tablet_id(); VerticalCompactionContextStats context_stats; Defer log_context_stats {[&] { @@ -537,29 +547,33 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t progress_cb(column_groups.size(), 0); } - // Calculate total rows for density calculation after compaction + // Segment-range vertical compaction only sees part of a rowset. Do not use partial rows to + // update tablet-level density or drive sparse optimization. int64_t total_rows = 0; - for (const auto& rs_reader : src_rowset_readers) { - total_rows += rs_reader->rowset()->rowset_meta()->num_rows(); - } - // Use historical density for sparse wide table optimization // density = (total_cells - null_cells) / total_cells, smaller means more sparse // When density <= threshold, enable sparse optimization // threshold = 0 means disable, 1 means always enable (default) bool enable_sparse_optimization = false; - if (config::sparse_column_compaction_threshold_percent > 0 && - tablet->keys_type() == KeysType::UNIQUE_KEYS) { - double density = tablet->compaction_density.load(); - enable_sparse_optimization = density <= config::sparse_column_compaction_threshold_percent; + if (!segment_range.has_value()) { + for (const auto& rs_reader : src_rowset_readers) { + total_rows += rs_reader->rowset()->rowset_meta()->num_rows(); + } - LOG(INFO) << "Vertical compaction sparse optimization check: tablet_id=" - << tablet->tablet_id() << ", density=" << density - << ", threshold=" << config::sparse_column_compaction_threshold_percent - << ", total_rows=" << total_rows - << ", num_columns=" << tablet_schema.num_columns() - << ", total_cells=" << total_rows * tablet_schema.num_columns() - << ", enable_sparse_optimization=" << enable_sparse_optimization; + if (config::sparse_column_compaction_threshold_percent > 0 && + tablet->keys_type() == KeysType::UNIQUE_KEYS) { + double density = tablet->compaction_density.load(); + enable_sparse_optimization = + density <= config::sparse_column_compaction_threshold_percent; + + LOG(INFO) << "Vertical compaction sparse optimization check: tablet_id=" + << tablet->tablet_id() << ", density=" << density + << ", threshold=" << config::sparse_column_compaction_threshold_percent + << ", total_rows=" << total_rows + << ", num_columns=" << tablet_schema.num_columns() + << ", total_cells=" << total_rows * tablet_schema.num_columns() + << ", enable_sparse_optimization=" << enable_sparse_optimization; + } } RowSourcesBuffer row_sources_buf(tablet->tablet_id(), dst_rowset_writer->context().tablet_path, @@ -594,7 +608,7 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t } } } - if (need_footer_collection) { + if (!segment_range.has_value() && need_footer_collection) { for (const auto& rs_reader : src_rowset_readers) { auto beta_rowset = std::dynamic_pointer_cast(rs_reader->rowset()); if (!beta_rowset) { @@ -608,7 +622,8 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t << ", rowset_id: " << beta_rowset->rowset_id() << ", status: " << st; continue; } - for (const auto& segment : segments) { + for (int64_t segment_idx = 0; segment_idx < segments.size(); ++segment_idx) { + const auto& segment = segments[segment_idx]; int64_t row_count = segment->num_rows(); auto collect_st = segment->traverse_column_meta_pbs( [&](const segment_v2::ColumnMetaPB& meta) { @@ -630,7 +645,7 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t // Pre-compute per-row estimate for each column group from footer data. std::vector group_per_row_from_footer(column_groups.size(), 0); - std::vector group_footer_fallback(column_groups.size(), false); + std::vector group_footer_fallback(column_groups.size(), segment_range.has_value()); for (size_t i = 0; i < column_groups.size(); ++i) { int64_t group_per_row = 0; bool need_fallback = false; @@ -708,7 +723,7 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t tablet, reader_type, tablet_schema, is_key, column_groups[i], &row_sources_buf, src_rowset_readers, dst_rowset_writer, max_rows_per_segment, group_stats_ptr, key_group_cluster_key_idxes, batch_size, &sample_info, &context_stats, - enable_sparse_optimization); + enable_sparse_optimization, segment_range); { std::unique_lock lock(sample_info_lock); sample_infos[i] = sample_info; @@ -739,7 +754,7 @@ Status Merger::vertical_merge_rowsets(BaseTabletSPtr tablet, ReaderType reader_t // Calculate and store density for next compaction's sparse optimization threshold // density = (total_cells - total_null_count) / total_cells // Smaller density means more sparse - { + if (!segment_range.has_value()) { std::unique_lock lock(sample_info_lock); int64_t total_null_count = 0; for (const auto& info : sample_infos) { diff --git a/be/src/storage/merger.h b/be/src/storage/merger.h index 5185c97e6b0dc9..a53d795500cd64 100644 --- a/be/src/storage/merger.h +++ b/be/src/storage/merger.h @@ -18,6 +18,8 @@ #pragma once #include +#include +#include #include #include "common/status.h" @@ -63,15 +65,17 @@ class Merger { // return OK and set statistics into `*stats_output`. // return others on error - static Status vmerge_rowsets(BaseTabletSPtr tablet, ReaderType reader_type, - const TabletSchema& cur_tablet_schema, - const std::vector& src_rowset_readers, - RowsetWriter* dst_rowset_writer, Statistics* stats_output); + static Status vmerge_rowsets( + BaseTabletSPtr tablet, ReaderType reader_type, const TabletSchema& cur_tablet_schema, + const std::vector& src_rowset_readers, + RowsetWriter* dst_rowset_writer, Statistics* stats_output, + std::optional> segment_range = std::nullopt); static Status vertical_merge_rowsets( BaseTabletSPtr tablet, ReaderType reader_type, const TabletSchema& tablet_schema, const std::vector& src_rowset_readers, RowsetWriter* dst_rowset_writer, uint32_t max_rows_per_segment, int64_t merge_way_num, - Statistics* stats_output, VerticalCompactionProgressCallback progress_cb = nullptr); + Statistics* stats_output, VerticalCompactionProgressCallback progress_cb = nullptr, + std::optional> segment_range = std::nullopt); // for vertical compaction static void vertical_split_columns(const TabletSchema& tablet_schema, @@ -86,7 +90,8 @@ class Merger { RowsetWriter* dst_rowset_writer, uint32_t max_rows_per_segment, Statistics* stats_output, std::vector key_group_cluster_key_idxes, int64_t batch_size, CompactionSampleInfo* sample_info, - VerticalCompactionContextStats* context_stats, bool enable_sparse_optimization = false); + VerticalCompactionContextStats* context_stats, bool enable_sparse_optimization = false, + std::optional> segment_range = std::nullopt); // for segcompaction static Status vertical_compact_one_group( diff --git a/be/src/storage/rowid_conversion.h b/be/src/storage/rowid_conversion.h index 4a516f0b60e365..213ee4c4ef3722 100644 --- a/be/src/storage/rowid_conversion.h +++ b/be/src/storage/rowid_conversion.h @@ -21,6 +21,7 @@ #include #include "common/cast_set.h" +#include "common/check.h" #include "runtime/thread_context.h" #include "storage/olap_common.h" #include "storage/utils.h" @@ -41,6 +42,15 @@ class RowIdConversion { // resize segment rowid map to its rows num Status init_segment_map(const RowsetId& src_rowset_id, const std::vector& num_rows) { for (size_t i = 0; i < num_rows.size(); i++) { + auto src_segment = std::pair {src_rowset_id, cast_set(i)}; + auto iter = _segment_to_id_map.find(src_segment); + // Each segment-group reader initializes all source segments, so reuse existing maps. + if (iter != _segment_to_id_map.end()) { + DORIS_CHECK_LT(iter->second, _segments_rowid_map.size()); + DORIS_CHECK_EQ(_segments_rowid_map[iter->second].size(), num_rows[i]); + continue; + } + constexpr size_t RESERVED_MEMORY = 10 * 1024 * 1024; // 10M if (doris::GlobalMemoryArbitrator::is_exceed_hard_mem_limit(RESERVED_MEMORY)) { return Status::MemoryLimitExceeded(fmt::format( @@ -60,9 +70,10 @@ class RowIdConversion { ->consumption())); } - uint32_t id = static_cast(_segments_rowid_map.size()); - _segment_to_id_map.emplace(std::pair {src_rowset_id, i}, id); - _id_to_segment_map.emplace_back(src_rowset_id, i); + uint32_t id = cast_set(_segments_rowid_map.size()); + auto insert_result = _segment_to_id_map.emplace(src_segment, id); + DORIS_CHECK(insert_result.second); + _id_to_segment_map.push_back(src_segment); std::vector> vec( num_rows[i], std::pair(UINT32_MAX, UINT32_MAX)); diff --git a/be/src/storage/rowset/rowset_meta.h b/be/src/storage/rowset/rowset_meta.h index 51b11494862434..c35a8195c608b6 100644 --- a/be/src/storage/rowset/rowset_meta.h +++ b/be/src/storage/rowset/rowset_meta.h @@ -159,6 +159,22 @@ class RowsetMeta : public MetadataAdder { auto& get_num_segment_rows() const { return _rowset_meta_pb.num_segment_rows(); } + void set_segment_group_sizes(const std::vector& segment_group_sizes) { + DORIS_CHECK_GT(segment_group_sizes.size(), 1); + int64_t segment_count = 0; + for (const auto group_size : segment_group_sizes) { + DORIS_CHECK_GT(group_size, 0); + segment_count += group_size; + } + DORIS_CHECK_EQ(segment_count, num_segments()); + _rowset_meta_pb.mutable_segment_group_sizes()->Assign(segment_group_sizes.cbegin(), + segment_group_sizes.cend()); + } + + void clear_segment_group_sizes() { _rowset_meta_pb.clear_segment_group_sizes(); } + + const auto& segment_group_sizes() const { return _rowset_meta_pb.segment_group_sizes(); } + int64_t total_disk_size() const { return _rowset_meta_pb.total_disk_size(); } void set_total_disk_size(int64_t total_disk_size) { @@ -293,15 +309,17 @@ class RowsetMeta : public MetadataAdder { // 1. the rowset contains more than one segment // 2. the rowset's start version == end version (non-singleton rowset was generated by compaction process // which always produces non-overlapped segments) - // 3. segments_overlap() flag is not NONOVERLAPPING (OVERLAP_UNKNOWN and OVERLAPPING are OK) + // 3. segments_overlap() flag is not NONOVERLAPPING (OVERLAP_UNKNOWN, OVERLAPPING, + // and NONOVERLAPPING_WITHIN_GROUP are considered overlapping) bool is_segments_overlapping() const { return num_segments() > 1 && is_singleton_delta() && segments_overlap() != NONOVERLAPPING; } bool produced_by_compaction() const { - return has_version() && - (start_version() < end_version() || - (start_version() == end_version() && segments_overlap() == NONOVERLAPPING)); + return has_version() && (start_version() < end_version() || + (start_version() == end_version() && + (segments_overlap() == NONOVERLAPPING || + segments_overlap() == NONOVERLAPPING_WITHIN_GROUP))); } // get the compaction score of this rowset. diff --git a/be/src/storage/rowset/vertical_beta_rowset_writer.cpp b/be/src/storage/rowset/vertical_beta_rowset_writer.cpp index c6e7a12b91bc6a..4655726ee0ccc6 100644 --- a/be/src/storage/rowset/vertical_beta_rowset_writer.cpp +++ b/be/src/storage/rowset/vertical_beta_rowset_writer.cpp @@ -139,8 +139,7 @@ Status VerticalBetaRowsetWriter::_flush_columns(segment_v2::SegmentWriter* se key_bounds.set_min_key(min_key.to_string()); key_bounds.set_max_key(max_key.to_string()); this->_segments_encoded_key_bounds.emplace_back(std::move(key_bounds)); - this->_segment_num_rows.resize(_cur_writer_idx + 1); - this->_segment_num_rows[_cur_writer_idx] = _segment_writers[_cur_writer_idx]->row_count(); + this->_segment_num_rows.emplace_back(segment_writer->row_count()); } return Status::OK(); } @@ -204,6 +203,7 @@ template requires std::is_base_of_v Status VerticalBetaRowsetWriter::final_flush() { for (auto& segment_writer : _segment_writers) { + DCHECK(segment_writer); uint64_t segment_size = 0; //uint64_t footer_position = 0; segment_v2::SegmentIndexFileCacheInfo index_file_cache_info; @@ -219,6 +219,8 @@ Status VerticalBetaRowsetWriter::final_flush() { segment_writer.reset(); _record_segment_index_file_cache_preload(segment_id, index_file_cache_info); } + _segment_writers.clear(); + _cur_writer_idx = 0; return Status::OK(); } diff --git a/be/test/cloud/cloud_compaction_test.cpp b/be/test/cloud/cloud_compaction_test.cpp index 62141d5ad883c5..0088086d8034a2 100644 --- a/be/test/cloud/cloud_compaction_test.cpp +++ b/be/test/cloud/cloud_compaction_test.cpp @@ -24,6 +24,7 @@ #include #include #include +#include #include #include @@ -51,6 +52,20 @@ namespace doris { class TabletMap; +namespace { + +void expect_segment_group_merge_ranges(const std::vector& actual, + const std::vector& expected) { + ASSERT_EQ(actual.size(), expected.size()); + for (size_t i = 0; i < expected.size(); ++i) { + EXPECT_EQ(actual[i].segment_start, expected[i].segment_start); + EXPECT_EQ(actual[i].segment_end, expected[i].segment_end); + EXPECT_EQ(actual[i].merge_way_num, expected[i].merge_way_num); + } +} + +} // namespace + class CloudCompactionTest : public testing::Test { CloudCompactionTest() : _engine(CloudStorageEngine(EngineOptions {})) {} void SetUp() override { @@ -324,7 +339,7 @@ TEST_F(CloudCompactionTest, generate_cloud_compaction_tasks_clears_metrics_witho } static RowsetSharedPtr create_rowset(Version version, int num_segments, bool overlapping, - int data_size) { + int data_size, int num_key_columns = 1) { auto rs_meta = std::make_shared(); rs_meta->set_rowset_type(BETA_ROWSET); // important rs_meta->_rowset_meta_pb.set_start_version(version.first); @@ -332,6 +347,19 @@ static RowsetSharedPtr create_rowset(Version version, int num_segments, bool ove rs_meta->set_num_segments(num_segments); rs_meta->set_segments_overlap(overlapping ? OVERLAPPING : NONOVERLAPPING); rs_meta->set_total_disk_size(data_size); + TabletSchemaPB tablet_schema_pb; + tablet_schema_pb.set_keys_type(DUP_KEYS); + for (int i = 0; i < num_key_columns + 1; ++i) { + ColumnPB* column = tablet_schema_pb.add_column(); + column->set_unique_id(i); + column->set_name("c" + std::to_string(i)); + column->set_type("INT"); + column->set_is_key(i < num_key_columns); + column->set_is_nullable(false); + } + auto tablet_schema = std::make_shared(); + tablet_schema->init_from_pb(tablet_schema_pb); + rs_meta->set_tablet_schema(tablet_schema); RowsetSharedPtr rowset; Status st = RowsetFactory::create_rowset(nullptr, "", rs_meta, &rowset); if (!st.ok()) { @@ -1084,6 +1112,345 @@ TEST_F(CloudCompactionTest, should_cache_compaction_output) { LOG(INFO) << "should_cache_compaction_output done"; } +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_execution_path_conditions) { + auto old_enable = config::enable_cloud_single_rowset_compaction; + auto old_min_segments = config::cloud_single_rowset_compaction_min_segments; + auto old_group_size = config::cloud_single_rowset_compaction_segment_group_size; + Defer restore_config {[&] { + config::enable_cloud_single_rowset_compaction = old_enable; + config::cloud_single_rowset_compaction_min_segments = old_min_segments; + config::cloud_single_rowset_compaction_segment_group_size = old_group_size; + }}; + config::enable_cloud_single_rowset_compaction = true; + config::cloud_single_rowset_compaction_min_segments = 4; + config::cloud_single_rowset_compaction_segment_group_size = 2; + + RowsetSharedPtr candidate = create_rowset(Version(2, 2), 4, true, 1024); + ASSERT_TRUE(candidate != nullptr); + const auto& tablet_schema = *candidate->tablet_schema(); + EXPECT_TRUE(cloud::is_single_rowset_compaction_candidate(candidate)); + EXPECT_TRUE(cloud::should_use_single_rowset_grouped_compaction({candidate}, tablet_schema, + CUMULATIVE_SIZE_BASED_POLICY)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction({candidate}, tablet_schema, + CUMULATIVE_TIME_SERIES_POLICY)); + + TabletSchemaPB cluster_key_schema_pb; + tablet_schema.to_schema_pb(&cluster_key_schema_pb); + cluster_key_schema_pb.set_keys_type(UNIQUE_KEYS); + cluster_key_schema_pb.add_cluster_key_uids(1); + TabletSchema cluster_key_schema; + cluster_key_schema.init_from_pb(cluster_key_schema_pb); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction({candidate}, cluster_key_schema, + CUMULATIVE_SIZE_BASED_POLICY)); + + config::enable_cloud_single_rowset_compaction = false; + EXPECT_TRUE(cloud::is_single_rowset_compaction_candidate(candidate)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction({candidate}, tablet_schema, + CUMULATIVE_SIZE_BASED_POLICY)); + config::enable_cloud_single_rowset_compaction = true; + + RowsetSharedPtr non_overlapping = create_rowset(Version(3, 3), 4, false, 1024); + ASSERT_TRUE(non_overlapping != nullptr); + EXPECT_FALSE(cloud::is_single_rowset_compaction_candidate(non_overlapping)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction( + {non_overlapping}, tablet_schema, CUMULATIVE_SIZE_BASED_POLICY)); + + RowsetSharedPtr too_few_segments = create_rowset(Version(4, 4), 3, true, 1024); + ASSERT_TRUE(too_few_segments != nullptr); + EXPECT_FALSE(cloud::is_single_rowset_compaction_candidate(too_few_segments)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction( + {too_few_segments}, tablet_schema, CUMULATIVE_SIZE_BASED_POLICY)); + + RowsetSharedPtr grouped_candidate = create_rowset(Version(5, 5), 8, true, 1024); + ASSERT_TRUE(grouped_candidate != nullptr); + grouped_candidate->rowset_meta()->set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + grouped_candidate->rowset_meta()->set_segment_group_sizes({2, 2, 2, 2}); + EXPECT_TRUE(cloud::is_single_rowset_compaction_candidate(grouped_candidate)); + + RowsetSharedPtr grouped_with_too_few_groups = create_rowset(Version(6, 6), 8, true, 1024); + ASSERT_TRUE(grouped_with_too_few_groups != nullptr); + grouped_with_too_few_groups->rowset_meta()->set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + grouped_with_too_few_groups->rowset_meta()->set_segment_group_sizes({3, 3, 2}); + EXPECT_FALSE(cloud::is_single_rowset_compaction_candidate(grouped_with_too_few_groups)); + + RowsetSharedPtr no_key_columns = create_rowset(Version(7, 7), 4, true, 1024, 0); + ASSERT_TRUE(no_key_columns != nullptr); + EXPECT_TRUE(cloud::is_single_rowset_compaction_candidate(no_key_columns)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction( + {no_key_columns}, *no_key_columns->tablet_schema(), CUMULATIVE_SIZE_BASED_POLICY)); + + RowsetSharedPtr with_delete_predicate = create_rowset(Version(8, 8), 4, true, 1024); + ASSERT_TRUE(with_delete_predicate != nullptr); + DeletePredicatePB delete_predicate; + auto* in_predicate = delete_predicate.add_in_predicates(); + in_predicate->set_column_name("c1"); + in_predicate->add_values("1"); + with_delete_predicate->rowset_meta()->set_delete_predicate(std::move(delete_predicate)); + EXPECT_FALSE(cloud::is_single_rowset_compaction_candidate(with_delete_predicate)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction( + {with_delete_predicate}, tablet_schema, CUMULATIVE_SIZE_BASED_POLICY)); + + RowsetSharedPtr another_candidate = create_rowset(Version(9, 9), 4, true, 1024); + ASSERT_TRUE(another_candidate != nullptr); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction( + {candidate, another_candidate}, tablet_schema, CUMULATIVE_SIZE_BASED_POLICY)); + EXPECT_FALSE(cloud::should_use_single_rowset_grouped_compaction({}, tablet_schema, + CUMULATIVE_SIZE_BASED_POLICY)); + + CloudTabletSPtr tablet = std::make_shared(_engine, _tablet_meta); + CloudCumulativeCompaction compaction(_engine, tablet); + compaction._input_rowsets = {candidate}; + compaction._cur_tablet_schema = candidate->tablet_schema(); + compaction._single_rowset_compaction_segment_group_size = + config::cloud_single_rowset_compaction_segment_group_size; + Compaction::MergeInputRowsetsResult result; + ASSERT_TRUE(compaction.prepare_merge_input_rowsets(&result).ok()); + EXPECT_TRUE(compaction._single_rowset_compaction_segment_group_size.has_value()); + EXPECT_TRUE(result.is_segment_grouped); + EXPECT_EQ(result.segment_group_size, config::cloud_single_rowset_compaction_segment_group_size); + + _tablet_meta->set_compaction_policy(std::string(CUMULATIVE_TIME_SERIES_POLICY)); + CloudCumulativeCompaction time_series_compaction(_engine, tablet); + time_series_compaction._input_rowsets = {candidate}; + time_series_compaction._cur_tablet_schema = candidate->tablet_schema(); + Compaction::MergeInputRowsetsResult time_series_result; + ASSERT_TRUE(time_series_compaction.prepare_merge_input_rowsets(&time_series_result).ok()); + EXPECT_FALSE(time_series_compaction._single_rowset_compaction_segment_group_size.has_value()); + EXPECT_FALSE(time_series_result.is_segment_grouped); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_builds_logical_group_ranges) { + RowsetMeta overlapping_meta; + overlapping_meta.set_num_segments(5); + overlapping_meta.set_segments_overlap(OVERLAPPING); + + const auto overlapping_ranges = cloud::build_segment_group_merge_ranges(overlapping_meta, 2); + expect_segment_group_merge_ranges(overlapping_ranges, + {{.segment_start = 0, .segment_end = 2, .merge_way_num = 2}, + {.segment_start = 2, .segment_end = 4, .merge_way_num = 2}, + {.segment_start = 4, .segment_end = 5, .merge_way_num = 1}}); + + const auto single_overlapping_range = + cloud::build_segment_group_merge_ranges(overlapping_meta, 10); + expect_segment_group_merge_ranges(single_overlapping_range, + {{.segment_start = 0, .segment_end = 5, .merge_way_num = 5}}); + + overlapping_meta.set_segments_overlap(NONOVERLAPPING); + const auto nonoverlapping_ranges = cloud::build_segment_group_merge_ranges(overlapping_meta, 2); + expect_segment_group_merge_ranges(nonoverlapping_ranges, + {{.segment_start = 0, .segment_end = 2, .merge_way_num = 2}, + {.segment_start = 2, .segment_end = 4, .merge_way_num = 2}, + {.segment_start = 4, .segment_end = 5, .merge_way_num = 1}}); + + overlapping_meta.set_segments_overlap(OVERLAP_UNKNOWN); + const auto unknown_overlap_ranges = + cloud::build_segment_group_merge_ranges(overlapping_meta, 2); + expect_segment_group_merge_ranges(unknown_overlap_ranges, + {{.segment_start = 0, .segment_end = 2, .merge_way_num = 2}, + {.segment_start = 2, .segment_end = 4, .merge_way_num = 2}, + {.segment_start = 4, .segment_end = 5, .merge_way_num = 1}}); + + RowsetMeta grouped_meta; + grouped_meta.set_num_segments(5); + grouped_meta.set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + grouped_meta.set_segment_group_sizes({2, 2, 1}); + + const auto grouped_ranges = cloud::build_segment_group_merge_ranges(grouped_meta, 2); + expect_segment_group_merge_ranges(grouped_ranges, + {{.segment_start = 0, .segment_end = 4, .merge_way_num = 2}, + {.segment_start = 4, .segment_end = 5, .merge_way_num = 1}}); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_builds_group_range_boundaries) { + RowsetMeta grouped_meta; + grouped_meta.set_num_segments(5); + grouped_meta.set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + grouped_meta.set_segment_group_sizes({2, 2, 1}); + + const auto single_range = cloud::build_segment_group_merge_ranges(grouped_meta, 10); + expect_segment_group_merge_ranges(single_range, + {{.segment_start = 0, .segment_end = 5, .merge_way_num = 3}}); + + grouped_meta.set_num_segments(10); + grouped_meta.set_segment_group_sizes({1, 2, 3, 4}); + const auto exact_ranges = cloud::build_segment_group_merge_ranges(grouped_meta, 2); + expect_segment_group_merge_ranges( + exact_ranges, {{.segment_start = 0, .segment_end = 3, .merge_way_num = 2}, + {.segment_start = 3, .segment_end = 10, .merge_way_num = 2}}); + + grouped_meta.set_num_segments(15); + grouped_meta.set_segment_group_sizes({3, 1, 4, 2, 5}); + const auto irregular_ranges = cloud::build_segment_group_merge_ranges(grouped_meta, 2); + expect_segment_group_merge_ranges( + irregular_ranges, {{.segment_start = 0, .segment_end = 4, .merge_way_num = 2}, + {.segment_start = 4, .segment_end = 10, .merge_way_num = 2}, + {.segment_start = 10, .segment_end = 15, .merge_way_num = 1}}); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_rejects_invalid_range_input) { + RowsetMeta rowset_meta; + rowset_meta.set_num_segments(5); + rowset_meta.set_segments_overlap(OVERLAPPING); + EXPECT_DEATH(static_cast(cloud::build_segment_group_merge_ranges(rowset_meta, 1)), ""); + + RowsetMeta empty_rowset_meta; + empty_rowset_meta.set_segments_overlap(OVERLAPPING); + EXPECT_DEATH(static_cast(cloud::build_segment_group_merge_ranges(empty_rowset_meta, 2)), + ""); + + RowsetMetaPB invalid_group_layout_pb; + invalid_group_layout_pb.set_rowset_id(1); + invalid_group_layout_pb.set_num_segments(5); + invalid_group_layout_pb.set_segments_overlap_pb(NONOVERLAPPING_WITHIN_GROUP); + + RowsetMeta empty_group_layout; + ASSERT_TRUE(empty_group_layout.init_from_pb(invalid_group_layout_pb)); + EXPECT_DEATH(static_cast(cloud::build_segment_group_merge_ranges(empty_group_layout, 2)), + ""); + + invalid_group_layout_pb.add_segment_group_sizes(2); + invalid_group_layout_pb.add_segment_group_sizes(2); + RowsetMeta invalid_group_layout; + ASSERT_TRUE(invalid_group_layout.init_from_pb(invalid_group_layout_pb)); + EXPECT_DEATH( + static_cast(cloud::build_segment_group_merge_ranges(invalid_group_layout, 2)), + ""); + + invalid_group_layout_pb.clear_segment_group_sizes(); + invalid_group_layout_pb.add_segment_group_sizes(2); + invalid_group_layout_pb.add_segment_group_sizes(0); + invalid_group_layout_pb.add_segment_group_sizes(3); + RowsetMeta zero_sized_group_layout; + ASSERT_TRUE(zero_sized_group_layout.init_from_pb(invalid_group_layout_pb)); + EXPECT_DEATH( + static_cast(cloud::build_segment_group_merge_ranges(zero_sized_group_layout, 2)), + ""); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_calculates_cumulative_point) { + CloudTabletSPtr tablet = std::make_shared(_engine, _tablet_meta); + CloudCumulativeCompaction compaction(_engine, tablet); + compaction._input_rowsets = {create_rowset(Version(2, 2), 5, true, 1024)}; + compaction._output_rowset = create_rowset(Version(2, 2), 5, false, 1024); + ASSERT_TRUE(compaction._input_rowsets.front() != nullptr); + ASSERT_TRUE(compaction._output_rowset != nullptr); + compaction._single_rowset_compaction_segment_group_size = 2; + + EXPECT_TRUE(compaction.should_calculate_new_cumulative_point(2)); + EXPECT_FALSE(compaction.should_calculate_new_cumulative_point(1)); + + compaction._output_rowset->rowset_meta()->set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + compaction._output_rowset->rowset_meta()->set_segment_group_sizes({2, 2, 1}); + EXPECT_FALSE(compaction.should_calculate_new_cumulative_point(2)); + + compaction._single_rowset_compaction_segment_group_size.reset(); + EXPECT_TRUE(compaction.should_calculate_new_cumulative_point(1)); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_uses_selection_snapshot) { + auto old_enable = config::enable_cloud_single_rowset_compaction; + auto old_min_segments = config::cloud_single_rowset_compaction_min_segments; + auto old_group_size = config::cloud_single_rowset_compaction_segment_group_size; + Defer restore_config {[&] { + config::enable_cloud_single_rowset_compaction = old_enable; + config::cloud_single_rowset_compaction_min_segments = old_min_segments; + config::cloud_single_rowset_compaction_segment_group_size = old_group_size; + }}; + config::enable_cloud_single_rowset_compaction = true; + config::cloud_single_rowset_compaction_min_segments = 4; + config::cloud_single_rowset_compaction_segment_group_size = 2; + + std::vector rowsets; + auto grouped_rowset = create_rowset(Version(2, 2), 4, true, 1024); + ASSERT_TRUE(grouped_rowset != nullptr); + rowsets.push_back(grouped_rowset); + for (int64_t version = 3; version <= 13; ++version) { + auto rowset = create_rowset(Version(version, version), 1, false, 1024); + ASSERT_TRUE(rowset != nullptr); + rowsets.push_back(std::move(rowset)); + } + + TabletSchemaPB tablet_schema_pb; + grouped_rowset->tablet_schema()->to_schema_pb(&tablet_schema_pb); + _tablet_meta->mutable_tablet_schema()->init_from_pb(tablet_schema_pb); + _tablet_meta->set_compaction_policy(std::string(CUMULATIVE_SIZE_BASED_POLICY)); + CloudTabletSPtr tablet = std::make_shared(_engine, _tablet_meta); + { + std::unique_lock wlock(tablet->get_header_lock()); + tablet->add_rowsets(std::move(rowsets), false, wlock, false); + } + + for (const int32_t invalid_group_size : {0, 1}) { + config::cloud_single_rowset_compaction_segment_group_size = invalid_group_size; + CloudCumulativeCompaction invalid_config_compaction(_engine, tablet); + ASSERT_TRUE(invalid_config_compaction.pick_rowsets_to_compact().ok()); + EXPECT_FALSE( + invalid_config_compaction._single_rowset_compaction_segment_group_size.has_value()); + } + config::cloud_single_rowset_compaction_segment_group_size = 2; + + CloudCumulativeCompaction compaction(_engine, tablet); + ASSERT_TRUE(compaction.pick_rowsets_to_compact().ok()); + ASSERT_EQ(compaction._input_rowsets.size(), 1); + EXPECT_EQ(compaction._input_rowsets.front(), grouped_rowset); + ASSERT_TRUE(compaction._single_rowset_compaction_segment_group_size.has_value()); + EXPECT_EQ(*compaction._single_rowset_compaction_segment_group_size, 2); + + config::enable_cloud_single_rowset_compaction = false; + config::cloud_single_rowset_compaction_min_segments = 5; + config::cloud_single_rowset_compaction_segment_group_size = 3; + + Compaction::MergeInputRowsetsResult result; + ASSERT_TRUE(compaction.prepare_merge_input_rowsets(&result).ok()); + EXPECT_TRUE(compaction._single_rowset_compaction_segment_group_size.has_value()); + EXPECT_TRUE(result.is_segment_grouped); + EXPECT_EQ(result.segment_group_size, 2); +} + +TEST_F(CloudCompactionTest, single_rowset_grouped_compaction_honors_notready_policy_filter) { + auto old_enable = config::enable_cloud_single_rowset_compaction; + auto old_min_segments = config::cloud_single_rowset_compaction_min_segments; + auto old_enable_empty_rowset_compaction = config::enable_empty_rowset_compaction; + Defer restore_config {[&] { + config::enable_cloud_single_rowset_compaction = old_enable; + config::cloud_single_rowset_compaction_min_segments = old_min_segments; + config::enable_empty_rowset_compaction = old_enable_empty_rowset_compaction; + }}; + config::enable_cloud_single_rowset_compaction = true; + config::cloud_single_rowset_compaction_min_segments = 4; + config::enable_empty_rowset_compaction = false; + + std::vector rowsets; + // Keep enough older inputs mergeable after the NOTREADY policy filters versions 11 through 20. + for (int64_t version = 2; version <= 19; ++version) { + auto rowset = create_rowset(Version(version, version), 1, false, 1024); + ASSERT_TRUE(rowset != nullptr); + rowsets.push_back(std::move(rowset)); + } + auto filtered_grouped_rowset = create_rowset(Version(20, 20), 4, true, 1024); + ASSERT_TRUE(filtered_grouped_rowset != nullptr); + rowsets.push_back(filtered_grouped_rowset); + + TabletSchemaPB tablet_schema_pb; + filtered_grouped_rowset->tablet_schema()->to_schema_pb(&tablet_schema_pb); + _tablet_meta->mutable_tablet_schema()->init_from_pb(tablet_schema_pb); + _tablet_meta->set_compaction_policy(std::string(CUMULATIVE_SIZE_BASED_POLICY)); + _tablet_meta->set_tablet_state(TABLET_NOTREADY); + CloudTabletSPtr tablet = std::make_shared(_engine, _tablet_meta); + tablet->set_alter_version(1); + { + std::unique_lock wlock(tablet->get_header_lock()); + tablet->add_rowsets(std::move(rowsets), false, wlock, false); + } + + CloudCumulativeCompaction compaction(_engine, tablet); + ASSERT_TRUE(compaction.pick_rowsets_to_compact().ok()); + ASSERT_EQ(compaction._input_rowsets.size(), 9); + EXPECT_EQ(compaction._input_rowsets.front()->version(), Version(2, 2)); + EXPECT_EQ(compaction._input_rowsets.back()->version(), Version(10, 10)); + EXPECT_FALSE(compaction._single_rowset_compaction_segment_group_size.has_value()); +} + TEST_F(CloudCompactionTest, test_truncate_rowsets_by_txn_size_empty_input) { std::vector rowsets; int64_t kept_size = 100; diff --git a/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp b/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp index 6e9fa1e1aaa13e..79ec7f14d20fa9 100644 --- a/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp +++ b/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp @@ -185,6 +185,25 @@ TEST_F(TestCloudSizeBasedCumulativeCompactionPolicy, EXPECT_EQ(Version(-1, -1), last_delete_version); } +TEST_F(TestCloudSizeBasedCumulativeCompactionPolicy, + overlapping_output_uses_policy_to_advance_cumulative_point) { + _tablet_meta->set_tablet_state(TABLET_RUNNING); + _tablet_meta->set_time_series_compaction_level_threshold(1); + CloudTablet tablet(_engine, _tablet_meta); + tablet._base_size = 100; + + auto output_rowset = create_rowset(Version(3, 3), 5, true, 256 * 1024 * 1024); + Version last_delete_version {-1, -1}; + + CloudSizeBasedCumulativeCompactionPolicy size_based_policy; + EXPECT_EQ(4, size_based_policy.new_cumulative_point(&tablet, output_rowset, last_delete_version, + 2)); + + CloudTimeSeriesCumulativeCompactionPolicy time_series_policy; + EXPECT_EQ(4, time_series_policy.new_cumulative_point(&tablet, output_rowset, + last_delete_version, 2)); +} + // Test case: Empty rowset compaction with skip_trim TEST_F(TestCloudSizeBasedCumulativeCompactionPolicy, pick_input_rowsets_empty_rowset_compaction) { // Save original config values diff --git a/be/test/cloud/cloud_snapshot_mgr_test.cpp b/be/test/cloud/cloud_snapshot_mgr_test.cpp index 5b097be04d4685..64e039b46f5131 100644 --- a/be/test/cloud/cloud_snapshot_mgr_test.cpp +++ b/be/test/cloud/cloud_snapshot_mgr_test.cpp @@ -101,6 +101,9 @@ TEST_F(CloudSnapshotMgrTest, TestConvertRowsets) { rowset_meta->set_tablet_id(1000); rowset_meta->set_txn_id(2000); rowset_meta->set_num_segments(3); + rowset_meta->set_segments_overlap_pb(NONOVERLAPPING_WITHIN_GROUP); + rowset_meta->add_segment_group_sizes(1); + rowset_meta->add_segment_group_sizes(2); rowset_meta->set_num_rows(100); rowset_meta->set_start_version(100); rowset_meta->set_end_version(101); @@ -148,6 +151,10 @@ TEST_F(CloudSnapshotMgrTest, TestConvertRowsets) { EXPECT_EQ(output_meta_pb.rs_metas(0).tablet_schema().index(0).index_id(), 1001); EXPECT_EQ(output_meta_pb.rs_metas(0).tablet_schema().index(1).index_id(), 1002); EXPECT_EQ(output_meta_pb.rs_metas(0).num_segments(), 3); + EXPECT_EQ(output_meta_pb.rs_metas(0).segments_overlap_pb(), NONOVERLAPPING_WITHIN_GROUP); + ASSERT_EQ(output_meta_pb.rs_metas(0).segment_group_sizes_size(), 2); + EXPECT_EQ(output_meta_pb.rs_metas(0).segment_group_sizes(0), 1); + EXPECT_EQ(output_meta_pb.rs_metas(0).segment_group_sizes(1), 2); EXPECT_EQ(output_meta_pb.rs_metas(0).num_rows(), 100); EXPECT_EQ(output_meta_pb.rs_metas(0).start_version(), 100); EXPECT_EQ(output_meta_pb.rs_metas(0).end_version(), 101); diff --git a/be/test/storage/cloud_file_cache_write_index_only_test.cpp b/be/test/storage/cloud_file_cache_write_index_only_test.cpp index 49f43e80cf16b0..ba937a64a83efd 100644 --- a/be/test/storage/cloud_file_cache_write_index_only_test.cpp +++ b/be/test/storage/cloud_file_cache_write_index_only_test.cpp @@ -699,21 +699,26 @@ TEST_F(CloudFileCacheWriteIndexOnlyTest, auto writer_result = RowsetFactory::create_rowset_writer(*_engine, context, true); ASSERT_TRUE(writer_result.has_value()) << writer_result.error(); auto rowset_writer = std::move(writer_result).value(); + EXPECT_EQ(rowset_writer->get_allocated_segment_id(), 0); std::vector key_column_ids = {0}; auto key_block = create_column_block(tablet_schema, key_column_ids, 8, 1); auto st = rowset_writer->add_columns(&key_block, key_column_ids, true, 4, false); ASSERT_TRUE(st.ok()) << st; + EXPECT_EQ(rowset_writer->get_allocated_segment_id(), 1); auto second_key_block = create_column_block(tablet_schema, key_column_ids, 8, 100); st = rowset_writer->add_columns(&second_key_block, key_column_ids, true, 4, false); ASSERT_TRUE(st.ok()) << st; + EXPECT_EQ(rowset_writer->get_allocated_segment_id(), 2); st = rowset_writer->flush_columns(true); ASSERT_TRUE(st.ok()) << st; + EXPECT_EQ(rowset_writer->get_allocated_segment_id(), 2); std::vector value_column_ids = {1}; auto value_block = create_column_block(tablet_schema, value_column_ids, 16, 1); st = rowset_writer->add_columns(&value_block, value_column_ids, false, UINT32_MAX, false); ASSERT_TRUE(st.ok()) << st; + EXPECT_EQ(rowset_writer->get_allocated_segment_id(), 2); st = rowset_writer->flush_columns(false); ASSERT_TRUE(st.ok()) << st; st = rowset_writer->final_flush(); diff --git a/be/test/storage/iterator/vertical_block_reader_test.cpp b/be/test/storage/iterator/vertical_block_reader_test.cpp new file mode 100644 index 00000000000000..52ee75b2d2fcba --- /dev/null +++ b/be/test/storage/iterator/vertical_block_reader_test.cpp @@ -0,0 +1,136 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include "storage/iterator/vertical_block_reader.h" + +#include + +#include "storage/rowset/rowset_meta.h" + +namespace doris { + +class VerticalBlockReaderTestAccessor { +public: + static void append_grouped_iterator_init_flags(const RowsetMeta& rowset_meta, + std::pair segment_offsets, + size_t added_iterators, + std::vector* iterator_init_flags) { + VerticalBlockReader::_append_grouped_iterator_init_flags( + rowset_meta, segment_offsets, added_iterators, iterator_init_flags); + } + + static std::vector grouped_iterator_init_flags( + const RowsetMeta& rowset_meta, std::pair segment_offsets, + size_t added_iterators) { + std::vector iterator_init_flags; + append_grouped_iterator_init_flags(rowset_meta, segment_offsets, added_iterators, + &iterator_init_flags); + return iterator_init_flags; + } +}; + +TEST(VerticalBlockReaderTest, GroupedIteratorInitFlagsForWholeRowset) { + RowsetMeta rowset_meta; + rowset_meta.set_num_segments(5); + rowset_meta.set_segment_group_sizes({2, 2, 1}); + + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {0, 0}, 5), + std::vector({true, false, true, false, true})); +} + +TEST(VerticalBlockReaderTest, GroupedIteratorInitFlagsForSegmentRange) { + RowsetMeta rowset_meta; + rowset_meta.set_num_segments(8); + rowset_meta.set_segment_group_sizes({3, 2, 3}); + + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {0, 3}, 3), + std::vector({true, false, false})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {1, 3}, 2), + std::vector({true, false})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {1, 2}, 1), + std::vector({true})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {1, 4}, 3), + std::vector({true, false, true})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {3, 5}, 2), + std::vector({true, false})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {4, 7}, 3), + std::vector({true, true, false})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {5, 8}, 3), + std::vector({true, false, false})); + EXPECT_EQ(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {7, 8}, 1), + std::vector({true})); +} + +TEST(VerticalBlockReaderTest, GroupedIteratorInitFlagsAppendToExistingFlags) { + RowsetMeta rowset_meta; + rowset_meta.set_num_segments(5); + rowset_meta.set_segment_group_sizes({2, 2, 1}); + std::vector iterator_init_flags {false, true}; + + VerticalBlockReaderTestAccessor::append_grouped_iterator_init_flags(rowset_meta, {1, 4}, 3, + &iterator_init_flags); + EXPECT_EQ(iterator_init_flags, std::vector({false, true, true, true, false})); + + VerticalBlockReaderTestAccessor::append_grouped_iterator_init_flags(rowset_meta, {4, 5}, 1, + &iterator_init_flags); + EXPECT_EQ(iterator_init_flags, std::vector({false, true, true, true, false, true})); +} + +TEST(VerticalBlockReaderTest, GroupedIteratorInitFlagsRejectInvalidInput) { + RowsetMeta rowset_meta; + rowset_meta.set_num_segments(5); + rowset_meta.set_segment_group_sizes({2, 2, 1}); + + RowsetMetaPB invalid_group_layout_pb; + invalid_group_layout_pb.set_rowset_id(1); + invalid_group_layout_pb.set_num_segments(5); + invalid_group_layout_pb.add_segment_group_sizes(2); + invalid_group_layout_pb.add_segment_group_sizes(2); + RowsetMeta invalid_group_layout; + ASSERT_TRUE(invalid_group_layout.init_from_pb(invalid_group_layout_pb)); + +#ifndef NDEBUG + GTEST_FLAG_SET(death_test_style, "threadsafe"); + EXPECT_DEATH( + { + static_cast(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags( + rowset_meta, {1, 4}, 2)); + }, + "Check failed"); + EXPECT_DEATH( + { + static_cast(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags( + rowset_meta, {0, 6}, 6)); + }, + "Check failed"); + EXPECT_DEATH( + { + static_cast(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags( + invalid_group_layout, {0, 0}, 5)); + }, + "Check failed"); +#else + EXPECT_ANY_THROW(static_cast( + VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {1, 4}, 2))); + EXPECT_ANY_THROW(static_cast( + VerticalBlockReaderTestAccessor::grouped_iterator_init_flags(rowset_meta, {0, 6}, 6))); + EXPECT_ANY_THROW(static_cast(VerticalBlockReaderTestAccessor::grouped_iterator_init_flags( + invalid_group_layout, {0, 0}, 5))); +#endif +} + +} // namespace doris diff --git a/be/test/storage/pb_convert_test.cpp b/be/test/storage/pb_convert_test.cpp index 9e03b023a4b74c..d7548b5b11cd15 100644 --- a/be/test/storage/pb_convert_test.cpp +++ b/be/test/storage/pb_convert_test.cpp @@ -374,6 +374,30 @@ TEST(PbConvert, test_rvalue_overloads) { << "\n diff=" << print(set_diff(tablet_meta_cloud_set_fields, tablet_meta_out_set_fields)); } +TEST(PbConvert, test_rowset_segment_group_sizes_conversion) { + RowsetMetaPB rs; + rs.set_segments_overlap_pb(NONOVERLAPPING_WITHIN_GROUP); + rs.add_segment_group_sizes(3); + rs.add_segment_group_sizes(2); + rs.add_segment_group_sizes(5); + + RowsetMetaCloudPB rs_cloud; + doris_rowset_meta_to_cloud(&rs_cloud, rs); + ASSERT_EQ(rs_cloud.segments_overlap_pb(), NONOVERLAPPING_WITHIN_GROUP); + ASSERT_EQ(rs_cloud.segment_group_sizes_size(), 3); + EXPECT_EQ(rs_cloud.segment_group_sizes(0), 3); + EXPECT_EQ(rs_cloud.segment_group_sizes(1), 2); + EXPECT_EQ(rs_cloud.segment_group_sizes(2), 5); + + RowsetMetaPB rs_out; + cloud_rowset_meta_to_doris(&rs_out, rs_cloud); + ASSERT_EQ(rs_out.segments_overlap_pb(), NONOVERLAPPING_WITHIN_GROUP); + ASSERT_EQ(rs_out.segment_group_sizes_size(), 3); + EXPECT_EQ(rs_out.segment_group_sizes(0), 3); + EXPECT_EQ(rs_out.segment_group_sizes(1), 2); + EXPECT_EQ(rs_out.segment_group_sizes(2), 5); +} + TEST(PbConvert, test_return_value_overloads) { // rowset meta RowsetMetaPB rs; diff --git a/be/test/storage/rowid_conversion_test.cpp b/be/test/storage/rowid_conversion_test.cpp index c8b3bad9336a04..2df224f2766c95 100644 --- a/be/test/storage/rowid_conversion_test.cpp +++ b/be/test/storage/rowid_conversion_test.cpp @@ -29,11 +29,16 @@ #include #include +#include #include #include #include #include +#include "cloud/cloud_cumulative_compaction.h" +#include "cloud/cloud_storage_engine.h" +#include "cloud/cloud_tablet.h" +#include "common/config.h" #include "common/status.h" #include "core/block/block.h" #include "core/block/column_with_type_and_name.h" @@ -43,7 +48,10 @@ #include "io/io_common.h" #include "json2pb/json_to_pb.h" #include "runtime/exec_env.h" +#include "storage/compaction_task_tracker.h" #include "storage/delete/delete_handler.h" +#include "storage/index/index_writer.h" +#include "storage/index/inverted/inverted_index_desc.h" #include "storage/merger.h" #include "storage/options.h" #include "storage/rowset/beta_rowset.h" @@ -54,10 +62,12 @@ #include "storage/rowset/rowset_reader_context.h" #include "storage/rowset/rowset_writer.h" #include "storage/rowset/rowset_writer_context.h" +#include "storage/segment/segment.h" #include "storage/storage_engine.h" #include "storage/tablet/tablet.h" #include "storage/tablet/tablet_meta.h" #include "storage/tablet/tablet_schema.h" +#include "util/defer_op.h" #include "util/uid_util.h" namespace doris { @@ -79,6 +89,14 @@ class TestRowIdConversion : public testing::TestWithParamcreate_directory(absolute_dir + "/tablet_path") .ok()); + + std::vector tmp_paths; + tmp_paths.emplace_back(absolute_dir, 1024000000); + auto tmp_file_dirs = std::make_unique(tmp_paths); + st = tmp_file_dirs->init(); + ASSERT_TRUE(st.ok()) << st; + ExecEnv::GetInstance()->set_tmp_file_dir(std::move(tmp_file_dirs)); + doris::EngineOptions options; auto engine = std::make_unique(options); engine_ref = engine.get(); @@ -86,12 +104,14 @@ class TestRowIdConversion : public testing::TestWithParamset_tmp_file_dir(nullptr); EXPECT_TRUE(io::global_local_filesystem()->delete_directory(absolute_dir).ok()); engine_ref = nullptr; ExecEnv::GetInstance()->set_storage_engine(nullptr); } - TabletSchemaSPtr create_schema(KeysType keys_type = DUP_KEYS) { + TabletSchemaSPtr create_schema(KeysType keys_type = DUP_KEYS, + bool with_inverted_index = false) { TabletSchemaSPtr tablet_schema = std::make_shared(); TabletSchemaPB tablet_schema_pb; tablet_schema_pb.set_keys_type(keys_type); @@ -135,6 +155,15 @@ class TestRowIdConversion : public testing::TestWithParamset_is_bf_column(false); } + if (with_inverted_index) { + tablet_schema_pb.set_inverted_index_storage_format(InvertedIndexStorageFormatPB::V2); + auto* index = tablet_schema_pb.add_index(); + index->set_index_id(1); + index->set_index_name("c2_idx"); + index->set_index_type(IndexType::INVERTED); + index->add_col_unique_id(2); + } + tablet_schema->init_from_pb(tablet_schema_pb); return tablet_schema; } @@ -185,7 +214,7 @@ class TestRowIdConversion : public testing::TestWithParam>> input_data; + for (int64_t segment_id = 0; segment_id < num_segments; ++segment_id) { + std::vector> segment_data; + for (int64_t row_id = 0; row_id < rows_per_segment; ++row_id) { + int64_t key = segment_id * rows_per_segment + row_id; + segment_data.emplace_back(key, key + 1); + } + input_data.push_back(std::move(segment_data)); + } + + CloudStorageEngine cloud_engine(EngineOptions {}); + for (bool is_vertical : {false, true}) { + SCOPED_TRACE(is_vertical ? "vertical merge" : "horizontal merge"); + + TabletSchemaSPtr tablet_schema = create_schema(UNIQUE_KEYS, true); + tablet_schema->set_schema_version(schema_version); + tablet_schema->set_db_id(1000); + RowsetSharedPtr input_rowset = create_rowset(tablet_schema, OVERLAPPING, input_data, 2); + ASSERT_TRUE(input_rowset != nullptr); + + TabletSharedPtr local_tablet = create_tablet(*tablet_schema, true); + auto writer_context = create_rowset_writer_context( + tablet_schema, NONOVERLAPPING, rows_per_segment, input_rowset->version()); + writer_context.tablet_id = local_tablet->tablet_id(); + writer_context.index_id = local_tablet->index_id(); + writer_context.partition_id = local_tablet->partition_id(); + writer_context.tablet_schema_hash = local_tablet->schema_hash(); + writer_context.tablet_uid = local_tablet->tablet_uid(); + writer_context.newest_write_timestamp = newest_write_timestamp; + writer_context.compaction_level = compaction_level; + writer_context.enable_unique_key_merge_on_write = true; + auto writer_result = + RowsetFactory::create_rowset_writer(*engine_ref, writer_context, is_vertical); + ASSERT_TRUE(writer_result.has_value()) << writer_result.error(); + + auto cloud_tablet = std::make_shared( + cloud_engine, std::make_shared(*local_tablet->tablet_meta())); + CloudCumulativeCompaction compaction(cloud_engine, cloud_tablet); + compaction._input_rowsets = {input_rowset}; + compaction._cur_tablet_schema = tablet_schema; + compaction._output_rs_writer = std::move(writer_result).value(); + compaction._is_vertical = is_vertical; + compaction._input_row_num = input_rowset->num_rows(); + compaction._input_rowsets_data_size = input_rowset->data_disk_size(); + compaction._stats.rowid_conversion = compaction._rowid_conversion.get(); + + auto* compaction_task_tracker = CompactionTaskTracker::instance(); + CompactionTaskInfo task_info; + task_info.compaction_id = compaction.compaction_id(); + compaction_task_tracker->register_task(std::move(task_info)); + Defer remove_tracker_task { + [compaction_task_tracker, compaction_id = compaction.compaction_id()] { + compaction_task_tracker->remove_task(compaction_id); + }}; + + Compaction::MergeInputRowsetsResult merge_result; + merge_result.is_segment_grouped = true; + merge_result.segment_group_size = segment_group_size; + ASSERT_TRUE(compaction.do_merge_input_rowsets({}, &merge_result).ok()); + const int64_t segment_group_count = + (num_segments + segment_group_size - 1) / segment_group_size; + ASSERT_EQ(merge_result.output_segment_group_sizes.size(), + static_cast(segment_group_count)); + for (const auto output_group_size : merge_result.output_segment_group_sizes) { + EXPECT_GT(output_group_size, 0); + } + if (is_vertical) { + constexpr int32_t default_num_columns_per_group = 5; + const int32_t num_columns_per_group = + config::vertical_compaction_num_columns_per_group != + default_num_columns_per_group + ? config::vertical_compaction_num_columns_per_group + : cloud_tablet->tablet_meta() + ->vertical_compaction_num_columns_per_group(); + std::vector> column_groups; + std::vector key_group_cluster_key_idxes; + Merger::vertical_split_columns(*tablet_schema, &column_groups, + &key_group_cluster_key_idxes, num_columns_per_group); + + const auto tracked_tasks = compaction_task_tracker->get_all_tasks(); + const auto task_it = std::find_if( + tracked_tasks.begin(), tracked_tasks.end(), [&](const auto& tracked_task) { + return tracked_task.compaction_id == compaction.compaction_id(); + }); + ASSERT_NE(task_it, tracked_tasks.end()); + const int64_t expected_total_groups = + static_cast(column_groups.size()) * segment_group_count; + EXPECT_EQ(task_it->vertical_total_groups, expected_total_groups); + EXPECT_EQ(task_it->vertical_completed_groups, expected_total_groups); + } + + RowsetSharedPtr output_rowset; + ASSERT_EQ(Status::OK(), compaction._output_rs_writer->build(output_rowset)); + ASSERT_TRUE(output_rowset != nullptr); + compaction._output_rowset = output_rowset; + compaction.update_output_rowset_after_build(merge_result); + EXPECT_EQ(compaction._stats.output_rows, input_rowset->num_rows()); + if (is_vertical) { + EXPECT_GT(output_rowset->num_segments(), + static_cast(merge_result.output_segment_group_sizes.size())); + } + EXPECT_EQ(output_rowset->rowset_meta()->get_num_segment_rows().size(), + output_rowset->num_segments()); + + EXPECT_EQ(output_rowset->rowset_meta()->segments_overlap(), NONOVERLAPPING_WITHIN_GROUP); + const auto output_rowset_pb = output_rowset->rowset_meta()->get_rowset_pb(); + ASSERT_EQ(output_rowset_pb.segment_group_sizes_size(), + merge_result.output_segment_group_sizes.size()); + int64_t output_segment_count = 0; + for (int i = 0; i < output_rowset_pb.segment_group_sizes_size(); ++i) { + EXPECT_EQ(output_rowset_pb.segment_group_sizes(i), + merge_result.output_segment_group_sizes[static_cast(i)]); + output_segment_count += output_rowset_pb.segment_group_sizes(i); + } + EXPECT_EQ(output_segment_count, output_rowset->num_segments()); + + RowsetReaderContext reader_context; + reader_context.tablet_schema = tablet_schema; + reader_context.need_ordered_result = false; + std::vector return_columns = {0, 1}; + reader_context.return_columns = &return_columns; + RowsetReaderSharedPtr output_reader; + create_and_init_rowset_reader(output_rowset.get(), reader_context, &output_reader); + + std::vector> output_data; + Status read_status; + do { + Block output_block = tablet_schema->create_block(); + read_status = output_reader->next_batch(&output_block); + const auto& columns = output_block.get_columns_with_type_and_name(); + ASSERT_EQ(columns.size(), return_columns.size()); + for (size_t row_id = 0; row_id < output_block.rows(); ++row_id) { + output_data.emplace_back(columns[0].column->get_int(row_id), + columns[1].column->get_int(row_id)); + } + } while (read_status.ok()); + ASSERT_TRUE(read_status.is()) << read_status; + ASSERT_EQ(output_data.size(), input_rowset->num_rows()); + + auto beta_rowset = std::dynamic_pointer_cast(output_rowset); + ASSERT_TRUE(beta_rowset != nullptr); + const auto& rowset_meta = output_rowset->rowset_meta(); + const auto rowset_meta_pb = rowset_meta->get_rowset_pb(); + const auto& segment_num_rows_from_meta = rowset_meta->get_num_segment_rows(); + + EXPECT_EQ(rowset_meta_pb.rowset_id_v2(), writer_context.rowset_id.to_string()); + EXPECT_EQ(rowset_meta->rowset_id().to_string(), writer_context.rowset_id.to_string()); + + EXPECT_EQ(rowset_meta->rowset_type(), BETA_ROWSET); + EXPECT_EQ(rowset_meta->rowset_state(), VISIBLE); + EXPECT_TRUE(rowset_meta->has_version()); + EXPECT_EQ(rowset_meta->version(), input_rowset->version()); + EXPECT_FALSE(rowset_meta->empty()); + EXPECT_EQ(rowset_meta->num_rows(), input_rowset->num_rows()); + EXPECT_EQ(rowset_meta->segments_overlap(), NONOVERLAPPING_WITHIN_GROUP); + EXPECT_TRUE(rowset_meta->is_segments_overlapping()); + + ASSERT_TRUE(rowset_meta->tablet_schema() != nullptr); + EXPECT_EQ(*rowset_meta->tablet_schema(), *tablet_schema); + EXPECT_TRUE(rowset_meta_pb.has_tablet_schema()); + TabletSchemaPB expected_tablet_schema_pb; + tablet_schema->to_schema_pb(&expected_tablet_schema_pb); + EXPECT_EQ(rowset_meta_pb.tablet_schema().SerializeAsString(), + expected_tablet_schema_pb.SerializeAsString()); + EXPECT_EQ(rowset_meta_pb.schema_version(), schema_version); + EXPECT_TRUE(rowset_meta_pb.has_has_variant_type_in_schema()); + EXPECT_FALSE(rowset_meta_pb.has_variant_type_in_schema()); + + std::vector output_segments; + ASSERT_TRUE(beta_rowset->load_segments(&output_segments).ok()); + ASSERT_EQ(rowset_meta->num_segments(), output_segments.size()); + ASSERT_EQ(segment_num_rows_from_meta.size(), output_segments.size()); + + const auto& segment_key_bounds_from_meta = rowset_meta->get_segments_key_bounds(); + EXPECT_FALSE(rowset_meta->is_segments_key_bounds_aggregated()); + EXPECT_FALSE(rowset_meta->is_segments_key_bounds_truncated()); + ASSERT_EQ(segment_key_bounds_from_meta.size(), output_segments.size()); + + const auto& inverted_index_file_info_from_meta = rowset_meta->inverted_index_file_info(); + EXPECT_TRUE(rowset_meta_pb.enable_inverted_index_file_info()); + ASSERT_EQ(inverted_index_file_info_from_meta.size(), output_segments.size()); + + EXPECT_FALSE(rowset_meta_pb.enable_segments_file_size()); + EXPECT_TRUE(rowset_meta_pb.segments_file_size().empty()); + + int64_t actual_data_disk_size = 0; + int64_t actual_index_disk_size = 0; + int64_t actual_num_rows = 0; + for (size_t segment_id = 0; segment_id < output_segments.size(); ++segment_id) { + EXPECT_EQ(output_segments[segment_id]->id(), segment_id); + EXPECT_EQ(segment_num_rows_from_meta[segment_id], + output_segments[segment_id]->num_rows()) + << "segment_id=" << segment_id; + EXPECT_EQ(segment_key_bounds_from_meta[segment_id].min_key(), + output_segments[segment_id]->min_key()) + << "segment_id=" << segment_id; + EXPECT_EQ(segment_key_bounds_from_meta[segment_id].max_key(), + output_segments[segment_id]->max_key()) + << "segment_id=" << segment_id; + + const auto& index_file_info = inverted_index_file_info_from_meta[segment_id]; + ASSERT_TRUE(index_file_info.has_index_size()) << "segment_id=" << segment_id; + const auto segment_path = output_rowset->segment_path(segment_id); + ASSERT_TRUE(segment_path.has_value()) << segment_path.error(); + int64_t segment_file_size = 0; + const auto segment_file_size_status = + rowset_meta->fs()->file_size(segment_path.value(), &segment_file_size); + ASSERT_TRUE(segment_file_size_status.ok()) << segment_file_size_status; + actual_data_disk_size += segment_file_size; + actual_num_rows += output_segments[segment_id]->num_rows(); + + const auto index_file_path = + segment_v2::InvertedIndexDescriptor::get_index_file_path_v2( + segment_v2::InvertedIndexDescriptor::get_index_file_path_prefix( + segment_path.value())); + int64_t index_file_size = 0; + const auto index_file_size_status = + rowset_meta->fs()->file_size(index_file_path, &index_file_size); + ASSERT_TRUE(index_file_size_status.ok()) << index_file_size_status; + EXPECT_EQ(index_file_info.index_size(), index_file_size) << "segment_id=" << segment_id; + actual_index_disk_size += index_file_size; + } + EXPECT_EQ(rowset_meta->num_rows(), actual_num_rows); + EXPECT_EQ(rowset_meta->data_disk_size(), actual_data_disk_size); + EXPECT_EQ(rowset_meta->index_disk_size(), actual_index_disk_size); + EXPECT_EQ(rowset_meta->total_disk_size(), + rowset_meta->data_disk_size() + rowset_meta->index_disk_size()); + + std::vector output_segment_num_rows; + OlapReaderStatistics reader_stats; + ASSERT_TRUE( + beta_rowset->get_segment_num_rows(&output_segment_num_rows, false, &reader_stats) + .ok()); + + RowIdConversion& rowid_conversion = *compaction._stats.rowid_conversion; + EXPECT_EQ(rowid_conversion.get_src_segment_to_id_map().size(), num_segments); + EXPECT_EQ(rowid_conversion.get_rowid_conversion_map().size(), num_segments); + EXPECT_EQ(rowid_conversion.get_rowid_conversion_map().size(), + rowid_conversion.get_src_segment_to_id_map().size()); + for (int64_t segment_id = 0; segment_id < num_segments; ++segment_id) { + for (int64_t row_id = 0; row_id < rows_per_segment; ++row_id) { + RowLocation src(input_rowset->rowset_id(), segment_id, row_id); + RowLocation dst; + ASSERT_EQ(rowid_conversion.get(src, &dst), 0) + << "segment_id=" << segment_id << ", row_id=" << row_id; + ASSERT_LT(dst.segment_id, output_segment_num_rows.size()); + ASSERT_LT(dst.row_id, output_segment_num_rows[dst.segment_id]); + + size_t output_row_id = dst.row_id; + for (uint32_t output_segment_id = 0; output_segment_id < dst.segment_id; + ++output_segment_id) { + output_row_id += output_segment_num_rows[output_segment_id]; + } + ASSERT_LT(output_row_id, output_data.size()); + EXPECT_EQ(output_data[output_row_id], input_data[segment_id][row_id]); + } + } + + if (is_vertical) { + auto second_writer_context = create_rowset_writer_context( + tablet_schema, NONOVERLAPPING, rows_per_segment, output_rowset->version()); + second_writer_context.tablet_id = local_tablet->tablet_id(); + second_writer_context.index_id = local_tablet->index_id(); + second_writer_context.partition_id = local_tablet->partition_id(); + second_writer_context.tablet_schema_hash = local_tablet->schema_hash(); + second_writer_context.tablet_uid = local_tablet->tablet_uid(); + second_writer_context.newest_write_timestamp = newest_write_timestamp; + second_writer_context.compaction_level = compaction_level; + second_writer_context.enable_unique_key_merge_on_write = true; + auto second_writer_result = + RowsetFactory::create_rowset_writer(*engine_ref, second_writer_context, true); + ASSERT_TRUE(second_writer_result.has_value()) << second_writer_result.error(); + + CloudCumulativeCompaction second_compaction(cloud_engine, cloud_tablet); + second_compaction._input_rowsets = {output_rowset}; + second_compaction._cur_tablet_schema = tablet_schema; + second_compaction._output_rs_writer = std::move(second_writer_result).value(); + second_compaction._is_vertical = true; + second_compaction._input_row_num = output_rowset->num_rows(); + second_compaction._input_rowsets_data_size = output_rowset->data_disk_size(); + second_compaction._stats.rowid_conversion = second_compaction._rowid_conversion.get(); + + Compaction::MergeInputRowsetsResult second_merge_result; + second_merge_result.is_segment_grouped = true; + second_merge_result.segment_group_size = output_rowset->num_segments(); + ASSERT_TRUE(second_compaction.do_merge_input_rowsets({}, &second_merge_result).ok()); + ASSERT_EQ(second_merge_result.output_segment_group_sizes.size(), 1); + + RowsetSharedPtr second_output_rowset; + ASSERT_EQ(Status::OK(), + second_compaction._output_rs_writer->build(second_output_rowset)); + ASSERT_TRUE(second_output_rowset != nullptr); + second_compaction._output_rowset = second_output_rowset; + second_compaction.update_output_rowset_after_build(second_merge_result); + EXPECT_EQ(second_output_rowset->rowset_meta()->segments_overlap(), NONOVERLAPPING); + EXPECT_TRUE(second_output_rowset->rowset_meta()->segment_group_sizes().empty()); + EXPECT_EQ(second_output_rowset->num_rows(), output_rowset->num_rows()); + EXPECT_GT(second_output_rowset->num_segments(), 1); + + auto second_beta_rowset = std::dynamic_pointer_cast(second_output_rowset); + ASSERT_TRUE(second_beta_rowset != nullptr); + std::vector second_output_segments; + ASSERT_TRUE(second_beta_rowset->load_segments(&second_output_segments).ok()); + ASSERT_EQ(second_output_segments.size(), second_output_rowset->num_segments()); + for (size_t segment_id = 0; segment_id < second_output_segments.size(); ++segment_id) { + const auto& segment = second_output_segments[segment_id]; + EXPECT_LE(segment->min_key(), segment->max_key()) << "segment_id=" << segment_id; + if (segment_id > 0) { + EXPECT_LT(second_output_segments[segment_id - 1]->max_key(), segment->min_key()) + << "previous_segment_id=" << segment_id - 1 + << ", segment_id=" << segment_id; + } + } + + RowsetReaderContext second_reader_context; + second_reader_context.tablet_schema = tablet_schema; + second_reader_context.need_ordered_result = false; + second_reader_context.return_columns = &return_columns; + RowsetReaderSharedPtr second_output_reader; + create_and_init_rowset_reader(second_output_rowset.get(), second_reader_context, + &second_output_reader); + + std::vector> second_output_data; + Status second_read_status; + do { + Block output_block = tablet_schema->create_block(); + second_read_status = second_output_reader->next_batch(&output_block); + const auto& columns = output_block.get_columns_with_type_and_name(); + for (size_t row_id = 0; row_id < output_block.rows(); ++row_id) { + second_output_data.emplace_back(columns[0].column->get_int(row_id), + columns[1].column->get_int(row_id)); + } + } while (second_read_status.ok()); + ASSERT_TRUE(second_read_status.is()) << second_read_status; + EXPECT_EQ(second_output_data, output_data); + } + } +} + INSTANTIATE_TEST_SUITE_P( Parameters, TestRowIdConversion, ::testing::ValuesIn(std::vector> { diff --git a/be/test/storage/rowset/rowset_meta_test.cpp b/be/test/storage/rowset/rowset_meta_test.cpp index 954b8e9ed22515..010791c69b9972 100644 --- a/be/test/storage/rowset/rowset_meta_test.cpp +++ b/be/test/storage/rowset/rowset_meta_test.cpp @@ -124,6 +124,27 @@ TEST_F(RowsetMetaTest, TestRowsetIdInit) { EXPECT_EQ(id.to_string(), "72057594037927935"); } +TEST_F(RowsetMetaTest, SegmentGroupsKeepRowsetOverlappingSemantics) { + RowsetMeta rowset_meta; + rowset_meta.set_version({10, 10}); + rowset_meta.set_num_segments(5); + rowset_meta.set_segments_overlap(NONOVERLAPPING_WITHIN_GROUP); + rowset_meta.set_segment_group_sizes({2, 3}); + + EXPECT_TRUE(rowset_meta.is_segments_overlapping()); + EXPECT_TRUE(rowset_meta.produced_by_compaction()); + EXPECT_EQ(rowset_meta.get_compaction_score(), 5); + EXPECT_EQ(rowset_meta.get_merge_way_num(), 5); + + auto rowset_meta_pb = rowset_meta.get_rowset_pb(); + ASSERT_EQ(rowset_meta_pb.segment_group_sizes_size(), 2); + EXPECT_EQ(rowset_meta_pb.segment_group_sizes(0), 2); + EXPECT_EQ(rowset_meta_pb.segment_group_sizes(1), 3); + + rowset_meta.clear_segment_group_sizes(); + EXPECT_EQ(rowset_meta.get_rowset_pb().segment_group_sizes_size(), 0); +} + TEST_F(RowsetMetaTest, TestNumSegmentRowsSetAndGet) { RowsetMeta rowset_meta; EXPECT_TRUE(rowset_meta.init_from_json(_json_rowset_meta)); diff --git a/gensrc/proto/olap_file.proto b/gensrc/proto/olap_file.proto index 498f0238efa303..4866aa955296ac 100644 --- a/gensrc/proto/olap_file.proto +++ b/gensrc/proto/olap_file.proto @@ -53,6 +53,9 @@ enum SegmentsOverlapPB { OVERLAP_UNKNOWN = 0; // this enum is added since Doris v0.11, so previous rowset's segment is unknown OVERLAPPING = 1; NONOVERLAPPING = 2; + // Segments are non-overlapping within each group, but groups may overlap with each other. + // The group layout is described by RowsetMetaPB.segment_group_sizes. + NONOVERLAPPING_WITHIN_GROUP = 3; } message KeyBoundsPB { @@ -145,6 +148,10 @@ message RowsetMetaPB { // Only applies to non-MOW rowsets to reduce meta size on cloud FDB. optional bool segments_key_bounds_aggregated = 57; + // Valid only when segments_overlap_pb is NONOVERLAPPING_WITHIN_GROUP. + // Each value is the number of consecutive output segments in one non-overlapping group. + repeated int32 segment_group_sizes = 59; + // For cloud // for data recycling optional int64 txn_expiration = 1000; @@ -254,6 +261,10 @@ message RowsetMetaCloudPB { // Only applies to non-MOW rowsets to reduce meta size on cloud FDB. optional bool segments_key_bounds_aggregated = 57; + // Valid only when segments_overlap_pb is NONOVERLAPPING_WITHIN_GROUP. + // Each value is the number of consecutive output segments in one non-overlapping group. + repeated int32 segment_group_sizes = 59; + // cloud // the field is a vector, rename it repeated int64 segments_file_size = 100; diff --git a/regression-test/suites/cloud_p0/compaction/test_cloud_single_rowset_grouped_compaction.groovy b/regression-test/suites/cloud_p0/compaction/test_cloud_single_rowset_grouped_compaction.groovy new file mode 100644 index 00000000000000..b98d51c243822a --- /dev/null +++ b/regression-test/suites/cloud_p0/compaction/test_cloud_single_rowset_grouped_compaction.groovy @@ -0,0 +1,563 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +import java.util.Base64 + +suite("test_cloud_single_rowset_grouped_compaction", "nonConcurrent") { + if (!isCloudMode()) { + logger.info("not cloud mode, skip this test") + return + } + + int initialInputSegmentsPerGroup = 2 + long compactionTimeoutMs = 90000L + def customBeConfig = [ + doris_scanner_row_bytes: 1, + enable_cloud_single_rowset_compaction: true, + cloud_single_rowset_compaction_min_segments: 2, + cloud_single_rowset_compaction_segment_group_size: initialInputSegmentsPerGroup, + cumulative_compaction_min_deltas: 2, + enable_aggregate_non_mow_key_bounds: false, + disable_auto_compaction: true, + vertical_compaction_max_segment_size: 1073741824, + compaction_batch_size: -1, + enable_rowid_conversion_correctness_check: false + ] + + setBeConfigTemporary(customBeConfig) { + def metaServiceEndpoint = context.config.metaServiceHttpAddress + + def getRowsetMeta = { tabletId, int version -> + def rowsetMeta = null + getSegmentFilesFromMs(metaServiceEndpoint, tabletId, version) { responseCode, body -> + assertEquals(200, responseCode) + rowsetMeta = parseJson(body) + } + assertNotNull(rowsetMeta) + return rowsetMeta + } + + def compareEncodedKeys = { String leftBase64, String rightBase64 -> + byte[] left = Base64.getDecoder().decode(leftBase64) + byte[] right = Base64.getDecoder().decode(rightBase64) + int commonLength = Math.min(left.length, right.length) + for (int i = 0; i < commonLength; ++i) { + int leftByte = Byte.toUnsignedInt(left[i]) + int rightByte = Byte.toUnsignedInt(right[i]) + if (leftByte != rightByte) { + return leftByte <=> rightByte + } + } + return left.length <=> right.length + } + + def logCaseSection = { String description -> + logger.info("========================================================================") + logger.info(description) + logger.info("========================================================================") + } + + def showTablet = { tableName, beHost, bePort, tabletId -> + sql "SELECT COUNT(*) FROM ${tableName}" + def (code, out, err) = be_show_tablet_status(beHost, bePort, tabletId) + logger.info("Show tablet status: code=${code}, out=${out}, err=${err}") + assertEquals(0, code) + return parseJson(out.trim()) + } + + def rowsetByVersion = { tabletJson, int version -> + def rowset = tabletJson.rowsets.find { it.startsWith("[${version}-${version}] ") } + assertNotNull(rowset) + return rowset + } + + def parseRowsetInfo = { rowset -> + def matcher = rowset =~ /\[[0-9]+-[0-9]+\]\s+([0-9]+)\s+DATA\s+([A-Z_]+)/ + assertTrue(matcher.find(), "unexpected rowset format: ${rowset}") + return [segments: matcher.group(1).toInteger(), overlap: matcher.group(2)] + } + + def waitForCompaction = { beHost, bePort, tabletId, timeoutMs -> + long deadline = System.currentTimeMillis() + timeoutMs + def lastStatus = null + while (System.currentTimeMillis() < deadline) { + def (code, out, err) = be_get_compaction_status(beHost, bePort, tabletId) + logger.info("Get compaction status: code=${code}, out=${out}, err=${err}") + assertEquals(0, code) + lastStatus = parseJson(out.trim()) + assertEquals("success", lastStatus.status.toLowerCase()) + if (!lastStatus.run_status) { + return + } + Thread.sleep(1000) + } + assertTrue(false, "compaction did not finish on ${beHost}:${bePort}, " + + "tablet=${tabletId}, timeoutMs=${timeoutMs}, last=${lastStatus}") + } + + def runCumulativeCompaction = { beHost, bePort, tabletId -> + def (code, out, err) = be_run_cumulative_compaction(beHost, bePort, tabletId) + logger.info("Run compaction: code=${code}, out=${out}, err=${err}") + assertEquals(0, code) + def compactJson = parseJson(out.trim()) + assertEquals("success", compactJson.status.toLowerCase()) + waitForCompaction(beHost, bePort, tabletId, compactionTimeoutMs) + } + + def readPointRows = { String tableName -> + def pointResult = sql "SELECT k, v FROM ${tableName} WHERE k = 100 ORDER BY v" + return pointResult.collect { row -> row.collect { column -> column.toString() } } + } + + def checkPointRows = { String tableName, List> expectedPointRows -> + def pointRows = readPointRows(tableName) + if (expectedPointRows != null) { + assertEquals(expectedPointRows, pointRows) + } + return pointRows + } + + def checkSingleRowsetGroupedCompaction = { + String tableName, String keyType, String valueColumn, String extraProperties, + int inputSegmentsPerGroup, int rowsPerLoadRound, int expectedRows, + List> expectedPointRows, + boolean expectMultipleOutputSegmentsPerGroup -> + sql "DROP TABLE IF EXISTS ${tableName}" + sql """ + CREATE TABLE ${tableName} ( + k INT, + ${valueColumn} + ) + ${keyType}(k) + DISTRIBUTED BY HASH(k) BUCKETS 1 + PROPERTIES ( + "replication_num" = "1", + "disable_auto_compaction" = "true"${extraProperties} + ) + """ + + StringBuilder content = new StringBuilder() + for (int i = 0; i < 2; i++) { + (1..rowsPerLoadRound).each { + content.append("${it},${it + i}\n") + } + } + streamLoad { + table "${tableName}" + set "column_separator", "," + inputStream new ByteArrayInputStream(content.toString().getBytes()) + time 30000 + check { result, exception, startTime, endTime -> + if (exception != null) { + throw exception + } + def json = parseJson(result) + assertEquals("success", json.Status.toLowerCase()) + assertEquals(rowsPerLoadRound * 2, json.NumberTotalRows) + assertEquals(0, json.NumberFilteredRows) + } + } + sql "sync" + + def pointRowsBeforeCompaction = checkPointRows(tableName, expectedPointRows) + if (expectedPointRows == null) { + assertEquals(1, pointRowsBeforeCompaction.size()) + } + + def tablets = sql_return_maparray "SHOW TABLETS FROM ${tableName}" + assertEquals(1, tablets.size()) + def tabletId = tablets[0].TabletId + def backendId = tablets[0].BackendId + def backends = sql_return_maparray "SHOW BACKENDS" + def backend = backends.find { it.BackendId == backendId } + assertNotNull(backend) + + def before = showTablet(tableName, backend.Host, backend.HttpPort, tabletId) + def inputRowset = rowsetByVersion(before, 2) + def inputInfo = parseRowsetInfo(inputRowset) + assertEquals("OVERLAPPING", inputInfo.overlap) + assertTrue(inputInfo.segments > inputSegmentsPerGroup, inputRowset) + + set_be_param("cloud_single_rowset_compaction_segment_group_size", + inputSegmentsPerGroup.toString()) + runCumulativeCompaction(backend.Host, backend.HttpPort, tabletId) + + def after = showTablet(tableName, backend.Host, backend.HttpPort, tabletId) + def outputRowset = rowsetByVersion(after, 2) + def outputInfo = parseRowsetInfo(outputRowset) + assertEquals("NONOVERLAPPING_WITHIN_GROUP", outputInfo.overlap) + def expectedOutputSegments = + (inputInfo.segments + inputSegmentsPerGroup - 1) + .intdiv(inputSegmentsPerGroup) + if (expectMultipleOutputSegmentsPerGroup) { + assertTrue(outputInfo.segments > expectedOutputSegments, outputRowset) + } else { + assertEquals(expectedOutputSegments, outputInfo.segments) + } + def countResult = sql "SELECT COUNT(*) FROM ${tableName}" + assertEquals(expectedRows, countResult[0][0]) + def pointRowsAfterCompaction = checkPointRows(tableName, expectedPointRows) + if (expectedPointRows == null && pointRowsAfterCompaction != pointRowsBeforeCompaction) { + logger.warn("Point query result changed after single rowset grouped compaction" + + ", table=${tableName}, before=${pointRowsBeforeCompaction}" + + ", after=${pointRowsAfterCompaction}") + } + + // Compact the grouped rowset as one range. This exercises VerticalBlockReader's + // NONOVERLAPPING_WITHIN_GROUP iterator initialization and must produce a fully + // non-overlapping rowset. + set_be_param("cloud_single_rowset_compaction_segment_group_size", + outputInfo.segments.toString()) + runCumulativeCompaction(backend.Host, backend.HttpPort, tabletId) + + def afterSecondCompaction = + showTablet(tableName, backend.Host, backend.HttpPort, tabletId) + def finalRowset = rowsetByVersion(afterSecondCompaction, 2) + def finalInfo = parseRowsetInfo(finalRowset) + assertEquals("NONOVERLAPPING", finalInfo.overlap) + if (expectMultipleOutputSegmentsPerGroup) { + assertTrue(finalInfo.segments > 1, finalRowset) + def finalRowsetMeta = getRowsetMeta(tabletId, 2) + assertEquals(finalInfo.segments, finalRowsetMeta.num_segments.toString().toInteger()) + assertFalse(finalRowsetMeta.segments_key_bounds_aggregated ?: false) + assertFalse(finalRowsetMeta.segments_key_bounds_truncated ?: false) + def segmentKeyBounds = finalRowsetMeta.segments_key_bounds + assertEquals(finalInfo.segments, segmentKeyBounds.size()) + segmentKeyBounds.eachWithIndex { keyBounds, int segmentId -> + assertTrue(compareEncodedKeys(keyBounds.min_key, keyBounds.max_key) <= 0, + "invalid key range at segment ${segmentId}") + if (segmentId > 0) { + def previousKeyBounds = segmentKeyBounds[segmentId - 1] + assertTrue(compareEncodedKeys( + previousKeyBounds.max_key, keyBounds.min_key) < 0, + "overlapping key ranges at segments ${segmentId - 1} and " + + "${segmentId}") + } + } + } + def finalCountResult = sql "SELECT COUNT(*) FROM ${tableName}" + assertEquals(expectedRows, finalCountResult[0][0]) + def finalPointRows = checkPointRows(tableName, expectedPointRows) + if (expectedPointRows == null && finalPointRows != pointRowsAfterCompaction) { + logger.warn("Point query result changed after repeated single rowset compaction" + + ", table=${tableName}, before=${pointRowsAfterCompaction}" + + ", after=${finalPointRows}") + } + } + + def sqlCacheOrigValue = sql("select @@enable_sql_cache")[0][0] + try { + sql "set enable_sql_cache=false" + GetDebugPoint().clearDebugPointsForAllBEs() + GetDebugPoint().enableDebugPointForAllBEs("MemTable.need_flush") + + // ==================================================== + logCaseSection("Standard cases: compact DUP, AGG, MOW, and MOR rowsets twice " + + "with 2 or 4 input segments per group") + [2, 4].each { int inputSegmentsPerGroup -> + checkSingleRowsetGroupedCompaction( + "test_cloud_single_rowset_grouped_compaction_g${inputSegmentsPerGroup}_dup", + "DUPLICATE KEY", "v INT", "", inputSegmentsPerGroup, 8192, 8192 * 2, + [["100", "100"], ["100", "101"]], false) + checkSingleRowsetGroupedCompaction( + "test_cloud_single_rowset_grouped_compaction_g${inputSegmentsPerGroup}_agg", + "AGGREGATE KEY", "v INT SUM", "", inputSegmentsPerGroup, 8192, 8192, + [["100", "201"]], false) + checkSingleRowsetGroupedCompaction( + "test_cloud_single_rowset_grouped_compaction_g${inputSegmentsPerGroup}_mow", + "UNIQUE KEY", "v INT", + ", \"enable_unique_key_merge_on_write\" = \"true\"", + inputSegmentsPerGroup, 8192, 8192, [["100", "101"]], false) + checkSingleRowsetGroupedCompaction( + "test_cloud_single_rowset_grouped_compaction_g${inputSegmentsPerGroup}_mor", + "UNIQUE KEY", "v INT", + ", \"enable_unique_key_merge_on_write\" = \"false\"", + inputSegmentsPerGroup, 8192, 8192, null, false) + } + + // ==================================================== + logCaseSection("MOW delete-bitmap case: compact each of three overlapping rowsets " + + "twice and verify cross-rowset delete bitmap conversion") + sql "DROP TABLE IF EXISTS test_cloud_single_rowset_compaction_mow_delete_bitmap" + sql """ + CREATE TABLE test_cloud_single_rowset_compaction_mow_delete_bitmap ( + k INT, + v INT + ) + UNIQUE KEY(k) + DISTRIBUTED BY HASH(k) BUCKETS 1 + PROPERTIES ( + "replication_num" = "1", + "disable_auto_compaction" = "true", + "enable_unique_key_merge_on_write" = "true" + ) + """ + + def loadMowRows = { + int startKey, int endKey, int valueBase, int duplicateRounds -> + StringBuilder content = new StringBuilder() + (startKey..endKey).each { int key -> + content.append("${key},${valueBase + key}\n") + } + // Append an inner key range to create duplicates within this rowset. Keep the + // offsets deterministic so a failure can be reproduced exactly. + int duplicateStartKey = startKey + 2048 + int duplicateEndKey = endKey - 3072 + for (int round = 0; round < duplicateRounds; ++round) { + (duplicateStartKey..duplicateEndKey).each { int key -> + content.append("${key},${valueBase + key + 1}\n") + } + } + streamLoad { + table "test_cloud_single_rowset_compaction_mow_delete_bitmap" + set "column_separator", "," + inputStream new ByteArrayInputStream(content.toString().getBytes()) + time 30000 + check { result, exception, startTime, endTime -> + if (exception != null) { + throw exception + } + def json = parseJson(result) + assertEquals("success", json.Status.toLowerCase()) + int originalRows = endKey - startKey + 1 + int duplicateRows = duplicateEndKey - duplicateStartKey + 1 + assertEquals(originalRows + duplicateRows * duplicateRounds, + json.NumberTotalRows) + assertEquals(0, json.NumberFilteredRows) + } + } + sql "sync" + } + + loadMowRows(1, 16384, 100000, 1) + loadMowRows(8193, 24576, 200000, 2) + loadMowRows(12289, 28672, 300000, 3) + + def readMowRows = { + def rows = sql """ + SELECT k, v + FROM test_cloud_single_rowset_compaction_mow_delete_bitmap + WHERE k IN (1, 8192, 8193, 12288, 12289, 16384, 16385, 24576, 24577, 28672) + ORDER BY k + """ + return rows.collect { row -> row.collect { column -> column.toString() } } + } + def expectedMowRows = [ + ["1", "100001"], + ["8192", "108193"], + ["8193", "208193"], + ["12288", "212289"], + ["12289", "312289"], + ["16384", "316385"], + ["16385", "316386"], + ["24576", "324577"], + ["24577", "324578"], + ["28672", "328672"] + ] + assertEquals(expectedMowRows, readMowRows()) + def mowCountBeforeCompaction = + sql "SELECT COUNT(*) FROM test_cloud_single_rowset_compaction_mow_delete_bitmap" + assertEquals(28672, mowCountBeforeCompaction[0][0]) + + def mowTablets = sql_return_maparray( + "SHOW TABLETS FROM test_cloud_single_rowset_compaction_mow_delete_bitmap") + assertEquals(1, mowTablets.size()) + def mowTabletId = mowTablets[0].TabletId + def mowBackendId = mowTablets[0].BackendId + def mowBackends = sql_return_maparray "SHOW BACKENDS" + def mowBackend = mowBackends.find { it.BackendId == mowBackendId } + assertNotNull(mowBackend) + + def checkMowQueryResult = { + def countResult = + sql "SELECT COUNT(*) FROM test_cloud_single_rowset_compaction_mow_delete_bitmap" + assertEquals(mowCountBeforeCompaction, countResult) + assertEquals(expectedMowRows, readMowRows()) + } + + def mowRowsetVersions = [2, 3, 4] + def mowMaxSegmentSizeByVersion = [2: 32768, 3: 8192, 4: 2048] + def finalMowSegmentCounts = [:] + def compactMowRowsetTwice = { int targetVersion -> + logger.info("Compact MOW rowset version ${targetVersion} twice") + def beforeFirstCompaction = showTablet( + "test_cloud_single_rowset_compaction_mow_delete_bitmap", + mowBackend.Host, mowBackend.HttpPort, mowTabletId) + def inputRowset = rowsetByVersion(beforeFirstCompaction, targetVersion) + def inputInfo = parseRowsetInfo(inputRowset) + assertEquals("OVERLAPPING", inputInfo.overlap) + assertTrue(inputInfo.segments > initialInputSegmentsPerGroup, inputRowset) + def untouchedVersions = mowRowsetVersions.findAll { it != targetVersion } + def untouchedRowsets = untouchedVersions.collectEntries { int version -> + [(version): rowsetByVersion(beforeFirstCompaction, version)] + } + + set_be_param("vertical_compaction_max_segment_size", + mowMaxSegmentSizeByVersion[targetVersion].toString()) + set_be_param("cloud_single_rowset_compaction_segment_group_size", + initialInputSegmentsPerGroup.toString()) + runCumulativeCompaction(mowBackend.Host, mowBackend.HttpPort, mowTabletId) + + def afterFirstCompaction = showTablet( + "test_cloud_single_rowset_compaction_mow_delete_bitmap", + mowBackend.Host, mowBackend.HttpPort, mowTabletId) + def groupedRowset = rowsetByVersion(afterFirstCompaction, targetVersion) + def groupedInfo = parseRowsetInfo(groupedRowset) + assertEquals("NONOVERLAPPING_WITHIN_GROUP", groupedInfo.overlap) + assertNotEquals(inputRowset, groupedRowset) + untouchedRowsets.each { int version, String rowset -> + assertEquals(rowset, rowsetByVersion(afterFirstCompaction, version)) + } + checkMowQueryResult() + + set_be_param("cloud_single_rowset_compaction_segment_group_size", + groupedInfo.segments.toString()) + runCumulativeCompaction(mowBackend.Host, mowBackend.HttpPort, mowTabletId) + + def afterSecondCompaction = showTablet( + "test_cloud_single_rowset_compaction_mow_delete_bitmap", + mowBackend.Host, mowBackend.HttpPort, mowTabletId) + def nonoverlappingRowset = rowsetByVersion(afterSecondCompaction, targetVersion) + def nonoverlappingInfo = parseRowsetInfo(nonoverlappingRowset) + assertEquals("NONOVERLAPPING", nonoverlappingInfo.overlap) + assertNotEquals(groupedRowset, nonoverlappingRowset) + finalMowSegmentCounts[targetVersion] = nonoverlappingInfo.segments + untouchedRowsets.each { int version, String rowset -> + assertEquals(rowset, rowsetByVersion(afterSecondCompaction, version)) + } + checkMowQueryResult() + } + + set_be_param("enable_rowid_conversion_correctness_check", "true") + set_be_param("compaction_batch_size", "512") + mowRowsetVersions.each { int version -> compactMowRowsetTwice(version) } + assertEquals(mowRowsetVersions.size(), finalMowSegmentCounts.values().toSet().size(), + "expected different final segment counts: ${finalMowSegmentCounts}") + + // ==================================================== + logCaseSection("Multi-segment case: verify disjoint key ranges after the second " + + "cumulative compaction") + set_be_param("cloud_single_rowset_compaction_segment_group_size", + initialInputSegmentsPerGroup.toString()) + set_be_param("vertical_compaction_max_segment_size", "2048") + set_be_param("compaction_batch_size", "512") + checkSingleRowsetGroupedCompaction( + "test_cloud_single_rowset_grouped_compact_multi_seg_dup", + "DUPLICATE KEY", "v INT", "", initialInputSegmentsPerGroup, 32768, 32768 * 2, + [["100", "100"], ["100", "101"]], true) + + // ==================================================== + logCaseSection("Schema-change case: verify CREATE INDEX clears the grouped rowset " + + "layout") + sql "DROP TABLE IF EXISTS test_cloud_grouped_compaction_schema_change" + sql """ + CREATE TABLE test_cloud_grouped_compaction_schema_change ( + k INT, + v INT + ) + DUPLICATE KEY(k) + DISTRIBUTED BY HASH(k) BUCKETS 1 + PROPERTIES ( + "replication_num" = "1", + "disable_auto_compaction" = "true" + ) + """ + + StringBuilder schemaChangeContent = new StringBuilder() + for (int i = 0; i < 2; i++) { + (1..8192).each { + schemaChangeContent.append("${it},${it + i}\n") + } + } + streamLoad { + table "test_cloud_grouped_compaction_schema_change" + set "column_separator", "," + inputStream new ByteArrayInputStream(schemaChangeContent.toString().getBytes()) + time 30000 + check { result, exception, startTime, endTime -> + if (exception != null) { + throw exception + } + def json = parseJson(result) + assertEquals("success", json.Status.toLowerCase()) + assertEquals(8192 * 2, json.NumberTotalRows) + assertEquals(0, json.NumberFilteredRows) + } + } + sql "sync" + + def schemaChangeTablets = + sql_return_maparray "SHOW TABLETS FROM test_cloud_grouped_compaction_schema_change" + assertEquals(1, schemaChangeTablets.size()) + def oldTabletId = schemaChangeTablets[0].TabletId + def oldBackendId = schemaChangeTablets[0].BackendId + def schemaChangeBackends = sql_return_maparray "SHOW BACKENDS" + def oldBackend = schemaChangeBackends.find { it.BackendId == oldBackendId } + assertNotNull(oldBackend) + + set_be_param("cloud_single_rowset_compaction_segment_group_size", + initialInputSegmentsPerGroup.toString()) + runCumulativeCompaction(oldBackend.Host, oldBackend.HttpPort, oldTabletId) + + def groupedTablet = showTablet("test_cloud_grouped_compaction_schema_change", + oldBackend.Host, oldBackend.HttpPort, oldTabletId) + def groupedRowset = rowsetByVersion(groupedTablet, 2) + def groupedInfo = parseRowsetInfo(groupedRowset) + assertEquals("NONOVERLAPPING_WITHIN_GROUP", groupedInfo.overlap) + + sql """ + CREATE INDEX idx_v ON test_cloud_grouped_compaction_schema_change(v) + USING INVERTED + """ + waitForSchemaChangeDone { + sql """ + SHOW ALTER TABLE COLUMN + WHERE TableName='test_cloud_grouped_compaction_schema_change' + ORDER BY CreateTime DESC LIMIT 1 + """ + time 90 + } + + schemaChangeTablets = + sql_return_maparray "SHOW TABLETS FROM test_cloud_grouped_compaction_schema_change" + assertEquals(1, schemaChangeTablets.size()) + def newTabletId = schemaChangeTablets[0].TabletId + assertNotEquals(oldTabletId, newTabletId) + def newBackendId = schemaChangeTablets[0].BackendId + schemaChangeBackends = sql_return_maparray "SHOW BACKENDS" + def newBackend = schemaChangeBackends.find { it.BackendId == newBackendId } + assertNotNull(newBackend) + + def rewrittenTablet = showTablet("test_cloud_grouped_compaction_schema_change", + newBackend.Host, newBackend.HttpPort, newTabletId) + def rewrittenRowset = rowsetByVersion(rewrittenTablet, 2) + def rewrittenInfo = parseRowsetInfo(rewrittenRowset) + assertNotEquals("NONOVERLAPPING_WITHIN_GROUP", rewrittenInfo.overlap) + def rewrittenRowsetMeta = getRowsetMeta(newTabletId, 2) + assertTrue((rewrittenRowsetMeta.segment_group_sizes ?: []).isEmpty()) + def rewrittenCount = + sql "SELECT COUNT(*) FROM test_cloud_grouped_compaction_schema_change" + assertEquals(8192 * 2, rewrittenCount[0][0]) + def rewrittenPointRows = + readPointRows("test_cloud_grouped_compaction_schema_change") + assertEquals([["100", "100"], ["100", "101"]], rewrittenPointRows) + } finally { + sql "set enable_sql_cache=${sqlCacheOrigValue}" + GetDebugPoint().clearDebugPointsForAllBEs() + } + } +} From ccd0af214f93e785b3c836d54e12725e46f85f8e Mon Sep 17 00:00:00 2001 From: meiyi Date: Mon, 14 Sep 2026 11:22:06 +0800 Subject: [PATCH 2/3] [fix](compaction) Count grouped compaction output segments on 4.1 --- be/src/cloud/cloud_cumulative_compaction.cpp | 10 +++++++--- .../storage/cloud_file_cache_write_index_only_test.cpp | 6 ------ 2 files changed, 7 insertions(+), 9 deletions(-) diff --git a/be/src/cloud/cloud_cumulative_compaction.cpp b/be/src/cloud/cloud_cumulative_compaction.cpp index c1b0d467bd21d4..f4ce0f81a02f65 100644 --- a/be/src/cloud/cloud_cumulative_compaction.cpp +++ b/be/src/cloud/cloud_cumulative_compaction.cpp @@ -838,7 +838,9 @@ Status CloudCumulativeCompaction::do_merge_input_rowsets( *input_rowset->rowset_meta(), segment_group_size); for (size_t range_index = 0; range_index < segment_ranges.size(); ++range_index) { const auto& range = segment_ranges[range_index]; - const int32_t output_segment_start = _output_rs_writer->get_allocated_segment_id(); + std::vector output_segment_num_rows; + RETURN_IF_ERROR(_output_rs_writer->get_segment_num_rows(&output_segment_num_rows)); + const size_t output_segment_start = output_segment_num_rows.size(); RowsetReaderSharedPtr rs_reader; RETURN_IF_ERROR(input_rowset->create_reader(&rs_reader)); @@ -861,8 +863,10 @@ Status CloudCumulativeCompaction::do_merge_input_rowsets( _stats.cloud_local_read_time += group_stats.cloud_local_read_time; _stats.cloud_remote_read_time += group_stats.cloud_remote_read_time; - const int32_t output_segment_end = _output_rs_writer->get_allocated_segment_id(); - const int32_t output_group_size = output_segment_end - output_segment_start; + RETURN_IF_ERROR(_output_rs_writer->get_segment_num_rows(&output_segment_num_rows)); + DORIS_CHECK_GE(output_segment_num_rows.size(), output_segment_start); + const int32_t output_group_size = + cast_set(output_segment_num_rows.size() - output_segment_start); if (output_group_size > 0) { result->output_segment_group_sizes.push_back(output_group_size); } diff --git a/be/test/storage/cloud_file_cache_write_index_only_test.cpp b/be/test/storage/cloud_file_cache_write_index_only_test.cpp index ba937a64a83efd..b7baae652028f5 100644 --- a/be/test/storage/cloud_file_cache_write_index_only_test.cpp +++ b/be/test/storage/cloud_file_cache_write_index_only_test.cpp @@ -699,26 +699,20 @@ TEST_F(CloudFileCacheWriteIndexOnlyTest, auto writer_result = RowsetFactory::create_rowset_writer(*_engine, context, true); ASSERT_TRUE(writer_result.has_value()) << writer_result.error(); auto rowset_writer = std::move(writer_result).value(); - EXPECT_EQ(rowset_writer->get_allocated_segment_id(), 0); - std::vector key_column_ids = {0}; auto key_block = create_column_block(tablet_schema, key_column_ids, 8, 1); auto st = rowset_writer->add_columns(&key_block, key_column_ids, true, 4, false); ASSERT_TRUE(st.ok()) << st; - EXPECT_EQ(rowset_writer->get_allocated_segment_id(), 1); auto second_key_block = create_column_block(tablet_schema, key_column_ids, 8, 100); st = rowset_writer->add_columns(&second_key_block, key_column_ids, true, 4, false); ASSERT_TRUE(st.ok()) << st; - EXPECT_EQ(rowset_writer->get_allocated_segment_id(), 2); st = rowset_writer->flush_columns(true); ASSERT_TRUE(st.ok()) << st; - EXPECT_EQ(rowset_writer->get_allocated_segment_id(), 2); std::vector value_column_ids = {1}; auto value_block = create_column_block(tablet_schema, value_column_ids, 16, 1); st = rowset_writer->add_columns(&value_block, value_column_ids, false, UINT32_MAX, false); ASSERT_TRUE(st.ok()) << st; - EXPECT_EQ(rowset_writer->get_allocated_segment_id(), 2); st = rowset_writer->flush_columns(false); ASSERT_TRUE(st.ok()) << st; st = rowset_writer->final_flush(); From cdbcbf9312ecece57b58671d10bc7a6d2a16ecaa Mon Sep 17 00:00:00 2001 From: meiyi Date: Mon, 14 Sep 2026 15:21:33 +0800 Subject: [PATCH 3/3] [fix](compaction) Partial pick: Avoid repeatedly compacting large cumulative rowset Partial pick of a1e076e8d6969feafd8411f7b19e22664f43c543 (#64954). Port only the last-rowset trimming guard so a small overlapping singleton can reach the existing eligibility checks and grouped compaction selection. Keep the current promotion-size early return. Add coverage for a six-segment DUP rowset below the promotion size, including grouped and non-overlapping layouts. Validation: git diff --check. No compilation or test execution, as requested. --- .../cloud_cumulative_compaction_policy.cpp | 13 ++++++- ...loud_cumulative_compaction_policy_test.cpp | 38 +++++++++++++++++++ 2 files changed, 50 insertions(+), 1 deletion(-) diff --git a/be/src/cloud/cloud_cumulative_compaction_policy.cpp b/be/src/cloud/cloud_cumulative_compaction_policy.cpp index 0f7371307c2322..2be0dde84e3cbc 100644 --- a/be/src/cloud/cloud_cumulative_compaction_policy.cpp +++ b/be/src/cloud/cloud_cumulative_compaction_policy.cpp @@ -18,6 +18,7 @@ #include "cloud/cloud_cumulative_compaction_policy.h" #include +#include #include #include #include @@ -213,6 +214,9 @@ int64_t CloudSizeBasedCumulativeCompactionPolicy::pick_input_rowsets( auto rs_begin = input_rowsets->begin(); size_t new_compaction_score = *compaction_score; + const bool can_handle_exhausted_input = + (config::prioritize_query_perf_in_compaction && tablet->keys_type() != DUP_KEYS) || + *compaction_score >= static_cast(max_compaction_score); while (rs_begin != input_rowsets->end()) { auto& rs_meta = (*rs_begin)->rowset_meta(); int64_t current_level = _level_size(rs_meta->total_disk_size()); @@ -222,9 +226,16 @@ int64_t CloudSizeBasedCumulativeCompactionPolicy::pick_input_rowsets( if (current_level <= remain_level) { break; } + + auto next = std::next(rs_begin); + // Keep the last suffix rowset for the singleton checks unless the exhausted-input + // fallback below can select a useful input. + if (next == input_rowsets->end() && !can_handle_exhausted_input) { + break; + } total_size -= rs_meta->total_disk_size(); new_compaction_score -= rs_meta->get_compaction_score(); - ++rs_begin; + rs_begin = next; } if (rs_begin == input_rowsets->end()) { // No suitable level size found in `input_rowsets` if (config::prioritize_query_perf_in_compaction && tablet->keys_type() != DUP_KEYS) { diff --git a/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp b/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp index 79ec7f14d20fa9..ad651374fdac9b 100644 --- a/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp +++ b/be/test/cloud/cloud_cumulative_compaction_policy_test.cpp @@ -204,6 +204,44 @@ TEST_F(TestCloudSizeBasedCumulativeCompactionPolicy, last_delete_version, 2)); } +TEST_F(TestCloudSizeBasedCumulativeCompactionPolicy, + pick_input_rowsets_small_single_overlapping_rowset_not_trimmed_empty) { + TTabletSchema schema; + schema.keys_type = TKeysType::DUP_KEYS; + TabletMetaSharedPtr tablet_meta(new TabletMeta(1, 2, 15673, 15674, 4, 5, schema, 6, {{7, 8}}, + UniqueId(9, 10), TTabletType::TABLET_TYPE_DISK, + TCompressionType::LZ4F)); + tablet_meta->set_tablet_state(TABLET_RUNNING); + CloudTablet tablet(_engine, tablet_meta); + tablet._base_size = 1024L * 1024 * 1024; + CloudSizeBasedCumulativeCompactionPolicy policy; + + // Match the regression case: one small rowset, six segments, min score 2, max score 1000. + auto rowset = create_rowset(Version(2, 2), 6, true, 6370); + ASSERT_NE(nullptr, rowset); + for (auto overlap : {OVERLAPPING, NONOVERLAPPING_WITHIN_GROUP, NONOVERLAPPING}) { + SCOPED_TRACE(static_cast(overlap)); + rowset->rowset_meta()->set_segments_overlap(overlap); + rowset->rowset_meta()->clear_segment_group_sizes(); + if (overlap == NONOVERLAPPING_WITHIN_GROUP) { + rowset->rowset_meta()->set_segment_group_sizes({2, 2, 2}); + } + std::vector input_rowsets; + Version last_delete_version {-1, -1}; + size_t compaction_score = 0; + policy.pick_input_rowsets(&tablet, {rowset}, 1000, 2, &input_rowsets, + &last_delete_version, &compaction_score, true); + if (overlap == NONOVERLAPPING) { + EXPECT_TRUE(input_rowsets.empty()); + EXPECT_EQ(0, compaction_score); + } else { + ASSERT_EQ(1, input_rowsets.size()); + EXPECT_EQ(rowset, input_rowsets.front()); + EXPECT_EQ(6, compaction_score); + } + } +} + // Test case: Empty rowset compaction with skip_trim TEST_F(TestCloudSizeBasedCumulativeCompactionPolicy, pick_input_rowsets_empty_rowset_compaction) { // Save original config values