1717
1818#include " cloud/cloud_cumulative_compaction.h"
1919
20+ #include < fmt/format.h>
21+ #include < fmt/ranges.h>
2022#include < gen_cpp/cloud.pb.h>
2123
2224#include < random>
@@ -47,9 +49,13 @@ using namespace ErrorCode;
4749namespace cloud {
4850
4951bool is_single_rowset_compaction_candidate (const RowsetSharedPtr& rowset) {
50- return !rowset->rowset_meta ()->has_delete_predicate () &&
51- rowset->rowset_meta ()->is_segments_overlapping () &&
52- rowset->num_segments () >= config::cloud_single_rowset_compaction_min_segments;
52+ const auto & rowset_meta = rowset->rowset_meta ();
53+ const int64_t overlap_unit_count =
54+ rowset_meta->segments_overlap () == NONOVERLAPPING_WITHIN_GROUP
55+ ? static_cast <int64_t >(rowset_meta->segment_group_sizes ().size ())
56+ : rowset->num_segments ();
57+ return !rowset_meta->has_delete_predicate () && rowset_meta->is_segments_overlapping () &&
58+ overlap_unit_count >= config::cloud_single_rowset_compaction_min_segments;
5359}
5460
5561bool should_use_single_rowset_grouped_compaction (const std::vector<RowsetSharedPtr>& input_rowsets,
@@ -61,6 +67,51 @@ bool should_use_single_rowset_grouped_compaction(const std::vector<RowsetSharedP
6167 is_single_rowset_compaction_candidate (input_rowsets.front ());
6268}
6369
70+ std::vector<SegmentGroupMergeRange> build_segment_group_merge_ranges (const RowsetMeta& rowset_meta,
71+ int64_t segment_group_size) {
72+ DORIS_CHECK_GT (segment_group_size, 1 );
73+ DORIS_CHECK_GT (rowset_meta.num_segments (), 0 );
74+
75+ std::vector<SegmentGroupMergeRange> ranges;
76+ if (rowset_meta.segments_overlap () == NONOVERLAPPING_WITHIN_GROUP ) {
77+ const auto & input_segment_group_sizes = rowset_meta.segment_group_sizes ();
78+ const int64_t input_group_count = cast_set<int64_t >(input_segment_group_sizes.size ());
79+ DORIS_CHECK_GT (input_group_count, 0 );
80+ ranges.reserve (cast_set<size_t >((input_group_count + segment_group_size - 1 ) /
81+ segment_group_size));
82+
83+ int64_t segment_end = 0 ;
84+ for (int64_t group_start = 0 ; group_start < input_group_count;
85+ group_start += segment_group_size) {
86+ const int64_t group_end = std::min (group_start + segment_group_size, input_group_count);
87+ const int64_t segment_start = segment_end;
88+ for (int64_t group_index = group_start; group_index < group_end; ++group_index) {
89+ const int32_t input_group_size =
90+ input_segment_group_sizes.Get (cast_set<int >(group_index));
91+ DORIS_CHECK_GT (input_group_size, 0 );
92+ segment_end += input_group_size;
93+ }
94+
95+ ranges.push_back ({.segment_start = segment_start,
96+ .segment_end = segment_end,
97+ .merge_way_num = group_end - group_start});
98+ }
99+ DORIS_CHECK_EQ (segment_end, rowset_meta.num_segments ());
100+ } else {
101+ ranges.reserve (cast_set<size_t >((rowset_meta.num_segments () + segment_group_size - 1 ) /
102+ segment_group_size));
103+ for (int64_t segment_start = 0 ; segment_start < rowset_meta.num_segments ();
104+ segment_start += segment_group_size) {
105+ const int64_t segment_end =
106+ std::min (segment_start + segment_group_size, rowset_meta.num_segments ());
107+ ranges.push_back ({.segment_start = segment_start,
108+ .segment_end = segment_end,
109+ .merge_way_num = segment_end - segment_start});
110+ }
111+ }
112+ return ranges;
113+ }
114+
64115} // namespace cloud
65116
66117bvar::Adder<uint64_t > cumu_output_size (" cumu_compaction" , " output_size" );
@@ -273,6 +324,18 @@ Status CloudCumulativeCompaction::execute_compact() {
273324 return st;
274325}
275326
327+ bool CloudCumulativeCompaction::should_calculate_new_cumulative_point (
328+ int64_t input_cumulative_point) const {
329+ if (!_single_rowset_compaction_segment_group_size.has_value ()) {
330+ return true ;
331+ }
332+
333+ DORIS_CHECK_EQ (_input_rowsets.size (), 1 );
334+ DORIS_CHECK (_output_rowset != nullptr );
335+ return _input_rowsets.front ()->start_version () == input_cumulative_point &&
336+ _output_rowset->rowset_meta ()->segments_overlap () == NONOVERLAPPING ;
337+ }
338+
276339Status CloudCumulativeCompaction::modify_rowsets () {
277340 // calculate new cumulative point
278341 int64_t input_cumulative_point;
@@ -290,7 +353,7 @@ Status CloudCumulativeCompaction::modify_rowsets() {
290353 }
291354 auto compaction_policy = cloud_tablet ()->tablet_meta ()->compaction_policy ();
292355 int64_t new_cumulative_point = input_cumulative_point;
293- if (!_single_rowset_compaction_segment_group_size. has_value ( )) {
356+ if (should_calculate_new_cumulative_point (input_cumulative_point )) {
294357 if (!_enable_parallel_cumu_compaction && input_tablet_state == TABLET_NOTREADY &&
295358 _output_rowset->start_version () > input_cumulative_point) {
296359 // Historical rowsets are absent from a schema-change target until conversion finishes.
@@ -672,14 +735,15 @@ Status CloudCumulativeCompaction::pick_rowsets_to_compact() {
672735 config::cumulative_compaction_min_deltas, &_input_rowsets,
673736 &_last_delete_version, &compaction_score);
674737
675- if (config::enable_cloud_single_rowset_compaction) {
738+ const int64_t segment_group_size =
739+ config::cloud_single_rowset_compaction_segment_group_size;
740+ if (config::enable_cloud_single_rowset_compaction && segment_group_size > 1 ) {
676741 for (const auto & rowset : _input_rowsets) {
677742 if (cloud::should_use_single_rowset_grouped_compaction (
678743 {rowset}, *cloud_tablet ()->tablet_schema (), compaction_policy)) {
679744 auto grouped_input_rowset = rowset;
680745 _input_rowsets = {std::move (grouped_input_rowset)};
681- _single_rowset_compaction_segment_group_size =
682- config::cloud_single_rowset_compaction_segment_group_size;
746+ _single_rowset_compaction_segment_group_size = segment_group_size;
683747 return Status::OK ();
684748 }
685749 }
@@ -754,11 +818,7 @@ Status CloudCumulativeCompaction::prepare_merge_input_rowsets(MergeInputRowsetsR
754818 }
755819
756820 const int64_t segment_group_size = *_single_rowset_compaction_segment_group_size;
757- if (segment_group_size <= 0 ) {
758- return Status::InvalidArgument (
759- " cloud_single_rowset_compaction_segment_group_size must be positive, value={}" ,
760- segment_group_size);
761- }
821+ DORIS_CHECK_GT (segment_group_size, 1 );
762822 result->is_segment_grouped = true ;
763823 result->segment_group_size = segment_group_size;
764824 return Status::OK ();
@@ -773,12 +833,10 @@ Status CloudCumulativeCompaction::do_merge_input_rowsets(
773833
774834 const int64_t segment_group_size = result->segment_group_size ;
775835 const auto & input_rowset = _input_rowsets.front ();
776- const int64_t segment_group_count =
777- (input_rowset->num_segments () + segment_group_size - 1 ) / segment_group_size;
778- for (int64_t segment_start = 0 ; segment_start < input_rowset->num_segments ();
779- segment_start += segment_group_size) {
780- const int64_t segment_end =
781- std::min (segment_start + segment_group_size, input_rowset->num_segments ());
836+ const auto segment_ranges = cloud::build_segment_group_merge_ranges (
837+ *input_rowset->rowset_meta (), segment_group_size);
838+ for (size_t range_index = 0 ; range_index < segment_ranges.size (); ++range_index) {
839+ const auto & range = segment_ranges[range_index];
782840 const int32_t output_segment_start = _output_rs_writer->get_allocated_segment_id ();
783841
784842 RowsetReaderSharedPtr rs_reader;
@@ -788,10 +846,10 @@ Status CloudCumulativeCompaction::do_merge_input_rowsets(
788846
789847 Merger::Statistics group_stats;
790848 group_stats.rowid_conversion = _stats.rowid_conversion ;
791- RETURN_IF_ERROR (execute_merge (group_readers, segment_end - segment_start , &group_stats,
792- std::make_pair (segment_start, segment_end),
793- {.total_ranges = segment_group_count ,
794- .range_index = segment_start / segment_group_size }));
849+ RETURN_IF_ERROR (execute_merge (group_readers, range. merge_way_num , &group_stats,
850+ std::make_pair (range. segment_start , range. segment_end ),
851+ {.total_ranges = cast_set< int64_t >(segment_ranges. size ()) ,
852+ .range_index = cast_set< int64_t >(range_index) }));
795853
796854 _stats.output_rows += group_stats.output_rows ;
797855 _stats.merged_rows += group_stats.merged_rows ;
@@ -805,7 +863,7 @@ Status CloudCumulativeCompaction::do_merge_input_rowsets(
805863 const int32_t output_segment_end = _output_rs_writer->get_allocated_segment_id ();
806864 const int32_t output_group_size = output_segment_end - output_segment_start;
807865 if (output_group_size > 0 ) {
808- ++ result->output_segment_group_count ;
866+ result->output_segment_group_sizes . push_back (output_group_size) ;
809867 }
810868 }
811869 return Status::OK ();
@@ -816,8 +874,9 @@ void CloudCumulativeCompaction::update_output_rowset_after_build(
816874 if (!result.is_segment_grouped ) {
817875 return ;
818876 }
819- if (result.output_segment_group_count > 1 ) {
820- _output_rowset->rowset_meta ()->set_segments_overlap (OVERLAPPING );
877+ if (result.output_segment_group_sizes .size () > 1 ) {
878+ _output_rowset->rowset_meta ()->set_segments_overlap (NONOVERLAPPING_WITHIN_GROUP );
879+ _output_rowset->rowset_meta ()->set_segment_group_sizes (result.output_segment_group_sizes );
821880 }
822881
823882 const auto & input_rowset = _input_rowsets.front ();
@@ -827,7 +886,9 @@ void CloudCumulativeCompaction::update_output_rowset_after_build(
827886 .tag (" input_segments" , input_rowset->num_segments ())
828887 .tag (" segment_group_size" , result.segment_group_size )
829888 .tag (" output_segments" , _output_rowset->num_segments ())
830- .tag (" output_groups" , result.output_segment_group_count );
889+ .tag (" output_groups" , result.output_segment_group_sizes .size ())
890+ .tag (" output_segment_group_sizes" ,
891+ fmt::format (" [{}]" , fmt::join (result.output_segment_group_sizes , " , " )));
831892}
832893
833894void CloudCumulativeCompaction::update_cumulative_point (int64_t input_cumulative_point,
0 commit comments