This is an automated email from the ASF dual-hosted git repository.
mymeiyi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new edb7b44dbef [improvement](compaction) extract functions of compaction
execution (#67466)
edb7b44dbef is described below
commit edb7b44dbefef7cfd46664829082e044e35df3cd
Author: meiyi <[email protected]>
AuthorDate: Wed Sep 16 14:46:08 2026 +0800
[improvement](compaction) extract functions of compaction execution (#67466)
Later, we will support parallel compaction. This PR extracts the cloud
compaction execution into reusable functions.
---
be/src/storage/compaction/compaction.cpp | 44 ++++++++++++++++++++++----------
be/src/storage/compaction/compaction.h | 19 ++++++++++++--
2 files changed, 48 insertions(+), 15 deletions(-)
diff --git a/be/src/storage/compaction/compaction.cpp
b/be/src/storage/compaction/compaction.cpp
index 015c6ee497a..2c91d5ee3a3 100644
--- a/be/src/storage/compaction/compaction.cpp
+++ b/be/src/storage/compaction/compaction.cpp
@@ -282,21 +282,26 @@ int64_t Compaction::merge_way_num() {
}
Status Compaction::merge_input_rowsets() {
- MergeInputRowsetsResult result;
- RETURN_IF_ERROR(prepare_merge_input_rowsets(&result));
+ MergeInputRowsetsContext context;
+ RETURN_IF_ERROR(prepare_merge_input_rowsets_execution(&context));
+ RETURN_IF_ERROR(execute_merge_input_rowsets(&context));
+ return finish_merge_input_rowsets_execution(&context);
+}
+
+Status
Compaction::prepare_merge_input_rowsets_execution(MergeInputRowsetsContext*
context) {
+ RETURN_IF_ERROR(prepare_merge_input_rowsets(&context->result));
- std::vector<RowsetReaderSharedPtr> input_rs_readers;
- input_rs_readers.reserve(_input_rowsets.size());
+ context->input_rs_readers.reserve(_input_rowsets.size());
for (auto& rowset : _input_rowsets) {
RowsetReaderSharedPtr rs_reader;
RETURN_IF_ERROR(rowset->create_reader(&rs_reader));
- input_rs_readers.push_back(std::move(rs_reader));
+ context->input_rs_readers.push_back(std::move(rs_reader));
}
RowsetWriterContext ctx;
// Propagate input rowset readers into the rowset writer context before
the writer is created.
// Variant nested-group compaction uses this metadata to enable the
streaming writer path.
- ctx.input_rs_readers = input_rs_readers;
+ ctx.input_rs_readers = context->input_rs_readers;
RETURN_IF_ERROR(construct_output_rowset_writer(ctx));
// write merged rows to output rowset
@@ -309,15 +314,22 @@ Status Compaction::merge_input_rowsets() {
_tablet->enable_unique_key_merge_on_write())) {
_stats.rowid_conversion = _rowid_conversion.get();
}
+ return Status::OK();
+}
+Status Compaction::execute_merge_input_rowsets(MergeInputRowsetsContext*
context) {
{
SCOPED_TIMER(_merge_rowsets_latency_timer);
// 1. Merge segment files and write bkd inverted index
- RETURN_IF_ERROR(do_merge_input_rowsets(input_rs_readers, &result));
+ RETURN_IF_ERROR(do_merge_input_rowsets(context->input_rs_readers,
&context->result));
// 2. Merge the remaining inverted index files of the string type
RETURN_IF_ERROR(do_inverted_index_compaction());
}
+ return Status::OK();
+}
+Status
Compaction::finish_merge_input_rowsets_execution(MergeInputRowsetsContext*
context) {
+ auto& result = context->result;
COUNTER_UPDATE(_merged_rows_counter, _stats.merged_rows);
COUNTER_UPDATE(_filtered_rows_counter, _stats.filtered_rows);
@@ -2123,16 +2135,15 @@ bool
CloudCompactionMixin::should_apply_cumulative_compaction_result(
return true;
}
-Status CloudCompactionMixin::execute_compact_impl(int64_t permits) {
- OlapStopWatch watch;
-
+Status CloudCompactionMixin::prepare_execute_compact(int64_t permits) {
RETURN_IF_ERROR(build_basic_info());
LOG(INFO) << "start " << compaction_name() << ". tablet=" <<
_tablet->tablet_id()
<< ", output_version=" << _output_version << ", permits: " <<
permits;
+ return Status::OK();
+}
- RETURN_IF_ERROR(merge_input_rowsets());
-
+Status CloudCompactionMixin::finish_execute_compact(int64_t
execution_start_time_us) {
DBUG_EXECUTE_IF("CloudFullCompaction::modify_rowsets.wrong_rowset_id", {
DCHECK(compaction_type() == ReaderType::READER_FULL_COMPACTION);
RowsetId id;
@@ -2158,11 +2169,18 @@ Status
CloudCompactionMixin::execute_compact_impl(int64_t permits) {
auto tablet = std::static_pointer_cast<CloudTablet>(_tablet);
tablet->local_read_time_us.fetch_add(_stats.cloud_local_read_time);
tablet->remote_read_time_us.fetch_add(_stats.cloud_remote_read_time);
- tablet->exec_compaction_time_us.fetch_add(watch.get_elapse_time_us());
+ tablet->exec_compaction_time_us.fetch_add(MonotonicMicros() -
execution_start_time_us);
return Status::OK();
}
+Status CloudCompactionMixin::execute_compact_impl(int64_t permits) {
+ const int64_t execution_start_time_us = MonotonicMicros();
+ RETURN_IF_ERROR(prepare_execute_compact(permits));
+ RETURN_IF_ERROR(merge_input_rowsets());
+ return finish_execute_compact(execution_start_time_us);
+}
+
int64_t CloudCompactionMixin::initiator() const {
return HashUtil::hash64(_uuid.data(), _uuid.size(), 0) &
std::numeric_limits<int64_t>::max();
}
diff --git a/be/src/storage/compaction/compaction.h
b/be/src/storage/compaction/compaction.h
index e2b8a9b68b5..ebf361f0b61 100644
--- a/be/src/storage/compaction/compaction.h
+++ b/be/src/storage/compaction/compaction.h
@@ -109,8 +109,19 @@ protected:
std::vector<int32_t> output_segment_group_sizes;
};
+ struct MergeInputRowsetsContext {
+ MergeInputRowsetsResult result;
+ std::vector<RowsetReaderSharedPtr> input_rs_readers;
+ };
+
Status merge_input_rowsets();
+ Status prepare_merge_input_rowsets_execution(MergeInputRowsetsContext*
context);
+
+ Status execute_merge_input_rowsets(MergeInputRowsetsContext* context);
+
+ Status finish_merge_input_rowsets_execution(MergeInputRowsetsContext*
context);
+
virtual Status prepare_merge_input_rowsets(MergeInputRowsetsResult*
/*result*/) {
return Status::OK();
}
@@ -292,6 +303,12 @@ protected:
// Caller must hold the tablet header lock.
bool should_apply_cumulative_compaction_result(int64_t
response_cumulative_compaction_cnt);
+ Status prepare_execute_compact(int64_t permits);
+
+ Status finish_execute_compact(int64_t execution_start_time_us);
+
+ int64_t get_compaction_permits();
+
CloudStorageEngine& _engine;
std::string _uuid;
@@ -309,8 +326,6 @@ private:
virtual Status modify_rowsets();
- int64_t get_compaction_permits();
-
void update_compaction_level();
bool should_cache_compaction_output();
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]