Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
209 changes: 198 additions & 11 deletions be/src/cloud/cloud_cumulative_compaction.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@

#include "cloud/cloud_cumulative_compaction.h"

#include <fmt/format.h>
#include <fmt/ranges.h>
#include <gen_cpp/cloud.pb.h>

#include <random>
Expand All @@ -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"
Expand All @@ -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<int64_t>(rowset_meta->segment_group_sizes().size())
: rowset->num_segments();
Comment on lines +54 to +57
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<RowsetSharedPtr>& 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<SegmentGroupMergeRange> 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<SegmentGroupMergeRange> 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<int64_t>(input_segment_group_sizes.size());
DORIS_CHECK_GT(input_group_count, 0);
ranges.reserve(cast_set<size_t>((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<int>(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<size_t>((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<uint64_t> 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");
Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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<CUMULATIVE_NO_SUITABLE_VERSION>(
"no suitable versions: input rowsets empty");
Expand Down Expand Up @@ -709,6 +813,89 @@ 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<RowsetReaderSharedPtr>& 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];
std::vector<uint32_t> 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));
std::vector<RowsetReaderSharedPtr> 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<int64_t>(segment_ranges.size()),
.range_index = cast_set<int64_t>(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;

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<int32_t>(output_segment_num_rows.size() - 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);
Expand Down
34 changes: 34 additions & 0 deletions be/src/cloud/cloud_cumulative_compaction.h
Original file line number Diff line number Diff line change
Expand Up @@ -20,15 +20,39 @@
#include <limits>
#include <memory>
#include <optional>
#include <string_view>
#include <vector>

#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<RowsetSharedPtr>& input_rowsets,
const TabletSchema& tablet_schema,
std::string_view compaction_policy);

std::vector<SegmentGroupMergeRange> 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);
Expand All @@ -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<RowsetReaderSharedPtr>& 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:
Expand All @@ -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<int64_t> _single_rowset_compaction_segment_group_size;
};

#include "common/compile_check_end.h"
Expand Down
13 changes: 12 additions & 1 deletion be/src/cloud/cloud_cumulative_compaction_policy.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
#include "cloud/cloud_cumulative_compaction_policy.h"

#include <algorithm>
#include <iterator>
#include <list>
#include <ostream>
#include <string>
Expand Down Expand Up @@ -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<size_t>(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());
Expand All @@ -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) {
Expand Down
8 changes: 7 additions & 1 deletion be/src/cloud/cloud_schema_change_job.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
2 changes: 2 additions & 0 deletions be/src/cloud/cloud_snapshot_mgr.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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()) {
Expand Down
Loading
Loading