Skip to content

Commit 3568feb

Browse files
committed
[feature](cloud) Support single rowset grouped compaction
Issue Number: None Related PR: None Problem Summary: Cloud cumulative compaction normally selects multiple rowsets, so a singleton rowset containing many overlapping segments can remain uncompacted. Add configurable selection and grouped merge support for eligible single rowsets, merging bounded segment ranges through the common compaction output lifecycle and preserving overlap semantics across non-empty output groups. Restrict the grouped path to size-based cumulative policy after normal policy filtering, exclude cluster-key tablets, retain the input cumulative point for this grouped path, and keep row ID conversion, vertical writer segment metadata, inverted index metadata, and vertical progress reporting consistent across all segment groups. Snapshot the grouped-mode decision and segment group size when input rowsets are selected so queued tasks preserve their selected execution semantics if dynamic settings change before merging begins. Add focused BE unit coverage and a cloud regression case for grouped compaction behavior. Support configurable cloud cumulative compaction of a single overlapping rowset in bounded segment groups. Queued tasks preserve their selected grouped execution mode when dynamic settings change. Cluster-key tablets continue to use the existing compaction path. - Test: Unit Test / Regression test - Added focused BE unit coverage and a cloud regression case. - Targeted CloudCompactionTest attempted with run-be-ut.sh; blocked during CMake configuration because OpenMP_C is unavailable. - Behavior changed: Yes. Eligible cloud single rowsets with overlapping segments can be compacted in bounded groups while retaining their previous cumulative point and task-local grouped settings. - Does this need documentation: No
1 parent f08b194 commit 3568feb

16 files changed

Lines changed: 1086 additions & 87 deletions

be/src/cloud/cloud_cumulative_compaction.cpp

Lines changed: 133 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -26,19 +26,43 @@
2626
#include "cloud/config.h"
2727
#include "common/config.h"
2828
#include "common/logging.h"
29+
#include "common/metrics/doris_metrics.h"
2930
#include "common/status.h"
3031
#include "cpp/sync_point.h"
3132
#include "service/backend_options.h"
3233
#include "storage/compaction/compaction.h"
3334
#include "storage/compaction/cumulative_compaction_policy.h"
3435
#include "storage/compaction/cumulative_compaction_time_series_policy.h"
36+
#include "storage/merger.h"
37+
#include "storage/rowset/rowset_reader.h"
38+
#include "storage/rowset/rowset_writer.h"
39+
#include "storage/tablet/tablet_schema.h"
3540
#include "util/debug_points.h"
3641
#include "util/trace.h"
3742
#include "util/uuid_generator.h"
3843

3944
namespace doris {
4045
using namespace ErrorCode;
4146

47+
namespace cloud {
48+
49+
bool 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;
53+
}
54+
55+
bool should_use_single_rowset_grouped_compaction(const std::vector<RowsetSharedPtr>& input_rowsets,
56+
const TabletSchema& tablet_schema,
57+
std::string_view compaction_policy) {
58+
return compaction_policy == CUMULATIVE_SIZE_BASED_POLICY &&
59+
tablet_schema.num_key_columns() > 0 && tablet_schema.cluster_key_uids().empty() &&
60+
config::enable_cloud_single_rowset_compaction && input_rowsets.size() == 1 &&
61+
is_single_rowset_compaction_candidate(input_rowsets.front());
62+
}
63+
64+
} // namespace cloud
65+
4266
bvar::Adder<uint64_t> cumu_output_size("cumu_compaction", "output_size");
4367
bvar::LatencyRecorder g_cu_compaction_hold_delete_bitmap_lock_time_ms(
4468
"cu_compaction_hold_delete_bitmap_lock_time_ms");
@@ -266,17 +290,19 @@ Status CloudCumulativeCompaction::modify_rowsets() {
266290
}
267291
auto compaction_policy = cloud_tablet()->tablet_meta()->compaction_policy();
268292
int64_t new_cumulative_point = input_cumulative_point;
269-
if (!_enable_parallel_cumu_compaction && input_tablet_state == TABLET_NOTREADY &&
270-
_output_rowset->start_version() > input_cumulative_point) {
271-
// Historical rowsets are absent from a schema-change target until conversion finishes.
272-
DORIS_CHECK_LE(input_cumulative_point, input_alter_version);
273-
DORIS_CHECK_GT(_output_rowset->start_version(), input_alter_version);
274-
} else if (!_enable_parallel_cumu_compaction ||
275-
_output_rowset->start_version() == input_cumulative_point) {
276-
new_cumulative_point =
277-
_engine.cumu_compaction_policy(compaction_policy)
278-
->new_cumulative_point(cloud_tablet(), _output_rowset, _last_delete_version,
279-
input_cumulative_point);
293+
if (!_single_rowset_compaction_segment_group_size.has_value()) {
294+
if (!_enable_parallel_cumu_compaction && input_tablet_state == TABLET_NOTREADY &&
295+
_output_rowset->start_version() > input_cumulative_point) {
296+
// Historical rowsets are absent from a schema-change target until conversion finishes.
297+
DORIS_CHECK_LE(input_cumulative_point, input_alter_version);
298+
DORIS_CHECK_GT(_output_rowset->start_version(), input_alter_version);
299+
} else if (!_enable_parallel_cumu_compaction ||
300+
_output_rowset->start_version() == input_cumulative_point) {
301+
new_cumulative_point =
302+
_engine.cumu_compaction_policy(compaction_policy)
303+
->new_cumulative_point(cloud_tablet(), _output_rowset,
304+
_last_delete_version, input_cumulative_point);
305+
}
280306
}
281307
// commit compaction job
282308
cloud::TabletJobInfoPB job;
@@ -607,6 +633,7 @@ Status CloudCumulativeCompaction::advance_cumulative_point_before_pick(
607633

608634
Status CloudCumulativeCompaction::pick_rowsets_to_compact() {
609635
_input_rowsets.clear();
636+
_single_rowset_compaction_segment_group_size.reset();
610637

611638
int64_t min_conflict_version = _min_conflict_version;
612639
int64_t max_conflict_version = _max_conflict_version;
@@ -645,6 +672,19 @@ Status CloudCumulativeCompaction::pick_rowsets_to_compact() {
645672
config::cumulative_compaction_min_deltas, &_input_rowsets,
646673
&_last_delete_version, &compaction_score);
647674

675+
if (config::enable_cloud_single_rowset_compaction) {
676+
for (const auto& rowset : _input_rowsets) {
677+
if (cloud::should_use_single_rowset_grouped_compaction(
678+
{rowset}, *cloud_tablet()->tablet_schema(), compaction_policy)) {
679+
auto grouped_input_rowset = rowset;
680+
_input_rowsets = {std::move(grouped_input_rowset)};
681+
_single_rowset_compaction_segment_group_size =
682+
config::cloud_single_rowset_compaction_segment_group_size;
683+
return Status::OK();
684+
}
685+
}
686+
}
687+
648688
if (_input_rowsets.empty()) {
649689
return Status::Error<CUMULATIVE_NO_SUITABLE_VERSION>(
650690
"no suitable versions: input rowsets empty");
@@ -708,6 +748,88 @@ Status CloudCumulativeCompaction::pick_rowsets_to_compact() {
708748
return Status::OK();
709749
}
710750

751+
Status CloudCumulativeCompaction::prepare_merge_input_rowsets(MergeInputRowsetsResult* result) {
752+
if (!_single_rowset_compaction_segment_group_size.has_value()) {
753+
return Status::OK();
754+
}
755+
756+
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+
}
762+
result->is_segment_grouped = true;
763+
result->segment_group_size = segment_group_size;
764+
return Status::OK();
765+
}
766+
767+
Status CloudCumulativeCompaction::do_merge_input_rowsets(
768+
const std::vector<RowsetReaderSharedPtr>& input_rs_readers,
769+
MergeInputRowsetsResult* result) {
770+
if (!result->is_segment_grouped) {
771+
return Compaction::do_merge_input_rowsets(input_rs_readers, result);
772+
}
773+
774+
const int64_t segment_group_size = result->segment_group_size;
775+
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());
782+
const int32_t output_segment_start = _output_rs_writer->get_allocated_segment_id();
783+
784+
RowsetReaderSharedPtr rs_reader;
785+
RETURN_IF_ERROR(input_rowset->create_reader(&rs_reader));
786+
std::vector<RowsetReaderSharedPtr> group_readers;
787+
group_readers.push_back(std::move(rs_reader));
788+
789+
Merger::Statistics group_stats;
790+
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}));
795+
796+
_stats.output_rows += group_stats.output_rows;
797+
_stats.merged_rows += group_stats.merged_rows;
798+
_stats.filtered_rows += group_stats.filtered_rows;
799+
_stats.bytes_read_from_local += group_stats.bytes_read_from_local;
800+
_stats.bytes_read_from_remote += group_stats.bytes_read_from_remote;
801+
_stats.cached_bytes_total += group_stats.cached_bytes_total;
802+
_stats.cloud_local_read_time += group_stats.cloud_local_read_time;
803+
_stats.cloud_remote_read_time += group_stats.cloud_remote_read_time;
804+
805+
const int32_t output_segment_end = _output_rs_writer->get_allocated_segment_id();
806+
const int32_t output_group_size = output_segment_end - output_segment_start;
807+
if (output_group_size > 0) {
808+
++result->output_segment_group_count;
809+
}
810+
}
811+
return Status::OK();
812+
}
813+
814+
void CloudCumulativeCompaction::update_output_rowset_after_build(
815+
const MergeInputRowsetsResult& result) {
816+
if (!result.is_segment_grouped) {
817+
return;
818+
}
819+
if (result.output_segment_group_count > 1) {
820+
_output_rowset->rowset_meta()->set_segments_overlap(OVERLAPPING);
821+
}
822+
823+
const auto& input_rowset = _input_rowsets.front();
824+
LOG_INFO("finish single rowset grouped compaction, tablet_id={}, version=[{}-{}]",
825+
_tablet->tablet_id(), input_rowset->start_version(), input_rowset->end_version())
826+
.tag("job_id", _uuid)
827+
.tag("input_segments", input_rowset->num_segments())
828+
.tag("segment_group_size", result.segment_group_size)
829+
.tag("output_segments", _output_rowset->num_segments())
830+
.tag("output_groups", result.output_segment_group_count);
831+
}
832+
711833
void CloudCumulativeCompaction::update_cumulative_point(int64_t input_cumulative_point,
712834
int64_t output_cumulative_point) {
713835
DORIS_CHECK_LT(input_cumulative_point, output_cumulative_point);

be/src/cloud/cloud_cumulative_compaction.h

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,14 +20,27 @@
2020
#include <limits>
2121
#include <memory>
2222
#include <optional>
23+
#include <string_view>
24+
#include <vector>
2325

2426
#include "cloud/cloud_storage_engine.h"
2527
#include "cloud/cloud_tablet.h"
2628
#include "storage/compaction/compaction.h"
2729
#include "storage/compaction_task_tracker.h"
30+
#include "storage/tablet/tablet_fwd.h"
2831

2932
namespace doris {
3033

34+
namespace cloud {
35+
36+
bool is_single_rowset_compaction_candidate(const RowsetSharedPtr& rowset);
37+
38+
bool should_use_single_rowset_grouped_compaction(const std::vector<RowsetSharedPtr>& input_rowsets,
39+
const TabletSchema& tablet_schema,
40+
std::string_view compaction_policy);
41+
42+
} // namespace cloud
43+
3144
class CloudCumulativeCompaction : public CloudCompactionMixin {
3245
public:
3346
CloudCumulativeCompaction(CloudStorageEngine& engine, CloudTabletSPtr tablet);
@@ -53,6 +66,13 @@ class CloudCumulativeCompaction : public CloudCompactionMixin {
5366

5467
Status pick_rowsets_to_compact();
5568

69+
Status prepare_merge_input_rowsets(MergeInputRowsetsResult* result) override;
70+
71+
Status do_merge_input_rowsets(const std::vector<RowsetReaderSharedPtr>& input_rs_readers,
72+
MergeInputRowsetsResult* result) override;
73+
74+
void update_output_rowset_after_build(const MergeInputRowsetsResult& result) override;
75+
5676
std::string_view compaction_name() const override { return "CloudCumulativeCompaction"; }
5777

5878
protected:
@@ -75,6 +95,7 @@ class CloudCumulativeCompaction : public CloudCompactionMixin {
7595
int64_t _cumulative_compaction_cnt = 0;
7696
int64_t _picked_cumulative_point = 0;
7797
Version _last_delete_version {-1, -1};
98+
std::optional<int64_t> _single_rowset_compaction_segment_group_size;
7899
};
79100

80101
} // namespace doris

be/src/cloud/config.cpp

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,9 @@ DEFINE_mInt32(max_base_compaction_task_num_per_disk, "2");
5555
DEFINE_mBool(prioritize_query_perf_in_compaction, "false");
5656
DEFINE_mInt32(compaction_max_rowset_count, "10000");
5757
DEFINE_mInt64(compaction_txn_max_size_bytes, "7340032"); // 7MB
58+
DEFINE_mBool(enable_cloud_single_rowset_compaction, "false");
59+
DEFINE_mInt32(cloud_single_rowset_compaction_min_segments, "512");
60+
DEFINE_mInt32(cloud_single_rowset_compaction_segment_group_size, "64");
5861

5962
DEFINE_mInt32(refresh_s3_info_interval_s, "60");
6063
DEFINE_mInt32(vacuum_stale_rowsets_interval_s, "300");

be/src/cloud/config.h

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,9 @@ DECLARE_mInt32(max_base_compaction_task_num_per_disk);
9696
DECLARE_mBool(prioritize_query_perf_in_compaction);
9797
DECLARE_mInt32(compaction_max_rowset_count);
9898
DECLARE_mInt64(compaction_txn_max_size_bytes);
99+
DECLARE_mBool(enable_cloud_single_rowset_compaction);
100+
DECLARE_mInt32(cloud_single_rowset_compaction_min_segments);
101+
DECLARE_mInt32(cloud_single_rowset_compaction_segment_group_size);
99102

100103
// CloudStorageEngine config
101104
DECLARE_mInt32(refresh_s3_info_interval_s);

be/src/storage/compaction/compaction.cpp

Lines changed: 45 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -282,6 +282,9 @@ int64_t Compaction::merge_way_num() {
282282
}
283283

284284
Status Compaction::merge_input_rowsets() {
285+
MergeInputRowsetsResult result;
286+
RETURN_IF_ERROR(prepare_merge_input_rowsets(&result));
287+
285288
std::vector<RowsetReaderSharedPtr> input_rs_readers;
286289
input_rs_readers.reserve(_input_rowsets.size());
287290
for (auto& rowset : _input_rowsets) {
@@ -307,38 +310,10 @@ Status Compaction::merge_input_rowsets() {
307310
_stats.rowid_conversion = _rowid_conversion.get();
308311
}
309312

310-
int64_t way_num = merge_way_num();
311-
312-
Status res;
313313
{
314314
SCOPED_TIMER(_merge_rowsets_latency_timer);
315315
// 1. Merge segment files and write bkd inverted index
316-
// TODO implement vertical compaction for seq map
317-
if (_is_vertical && !_tablet->tablet_schema()->has_seq_map()) {
318-
if (!_tablet->tablet_schema()->cluster_key_uids().empty()) {
319-
RETURN_IF_ERROR(update_delete_bitmap());
320-
}
321-
auto progress_cb = [compaction_id = this->_compaction_id](int64_t total,
322-
int64_t completed) {
323-
CompactionTaskTracker::instance()->update_progress(compaction_id, total, completed);
324-
};
325-
res = Merger::vertical_merge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema,
326-
input_rs_readers, _output_rs_writer.get(),
327-
cast_set<uint32_t>(get_avg_segment_rows()),
328-
way_num, &_stats, progress_cb);
329-
} else {
330-
if (!_tablet->tablet_schema()->cluster_key_uids().empty()) {
331-
return Status::InternalError(
332-
"mow table with cluster keys does not support non vertical compaction");
333-
}
334-
res = Merger::vmerge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema,
335-
input_rs_readers, _output_rs_writer.get(), &_stats);
336-
}
337-
338-
_tablet->last_compaction_status = res;
339-
if (!res.ok()) {
340-
return res;
341-
}
316+
RETURN_IF_ERROR(do_merge_input_rowsets(input_rs_readers, &result));
342317
// 2. Merge the remaining inverted index files of the string type
343318
RETURN_IF_ERROR(do_inverted_index_compaction());
344319
}
@@ -364,6 +339,7 @@ Status Compaction::merge_input_rowsets() {
364339

365340
//RETURN_IF_ERROR(_engine.meta_mgr().commit_rowset(*_output_rowset->rowset_meta().get()));
366341
set_delete_predicate_for_output_rowset();
342+
update_output_rowset_after_build(result);
367343

368344
_local_read_bytes_total = _stats.bytes_read_from_local;
369345
_remote_read_bytes_total = _stats.bytes_read_from_remote;
@@ -380,6 +356,46 @@ Status Compaction::merge_input_rowsets() {
380356
return check_correctness();
381357
}
382358

359+
Status Compaction::do_merge_input_rowsets(
360+
const std::vector<RowsetReaderSharedPtr>& input_rs_readers,
361+
MergeInputRowsetsResult* /*result*/) {
362+
return execute_merge(input_rs_readers, merge_way_num(), &_stats);
363+
}
364+
365+
Status Compaction::execute_merge(const std::vector<RowsetReaderSharedPtr>& input_rs_readers,
366+
int64_t merge_way_num, Merger::Statistics* stats,
367+
std::optional<std::pair<int64_t, int64_t>> segment_range,
368+
VerticalMergeProgressContext progress) {
369+
Status status;
370+
// TODO implement vertical compaction for seq map
371+
if (_is_vertical && !_tablet->tablet_schema()->has_seq_map()) {
372+
if (!_tablet->tablet_schema()->cluster_key_uids().empty() && !segment_range.has_value()) {
373+
RETURN_IF_ERROR(update_delete_bitmap());
374+
}
375+
auto progress_cb = [compaction_id = this->_compaction_id, progress](int64_t total,
376+
int64_t completed) {
377+
CompactionTaskTracker::instance()->update_progress(
378+
compaction_id, total * progress.total_ranges,
379+
total * progress.range_index + completed);
380+
};
381+
status = Merger::vertical_merge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema,
382+
input_rs_readers, _output_rs_writer.get(),
383+
cast_set<uint32_t>(get_avg_segment_rows()),
384+
merge_way_num, stats, progress_cb, segment_range);
385+
} else {
386+
if (!_tablet->tablet_schema()->cluster_key_uids().empty()) {
387+
return Status::InternalError(
388+
"mow table with cluster keys does not support non vertical compaction");
389+
}
390+
status = Merger::vmerge_rowsets(_tablet, compaction_type(), *_cur_tablet_schema,
391+
input_rs_readers, _output_rs_writer.get(), stats,
392+
segment_range);
393+
}
394+
395+
_tablet->last_compaction_status = status;
396+
return status;
397+
}
398+
383399
void Compaction::set_delete_predicate_for_output_rowset() {
384400
// Now we support delete in cumu compaction, to make all data in rowsets whose version
385401
// is below output_version to be delete in the future base compaction, we should carry

0 commit comments

Comments
 (0)