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 59cb8584c16 [improvement](compaction) extract functions of cloud cu
and base compaction record result (#67467)
59cb8584c16 is described below
commit 59cb8584c166f579cf1943384b5f8e1d698dcf7b
Author: meiyi <[email protected]>
AuthorDate: Wed Sep 16 14:41:39 2026 +0800
[improvement](compaction) extract functions of cloud cu and base compaction
record result (#67467)
Later, we will support parallel compaction. This PR extracts the cloud
cloud cu and base compaction record result into reusable functions.
---
be/src/cloud/cloud_base_compaction.cpp | 44 +++++++++++++++-------------
be/src/cloud/cloud_base_compaction.h | 4 +++
be/src/cloud/cloud_cumulative_compaction.cpp | 44 +++++++++++++++-------------
be/src/cloud/cloud_cumulative_compaction.h | 4 +++
4 files changed, 56 insertions(+), 40 deletions(-)
diff --git a/be/src/cloud/cloud_base_compaction.cpp
b/be/src/cloud/cloud_base_compaction.cpp
index 3570abe492f..8ea6c377a5a 100644
--- a/be/src/cloud/cloud_base_compaction.cpp
+++ b/be/src/cloud/cloud_base_compaction.cpp
@@ -286,26 +286,22 @@ Status CloudBaseCompaction::execute_compact() {
SCOPED_ATTACH_TASK(_mem_tracker);
- using namespace std::chrono;
- auto start = steady_clock::now();
- Status st;
- Defer defer_set_st([&] {
- cloud_tablet()->set_last_base_compaction_status(st.to_string());
- if (!st.ok()) {
-
cloud_tablet()->set_last_base_compaction_failure_time(UnixMillis());
- } else {
-
cloud_tablet()->set_last_base_compaction_success_time(UnixMillis());
- }
- });
- st = CloudCompactionMixin::execute_compact();
- if (!st.ok()) {
- LOG(WARNING) << "fail to do " << compaction_name() << ". res=" << st
- << ", tablet=" << _tablet->tablet_id()
- << ", output_version=" << _output_version;
- return st;
+ const auto execution_start_time = std::chrono::steady_clock::now();
+ Status status = CloudCompactionMixin::execute_compact();
+ if (!status.ok()) {
+ return record_compaction_failure(status);
}
+ record_compaction_success(execution_start_time);
+ return Status::OK();
+}
+
+void CloudBaseCompaction::record_compaction_success(
+ std::chrono::steady_clock::time_point execution_start_time) {
LOG_INFO("finish CloudBaseCompaction, tablet_id={}, cost={}ms
range=[{}-{}]",
- _tablet->tablet_id(),
duration_cast<milliseconds>(steady_clock::now() - start).count(),
+ _tablet->tablet_id(),
+ std::chrono::duration_cast<std::chrono::milliseconds>(
+ std::chrono::steady_clock::now() - execution_start_time)
+ .count(),
_input_rowsets.front()->start_version(),
_input_rowsets.back()->end_version())
.tag("job_id", _uuid)
.tag("input_rowsets", _input_rowsets.size())
@@ -331,8 +327,16 @@ Status CloudBaseCompaction::execute_compact() {
DorisMetrics::instance()->base_compaction_bytes_total->increment(_input_rowsets_total_size);
base_output_size << _output_rowset->total_disk_size();
- st = Status::OK();
- return st;
+ cloud_tablet()->set_last_base_compaction_status(Status::OK().to_string());
+ cloud_tablet()->set_last_base_compaction_success_time(UnixMillis());
+}
+
+Status CloudBaseCompaction::record_compaction_failure(Status status) {
+ LOG(WARNING) << "fail to do " << compaction_name() << ". res=" << status
+ << ", tablet=" << _tablet->tablet_id() << ", output_version="
<< _output_version;
+ cloud_tablet()->set_last_base_compaction_status(status.to_string());
+ cloud_tablet()->set_last_base_compaction_failure_time(UnixMillis());
+ return status;
}
Status CloudBaseCompaction::modify_rowsets() {
diff --git a/be/src/cloud/cloud_base_compaction.h
b/be/src/cloud/cloud_base_compaction.h
index 3950defbaed..696e0a6e020 100644
--- a/be/src/cloud/cloud_base_compaction.h
+++ b/be/src/cloud/cloud_base_compaction.h
@@ -17,6 +17,7 @@
#pragma once
+#include <chrono>
#include <memory>
#include <optional>
@@ -46,6 +47,9 @@ public:
private:
Status pick_rowsets_to_compact();
+ void record_compaction_success(std::chrono::steady_clock::time_point
execution_start_time);
+ Status record_compaction_failure(Status status);
+
std::string_view compaction_name() const override { return
"CloudBaseCompaction"; }
Status modify_rowsets() override;
diff --git a/be/src/cloud/cloud_cumulative_compaction.cpp
b/be/src/cloud/cloud_cumulative_compaction.cpp
index 10c00bebae1..05e553291b6 100644
--- a/be/src/cloud/cloud_cumulative_compaction.cpp
+++ b/be/src/cloud/cloud_cumulative_compaction.cpp
@@ -271,26 +271,22 @@ Status CloudCumulativeCompaction::execute_compact() {
SCOPED_ATTACH_TASK(_mem_tracker);
- using namespace std::chrono;
- auto start = steady_clock::now();
- Status st;
- Defer defer_set_st([&] {
- cloud_tablet()->set_last_cumu_compaction_status(st.to_string());
- if (!st.ok()) {
-
cloud_tablet()->set_last_cumu_compaction_failure_time(UnixMillis());
- } else {
-
cloud_tablet()->set_last_cumu_compaction_success_time(UnixMillis());
- }
- });
- st = CloudCompactionMixin::execute_compact();
- if (!st.ok()) {
- LOG(WARNING) << "fail to do " << compaction_name() << ". res=" << st
- << ", tablet=" << _tablet->tablet_id()
- << ", output_version=" << _output_version;
- return st;
+ const auto execution_start_time = std::chrono::steady_clock::now();
+ Status status = CloudCompactionMixin::execute_compact();
+ if (!status.ok()) {
+ return record_compaction_failure(status);
}
+ record_compaction_success(execution_start_time);
+ return Status::OK();
+}
+
+void CloudCumulativeCompaction::record_compaction_success(
+ std::chrono::steady_clock::time_point execution_start_time) {
LOG_INFO("finish CloudCumulativeCompaction, tablet_id={}, cost={}ms,
range=[{}-{}]",
- _tablet->tablet_id(),
duration_cast<milliseconds>(steady_clock::now() - start).count(),
+ _tablet->tablet_id(),
+ std::chrono::duration_cast<std::chrono::milliseconds>(
+ std::chrono::steady_clock::now() - execution_start_time)
+ .count(),
_input_rowsets.front()->start_version(),
_input_rowsets.back()->end_version())
.tag("job_id", _uuid)
.tag("input_rowsets", _input_rowsets.size())
@@ -320,8 +316,16 @@ Status CloudCumulativeCompaction::execute_compact() {
_input_rowsets_total_size);
cumu_output_size << _output_rowset->total_disk_size();
- st = Status::OK();
- return st;
+ cloud_tablet()->set_last_cumu_compaction_status(Status::OK().to_string());
+ cloud_tablet()->set_last_cumu_compaction_success_time(UnixMillis());
+}
+
+Status CloudCumulativeCompaction::record_compaction_failure(Status status) {
+ LOG(WARNING) << "fail to do " << compaction_name() << ". res=" << status
+ << ", tablet=" << _tablet->tablet_id() << ", output_version="
<< _output_version;
+ cloud_tablet()->set_last_cumu_compaction_status(status.to_string());
+ cloud_tablet()->set_last_cumu_compaction_failure_time(UnixMillis());
+ return status;
}
bool CloudCumulativeCompaction::should_calculate_new_cumulative_point(
diff --git a/be/src/cloud/cloud_cumulative_compaction.h
b/be/src/cloud/cloud_cumulative_compaction.h
index c111e2c0773..20836fa3627 100644
--- a/be/src/cloud/cloud_cumulative_compaction.h
+++ b/be/src/cloud/cloud_cumulative_compaction.h
@@ -17,6 +17,7 @@
#pragma once
+#include <chrono>
#include <limits>
#include <memory>
#include <optional>
@@ -82,6 +83,9 @@ private:
Status do_merge_input_rowsets(const std::vector<RowsetReaderSharedPtr>&
input_rs_readers,
MergeInputRowsetsResult* result) override;
+ void record_compaction_success(std::chrono::steady_clock::time_point
execution_start_time);
+ Status record_compaction_failure(Status status);
+
void update_output_rowset_after_build(const MergeInputRowsetsResult&
result) override;
bool should_calculate_new_cumulative_point(int64_t input_cumulative_point)
const;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]