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]

Reply via email to