Coverage Report

Created: 2026-09-27 17:33

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/cloud/cloud_compaction_action.cpp
Line
Count
Source
1
// Licensed to the Apache Software Foundation (ASF) under one
2
// or more contributor license agreements.  See the NOTICE file
3
// distributed with this work for additional information
4
// regarding copyright ownership.  The ASF licenses this file
5
// to you under the Apache License, Version 2.0 (the
6
// "License"); you may not use this file except in compliance
7
// with the License.  You may obtain a copy of the License at
8
//
9
//   http://www.apache.org/licenses/LICENSE-2.0
10
//
11
// Unless required by applicable law or agreed to in writing,
12
// software distributed under the License is distributed on an
13
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14
// KIND, either express or implied.  See the License for the
15
// specific language governing permissions and limitations
16
// under the License.
17
18
#include "cloud/cloud_compaction_action.h"
19
20
// IWYU pragma: no_include <bits/chrono.h>
21
#include <chrono> // IWYU pragma: keep
22
#include <exception>
23
#include <future>
24
#include <memory>
25
#include <mutex>
26
#include <sstream>
27
#include <string>
28
#include <thread>
29
#include <utility>
30
31
#include "absl/strings/substitute.h"
32
#include "cloud/cloud_base_compaction.h"
33
#include "cloud/cloud_compaction_action.h"
34
#include "cloud/cloud_cumulative_compaction.h"
35
#include "cloud/cloud_full_compaction.h"
36
#include "cloud/cloud_tablet.h"
37
#include "cloud/cloud_tablet_mgr.h"
38
#include "common/logging.h"
39
#include "common/metrics/doris_metrics.h"
40
#include "common/status.h"
41
#include "service/http/http_channel.h"
42
#include "service/http/http_headers.h"
43
#include "service/http/http_request.h"
44
#include "service/http/http_status.h"
45
#include "storage/compaction/base_compaction.h"
46
#include "storage/compaction/cumulative_compaction.h"
47
#include "storage/compaction/cumulative_compaction_policy.h"
48
#include "storage/compaction/cumulative_compaction_time_series_policy.h"
49
#include "storage/compaction/full_compaction.h"
50
#include "storage/compaction_task_tracker.h"
51
#include "storage/olap_define.h"
52
#include "storage/storage_engine.h"
53
#include "storage/tablet/tablet_manager.h"
54
#include "util/stopwatch.hpp"
55
56
namespace doris {
57
using namespace ErrorCode;
58
59
namespace {}
60
61
const static std::string HEADER_JSON = "application/json";
62
63
CloudCompactionAction::CloudCompactionAction(CompactionActionType ctype, ExecEnv* exec_env,
64
                                             CloudStorageEngine& engine, TPrivilegeHier::type hier,
65
                                             TPrivilegeType::type ptype)
66
3
        : HttpHandlerWithAuth(exec_env, hier, ptype), _engine(engine), _compaction_type(ctype) {}
67
68
/// check param and fetch tablet_id & table_id from req
69
531
static Status _check_param(HttpRequest* req, uint64_t* tablet_id, uint64_t* table_id) {
70
    // req tablet id and table id, we have to set only one of them.
71
531
    std::string req_tablet_id = req->param(TABLET_ID_KEY);
72
531
    std::string req_table_id = req->param(TABLE_ID_KEY);
73
531
    if (req_tablet_id == "") {
74
100
        if (req_table_id == "") {
75
            // both tablet id and table id are empty, return error.
76
0
            return Status::InternalError(
77
0
                    "tablet id and table id can not be empty at the same time!");
78
100
        } else {
79
100
            try {
80
100
                *table_id = std::stoull(req_table_id);
81
100
            } catch (const std::exception& e) {
82
0
                return Status::InternalError("convert table_id failed, {}", e.what());
83
0
            }
84
100
            return Status::OK();
85
100
        }
86
431
    } else {
87
431
        if (req_table_id == "") {
88
431
            try {
89
431
                *tablet_id = std::stoull(req_tablet_id);
90
431
            } catch (const std::exception& e) {
91
0
                return Status::InternalError("convert tablet_id failed, {}", e.what());
92
0
            }
93
431
            return Status::OK();
94
431
        } else {
95
            // both tablet id and table id are not empty, return err.
96
0
            return Status::InternalError("tablet id and table id can not be set at the same time!");
97
0
        }
98
431
    }
99
531
}
100
101
/// retrieve specific id from req
102
2.88k
static Status _check_param(HttpRequest* req, uint64_t* id_param, const std::string param_name) {
103
2.88k
    const auto& req_id_param = req->param(param_name);
104
2.88k
    if (!req_id_param.empty()) {
105
2.88k
        try {
106
2.88k
            *id_param = std::stoull(req_id_param);
107
2.88k
        } catch (const std::exception& e) {
108
0
            return Status::InternalError("convert {} failed, {}", param_name, e.what());
109
0
        }
110
2.88k
    }
111
112
2.88k
    return Status::OK();
113
2.88k
}
114
115
// for viewing the compaction status
116
2.42k
Status CloudCompactionAction::_handle_show_compaction(HttpRequest* req, std::string* json_result) {
117
2.42k
    uint64_t tablet_id = 0;
118
2.42k
    RETURN_NOT_OK_STATUS_WITH_WARN(_check_param(req, &tablet_id, TABLET_ID_KEY),
119
2.42k
                                   "check param failed");
120
2.42k
    if (tablet_id == 0) {
121
0
        return Status::InternalError("check param failed: missing tablet_id");
122
0
    }
123
124
2.42k
    LOG(INFO) << "begin to handle show compaction, tablet id: " << tablet_id;
125
126
    //TabletSharedPtr tablet = _engine.tablet_manager()->get_tablet(tablet_id);
127
2.42k
    CloudTabletSPtr tablet = DORIS_TRY(_engine.tablet_mgr().get_tablet(tablet_id));
128
2.42k
    if (tablet == nullptr) {
129
0
        return Status::NotFound("Tablet not found. tablet_id={}", tablet_id);
130
0
    }
131
132
2.42k
    tablet->get_compaction_status(json_result);
133
2.42k
    LOG(INFO) << "finished to handle show compaction, tablet id: " << tablet_id;
134
2.42k
    return Status::OK();
135
2.42k
}
136
137
531
Status CloudCompactionAction::_handle_run_compaction(HttpRequest* req, std::string* json_result) {
138
    // 1. param check
139
    // check req_tablet_id or req_table_id is not empty and can not be set together.
140
531
    uint64_t tablet_id = 0;
141
531
    uint64_t table_id = 0;
142
531
    RETURN_NOT_OK_STATUS_WITH_WARN(_check_param(req, &tablet_id, &table_id), "check param failed");
143
531
    LOG(INFO) << "begin to handle run compaction, tablet id: " << tablet_id
144
531
              << " table id: " << table_id;
145
146
    // check compaction_type equals 'base' or 'cumulative'
147
531
    auto& compaction_type = req->param(PARAM_COMPACTION_TYPE);
148
531
    if (compaction_type == PARAM_COMPACTION_ROW_BINLOG_TTL) {
149
0
        if (tablet_id == 0 || table_id != 0) {
150
0
            return Status::InvalidArgument("row_binlog_ttl requires a tablet_id");
151
0
        }
152
0
        RETURN_IF_ERROR(_engine.submit_row_binlog_ttl(tablet_id, true));
153
0
        *json_result = R"({"status":"Success","msg":"ROW binlog TTL task queued"})";
154
0
        return Status::OK();
155
0
    }
156
531
    if (compaction_type != PARAM_COMPACTION_BASE &&
157
531
        compaction_type != PARAM_COMPACTION_CUMULATIVE &&
158
531
        compaction_type != PARAM_COMPACTION_FULL) {
159
0
        return Status::NotSupported("The compaction type '{}' is not supported", compaction_type);
160
0
    }
161
531
    bool sync_delete_bitmap = compaction_type != PARAM_COMPACTION_FULL;
162
531
    CloudTabletSPtr tablet =
163
531
            DORIS_TRY(_engine.tablet_mgr().get_tablet(tablet_id, false, sync_delete_bitmap));
164
431
    if (tablet == nullptr) {
165
0
        return Status::NotFound("Tablet not found. tablet_id={}", tablet_id);
166
0
    }
167
168
431
    if (compaction_type == PARAM_COMPACTION_BASE) {
169
1
        tablet->set_last_base_compaction_schedule_time(UnixMillis());
170
430
    } else if (compaction_type == PARAM_COMPACTION_CUMULATIVE) {
171
371
        tablet->set_last_cumu_compaction_schedule_time(UnixMillis());
172
371
    } else if (compaction_type == PARAM_COMPACTION_FULL) {
173
59
        tablet->set_last_full_compaction_schedule_time(UnixMillis());
174
59
    }
175
176
431
    LOG(INFO) << "manual submit compaction task, tablet id: " << tablet_id
177
431
              << " table id: " << table_id;
178
    // 3. submit compaction task (trigger_method=1 for MANUAL)
179
431
    RETURN_IF_ERROR(_engine.submit_compaction_task(
180
431
            tablet,
181
431
            compaction_type == PARAM_COMPACTION_BASE         ? CompactionType::BASE_COMPACTION
182
431
            : compaction_type == PARAM_COMPACTION_CUMULATIVE ? CompactionType::CUMULATIVE_COMPACTION
183
431
                                                             : CompactionType::FULL_COMPACTION,
184
431
            /*trigger_method=*/1));
185
186
431
    LOG(INFO) << "Manual compaction task is successfully triggered, tablet id: " << tablet_id
187
410
              << " table id: " << table_id;
188
410
    *json_result =
189
410
            R"({"status": "Success", "msg": "compaction task is successfully triggered. Table id: )" +
190
410
            std::to_string(table_id) + ". Tablet id: " + std::to_string(tablet_id) + "\"}";
191
410
    return Status::OK();
192
431
}
193
194
Status CloudCompactionAction::_handle_run_status_compaction(HttpRequest* req,
195
457
                                                            std::string* json_result) {
196
457
    uint64_t tablet_id = 0;
197
457
    RETURN_NOT_OK_STATUS_WITH_WARN(_check_param(req, &tablet_id, TABLET_ID_KEY),
198
457
                                   "check param failed");
199
457
    LOG(INFO) << "begin to handle run status compaction, tablet id: " << tablet_id;
200
201
457
    if (tablet_id == 0) {
202
        // overall compaction status
203
0
        RETURN_IF_ERROR(_engine.get_compaction_status_json(json_result));
204
457
    } else {
205
457
        std::string json_template = R"({
206
457
            "status" : "Success",
207
457
            "run_status" : $0,
208
457
            "msg" : "$1",
209
457
            "tablet_id" : $2,
210
457
            "compact_type" : "$3"
211
457
        })";
212
213
457
        std::string msg = "compaction task for this tablet is not running";
214
457
        std::string compaction_type;
215
457
        bool run_status = false;
216
217
457
        if (_engine.has_cumu_compaction(tablet_id)) {
218
16
            msg = "compaction task for this tablet is running";
219
16
            compaction_type = "cumulative";
220
16
            run_status = true;
221
16
            *json_result =
222
16
                    absl::Substitute(json_template, run_status, msg, tablet_id, compaction_type);
223
16
            return Status::OK();
224
16
        }
225
226
441
        if (_engine.has_base_compaction(tablet_id)) {
227
0
            msg = "compaction task for this tablet is running";
228
0
            compaction_type = "base";
229
0
            run_status = true;
230
0
            *json_result =
231
0
                    absl::Substitute(json_template, run_status, msg, tablet_id, compaction_type);
232
0
            return Status::OK();
233
0
        }
234
235
441
        if (_engine.has_full_compaction(tablet_id)) {
236
26
            msg = "compaction task for this tablet is running";
237
26
            compaction_type = "full";
238
26
            run_status = true;
239
26
            *json_result =
240
26
                    absl::Substitute(json_template, run_status, msg, tablet_id, compaction_type);
241
26
            return Status::OK();
242
26
        }
243
        // not running any compaction
244
415
        *json_result = absl::Substitute(json_template, run_status, msg, tablet_id, compaction_type);
245
415
    }
246
457
    LOG(INFO) << "finished to handle run status compaction, tablet id: " << tablet_id;
247
415
    return Status::OK();
248
457
}
249
250
3.41k
void CloudCompactionAction::handle(HttpRequest* req) {
251
3.41k
    req->add_output_header(HttpHeaders::CONTENT_TYPE, HEADER_JSON.c_str());
252
253
3.41k
    if (_compaction_type == CompactionActionType::SHOW_INFO) {
254
2.42k
        std::string json_result;
255
2.42k
        Status st = _handle_show_compaction(req, &json_result);
256
2.42k
        if (!st.ok()) {
257
0
            HttpChannel::send_reply(req, HttpStatus::OK, st.to_json());
258
2.42k
        } else {
259
2.42k
            HttpChannel::send_reply(req, HttpStatus::OK, json_result);
260
2.42k
        }
261
2.42k
    } else if (_compaction_type == CompactionActionType::RUN_COMPACTION) {
262
531
        std::string json_result;
263
531
        Status st = _handle_run_compaction(req, &json_result);
264
531
        if (!st.ok()) {
265
121
            HttpChannel::send_reply(req, HttpStatus::OK, st.to_json());
266
410
        } else {
267
410
            HttpChannel::send_reply(req, HttpStatus::OK, json_result);
268
410
        }
269
531
    } else {
270
457
        std::string json_result;
271
457
        Status st = _handle_run_status_compaction(req, &json_result);
272
457
        if (!st.ok()) {
273
0
            HttpChannel::send_reply(req, HttpStatus::OK, st.to_json());
274
457
        } else {
275
457
            HttpChannel::send_reply(req, HttpStatus::OK, json_result);
276
457
        }
277
457
    }
278
3.41k
}
279
280
} // end namespace doris