Coverage Report

Created: 2026-08-06 11:34

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/storage/rowset_version_mgr.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 <brpc/controller.h>
19
#include <bthread/bthread.h>
20
#include <bthread/countdown_event.h>
21
#include <bthread/mutex.h>
22
#include <bthread/types.h>
23
#include <bvar/latency_recorder.h>
24
#include <gen_cpp/FrontendService_types.h>
25
#include <gen_cpp/HeartbeatService_types.h>
26
#include <gen_cpp/Types_types.h>
27
#include <gen_cpp/internal_service.pb.h>
28
#include <gen_cpp/olap_file.pb.h>
29
#include <glog/logging.h>
30
31
#include <cstdint>
32
#include <memory>
33
#include <mutex>
34
#include <optional>
35
#include <ranges>
36
#include <sstream>
37
#include <utility>
38
39
#include "cloud/config.h"
40
#include "common/status.h"
41
#include "cpp/sync_point.h"
42
#include "runtime/cluster_info.h"
43
#include "service/backend_options.h"
44
#include "service/internal_service.h"
45
#include "storage/olap_common.h"
46
#include "storage/rowset/rowset.h"
47
#include "storage/rowset/rowset_factory.h"
48
#include "storage/rowset/rowset_reader.h"
49
#include "storage/tablet/base_tablet.h"
50
#include "util/brpc_client_cache.h"
51
#include "util/client_cache.h"
52
#include "util/debug_points.h"
53
#include "util/thrift_rpc_helper.h"
54
#include "util/time.h"
55
56
namespace doris {
57
58
using namespace ErrorCode;
59
using namespace std::ranges;
60
61
static bvar::LatencyRecorder g_remote_fetch_tablet_rowsets_single_request_latency(
62
        "remote_fetch_rowsets_single_rpc");
63
static bvar::LatencyRecorder g_remote_fetch_tablet_rowsets_latency("remote_fetch_rowsets");
64
65
[[nodiscard]] Result<std::vector<Version>> BaseTablet::capture_consistent_versions_unlocked(
66
1.49M
        const Version& version_range, const CaptureRowsetOps& options) const {
67
1.49M
    std::vector<Version> version_path;
68
1.49M
    auto& version_tracker =
69
1.49M
            options.capture_row_binlog ? _row_binlog_version_tracker : _timestamped_version_tracker;
70
1.49M
    auto st = version_tracker.capture_consistent_versions(version_range, &version_path);
71
1.49M
    if (!st && !options.quiet) {
72
2
        auto missed_versions =
73
2
                get_missed_versions_unlocked(version_range.second, options.capture_row_binlog);
74
2
        if (missed_versions.empty()) {
75
0
            LOG(WARNING) << fmt::format(
76
0
                    "version already has been merged. version_range={}, max_version={}, "
77
0
                    "tablet_id={}",
78
0
                    version_range.to_string(), _tablet_meta->max_version().second, tablet_id());
79
0
            return ResultError(Status::Error<VERSION_ALREADY_MERGED>(
80
0
                    "missed versions is empty, version_range={}, max_version={}, tablet_id={}",
81
0
                    version_range.to_string(), _tablet_meta->max_version().second, tablet_id()));
82
0
        }
83
2
        LOG(WARNING) << fmt::format("missed version for version_range={}, tablet_id={}, st={}",
84
2
                                    version_range.to_string(), tablet_id(), st);
85
2
        _print_missed_versions(missed_versions);
86
2
        if (!options.skip_missing_versions) {
87
2
            return ResultError(std::move(st));
88
2
        }
89
2
        LOG(WARNING) << "force skipping missing version for tablet:" << tablet_id();
90
0
    }
91
1.49M
    DBUG_EXECUTE_IF("Tablet::capture_consistent_versions.inject_failure", {
92
1.49M
        auto tablet_id = dp->param<int64_t>("tablet_id", -1);
93
1.49M
        auto skip_by_option = dp->param<bool>("skip_by_option", false);
94
1.49M
        if (skip_by_option && !options.enable_fetch_rowsets_from_peers) {
95
1.49M
            return version_path;
96
1.49M
        }
97
1.49M
        if ((tablet_id != -1 && (tablet_id == _tablet_meta->tablet_id())) || tablet_id == -2) {
98
1.49M
            return ResultError(Status::Error<VERSION_ALREADY_MERGED>("version already merged"));
99
1.49M
        }
100
1.49M
    });
101
1.49M
    return version_path;
102
1.49M
}
103
104
[[nodiscard]] Result<CaptureRowsetResult> BaseTablet::capture_consistent_rowsets_unlocked(
105
1.49M
        const Version& version_range, const CaptureRowsetOps& options) const {
106
1.49M
    CaptureRowsetResult result;
107
1.49M
    auto& rowsets = result.rowsets;
108
1.49M
    auto maybe_versions = capture_consistent_versions_unlocked(version_range, options);
109
1.49M
    if (maybe_versions) {
110
1.49M
        const auto& version_paths = maybe_versions.value();
111
1.49M
        rowsets.reserve(version_paths.size());
112
113
1.49M
        auto rowset_for_version = [&](const Version& version,
114
6.17M
                                      bool include_stale) -> Result<RowsetSharedPtr> {
115
6.17M
            const auto& rs_version_map =
116
6.17M
                    options.capture_row_binlog ? _row_binlog_rs_version_map : _rs_version_map;
117
6.17M
            if (auto it = rs_version_map.find(version); it != rs_version_map.end()) {
118
6.16M
                return it->second;
119
6.16M
            } else {
120
3.39k
                VLOG_NOTICE << "fail to find Rowset in "
121
3.34k
                            << (options.capture_row_binlog ? "row_binlog_rs_version" : "rs_version")
122
3.34k
                            << " for version. tablet=" << tablet_id() << ", version='"
123
3.34k
                            << version.first << "-" << version.second;
124
3.39k
            }
125
3.39k
            if (!options.capture_row_binlog && include_stale) {
126
58
                if (auto it = _stale_rs_version_map.find(version);
127
58
                    it != _stale_rs_version_map.end()) {
128
58
                    return it->second;
129
58
                } else {
130
0
                    LOG(WARNING) << fmt::format(
131
0
                            "fail to find Rowset in stale_rs_version for version. tablet={}, "
132
0
                            "version={}-{}",
133
0
                            tablet_id(), version.first, version.second);
134
0
                }
135
58
            }
136
3.34k
            return ResultError(Status::Error<CAPTURE_ROWSET_ERROR>(
137
3.34k
                    "failed to find rowset for version={}", version.to_string()));
138
3.39k
        };
139
140
6.17M
        for (const auto& version : version_paths) {
141
6.17M
            auto ret = rowset_for_version(version, options.include_stale_rowsets);
142
6.17M
            if (!ret) {
143
0
                return ResultError(std::move(ret.error()));
144
0
            }
145
146
6.17M
            rowsets.push_back(std::move(ret.value()));
147
6.17M
        }
148
1.49M
        if (options.capture_row_binlog) {
149
0
            result.delete_bitmap = _tablet_meta->binlog_delvec_ptr();
150
1.49M
        } else if (keys_type() == KeysType::UNIQUE_KEYS && enable_unique_key_merge_on_write()) {
151
857k
            result.delete_bitmap = _tablet_meta->delete_bitmap_ptr();
152
857k
        }
153
1.49M
        return result;
154
1.49M
    }
155
156
18.4E
    if (!config::is_cloud_mode() || !options.enable_fetch_rowsets_from_peers) {
157
2
        return ResultError(std::move(maybe_versions.error()));
158
2
    }
159
18.4E
    auto ret = _remote_capture_rowsets(version_range);
160
18.4E
    if (!ret) {
161
0
        auto st = Status::Error<VERSION_ALREADY_MERGED>(
162
0
                "version already merged, meet error during remote capturing rowsets, "
163
0
                "error={}, version_range={}",
164
0
                ret.error().to_string(), version_range.to_string());
165
0
        return ResultError(std::move(st));
166
0
    }
167
18.4E
    return ret;
168
18.4E
}
169
170
[[nodiscard]] Result<std::vector<RowSetSplits>> BaseTablet::capture_rs_readers_unlocked(
171
5.41k
        const Version& version_range, const CaptureRowsetOps& options) const {
172
5.41k
    auto maybe_rs_list = capture_consistent_rowsets_unlocked(version_range, options);
173
5.41k
    if (!maybe_rs_list) {
174
1
        return ResultError(std::move(maybe_rs_list.error()));
175
1
    }
176
5.41k
    const auto& rs_list = maybe_rs_list.value().rowsets;
177
5.41k
    std::vector<RowSetSplits> rs_splits;
178
5.41k
    rs_splits.reserve(rs_list.size());
179
9.80k
    for (const auto& rs : rs_list) {
180
9.80k
        RowsetReaderSharedPtr rs_reader;
181
9.80k
        auto st = rs->create_reader(&rs_reader);
182
9.80k
        if (!st) {
183
0
            return ResultError(Status::Error<CAPTURE_ROWSET_READER_ERROR>(
184
0
                    "failed to create reader for rowset={}, reason={}", rs->rowset_id().to_string(),
185
0
                    st.to_string()));
186
0
        }
187
9.80k
        rs_splits.emplace_back(std::move(rs_reader));
188
9.80k
    }
189
5.41k
    return rs_splits;
190
5.41k
}
191
192
[[nodiscard]] Result<TabletReadSource> BaseTablet::capture_read_source(
193
1.48M
        const Version& version_range, const CaptureRowsetOps& options) {
194
1.48M
    std::shared_lock rdlock(get_header_lock());
195
1.48M
    auto maybe_result = capture_consistent_rowsets_unlocked(version_range, options);
196
1.48M
    if (!maybe_result) {
197
1
        return ResultError(std::move(maybe_result.error()));
198
1
    }
199
1.48M
    auto rowsets_result = std::move(maybe_result.value());
200
1.48M
    TabletReadSource read_source;
201
1.48M
    read_source.delete_bitmap = std::move(rowsets_result.delete_bitmap);
202
1.48M
    const auto& rowsets = rowsets_result.rowsets;
203
1.48M
    read_source.rs_splits.reserve(rowsets.size());
204
6.16M
    for (const auto& rs : rowsets) {
205
6.16M
        RowsetReaderSharedPtr rs_reader;
206
6.16M
        auto st = rs->create_reader(&rs_reader);
207
6.16M
        if (!st) {
208
0
            return ResultError(Status::Error<CAPTURE_ROWSET_READER_ERROR>(
209
0
                    "failed to create reader for rowset={}, reason={}", rs->rowset_id().to_string(),
210
0
                    st.to_string()));
211
0
        }
212
6.16M
        read_source.rs_splits.emplace_back(std::move(rs_reader));
213
6.16M
    }
214
1.48M
    return read_source;
215
1.48M
}
216
217
template <typename Fn, typename... Args>
218
0
bool call_bthread(bthread_t& th, const bthread_attr_t* attr, Fn&& fn, Args&&... args) {
219
0
    auto p_wrap_fn = new auto([=] { fn(args...); });
220
0
    auto call_back = [](void* ar) -> void* {
221
0
        auto f = reinterpret_cast<decltype(p_wrap_fn)>(ar);
222
0
        (*f)();
223
0
        delete f;
224
0
        return nullptr;
225
0
    };
226
0
    return bthread_start_background(&th, attr, call_back, p_wrap_fn) == 0;
227
0
}
228
229
struct GetRowsetsCntl : std::enable_shared_from_this<GetRowsetsCntl> {
230
    struct RemoteGetRowsetResult {
231
        std::vector<RowsetMetaSharedPtr> rowsets;
232
        std::unique_ptr<DeleteBitmap> delete_bitmap;
233
    };
234
235
0
    Status start_req_bg() {
236
0
        task_cnt = req_addrs.size();
237
0
        for (const auto& [ip, port] : req_addrs) {
238
0
            bthread_t tid;
239
0
            bthread_attr_t attr = BTHREAD_ATTR_NORMAL;
240
241
0
            bool succ = call_bthread(tid, &attr, [self = shared_from_this(), &ip, port]() {
242
0
                LOG(INFO) << "start to get tablet rowsets from peer BE, ip=" << ip;
243
0
                Defer defer_log {[&ip, port]() {
244
0
                    LOG(INFO) << "finish to get rowsets from peer BE, ip=" << ip
245
0
                              << ", port=" << port;
246
0
                }};
247
248
0
                PGetTabletRowsetsRequest req;
249
0
                req.set_tablet_id(self->tablet_id);
250
0
                req.set_version_start(self->version_range.first);
251
0
                req.set_version_end(self->version_range.second);
252
0
                if (self->delete_bitmap_keys.has_value()) {
253
0
                    req.mutable_delete_bitmap_keys()->CopyFrom(self->delete_bitmap_keys.value());
254
0
                }
255
0
                brpc::Controller cntl;
256
0
                cntl.set_timeout_ms(60000);
257
0
                cntl.set_max_retry(3);
258
0
                PGetTabletRowsetsResponse response;
259
0
                auto start_tm_us = MonotonicMicros();
260
0
#ifndef BE_TEST
261
0
                std::shared_ptr<PBackendService_Stub> stub =
262
0
                        ExecEnv::GetInstance()->brpc_internal_client_cache()->get_client(ip, port);
263
0
                if (stub == nullptr) {
264
0
                    self->result = ResultError(Status::InternalError(
265
0
                            "failed to fetch get_tablet_rowsets stub, ip={}, port={}", ip, port));
266
0
                    return;
267
0
                }
268
0
                stub->get_tablet_rowsets(&cntl, &req, &response, nullptr);
269
#else
270
                TEST_SYNC_POINT_CALLBACK("get_tablet_rowsets", &response);
271
#endif
272
0
                g_remote_fetch_tablet_rowsets_single_request_latency
273
0
                        << MonotonicMicros() - start_tm_us;
274
275
0
                std::unique_lock l(self->butex);
276
0
                if (self->done) {
277
0
                    return;
278
0
                }
279
0
                --self->task_cnt;
280
0
                auto resp_st = Status::create(response.status());
281
0
                DBUG_EXECUTE_IF("GetRowsetCntl::start_req_bg.inject_failure",
282
0
                                { resp_st = Status::InternalError("inject error"); });
283
0
                if (cntl.Failed() || !resp_st) {
284
0
                    if (self->task_cnt != 0) {
285
0
                        return;
286
0
                    }
287
0
                    std::stringstream reason;
288
0
                    reason << "failed to get rowsets from all replicas, tablet_id="
289
0
                           << self->tablet_id;
290
0
                    if (cntl.Failed()) {
291
0
                        reason << ", reason=[" << cntl.ErrorCode() << "] " << cntl.ErrorText();
292
0
                    } else {
293
0
                        reason << ", reason=" << resp_st.to_string();
294
0
                    }
295
0
                    self->result = ResultError(Status::InternalError(reason.str()));
296
0
                    self->done = true;
297
0
                    self->event.signal();
298
0
                    return;
299
0
                }
300
301
0
                Defer done_cb {[&]() {
302
0
                    self->done = true;
303
0
                    self->event.signal();
304
0
                }};
305
0
                std::vector<RowsetMetaSharedPtr> rs_metas;
306
0
                for (auto&& rs_pb : response.rowsets()) {
307
0
                    auto rs_meta = std::make_shared<RowsetMeta>();
308
0
                    if (!rs_meta->init_from_pb(rs_pb)) {
309
0
                        self->result =
310
0
                                ResultError(Status::InternalError("failed to init rowset from pb"));
311
0
                        return;
312
0
                    }
313
0
                    rs_metas.push_back(std::move(rs_meta));
314
0
                }
315
0
                CaptureRowsetResult result;
316
0
                self->result->rowsets = std::move(rs_metas);
317
318
0
                if (response.has_delete_bitmap()) {
319
0
                    self->result->delete_bitmap = std::make_unique<DeleteBitmap>(
320
0
                            DeleteBitmap::from_pb(response.delete_bitmap(), self->tablet_id));
321
0
                }
322
0
            });
323
324
0
            if (!succ) {
325
0
                return Status::InternalError(
326
0
                        "failed to create bthread when request rowsets for tablet={}", tablet_id);
327
0
            }
328
0
        }
329
0
        return Status::OK();
330
0
    }
331
332
0
    Result<RemoteGetRowsetResult> wait_for_ret() {
333
0
        event.wait();
334
0
        return std::move(result);
335
0
    }
336
337
    int64_t tablet_id;
338
    std::vector<std::pair<std::string, int32_t>> req_addrs;
339
    Version version_range;
340
    std::optional<DeleteBitmapPB> delete_bitmap_keys = std::nullopt;
341
342
private:
343
    size_t task_cnt;
344
345
    bthread::Mutex butex;
346
    bthread::CountdownEvent event {1};
347
    bool done = false;
348
349
    Result<RemoteGetRowsetResult> result;
350
};
351
352
Result<std::vector<std::pair<std::string, int32_t>>> get_peer_replicas_addresses(
353
0
        const int64_t tablet_id) {
354
0
    auto* cluster_info = ExecEnv::GetInstance()->cluster_info();
355
0
    DCHECK_NE(cluster_info, nullptr);
356
0
    auto master_addr = cluster_info->master_fe_addr;
357
0
    TGetTabletReplicaInfosRequest req;
358
0
    req.tablet_ids.push_back(tablet_id);
359
0
    TGetTabletReplicaInfosResult resp;
360
0
    auto st = ThriftRpcHelper::rpc<FrontendServiceClient>(
361
0
            master_addr.hostname, master_addr.port,
362
0
            [&](FrontendServiceConnection& client) { client->getTabletReplicaInfos(resp, req); });
363
0
    if (!st) {
364
0
        return ResultError(Status::InternalError(
365
0
                "failed to get tablet replica infos, rpc error={}, tablet_id={}", st.to_string(),
366
0
                tablet_id));
367
0
    }
368
369
0
    auto it = resp.tablet_replica_infos.find(tablet_id);
370
0
    if (it == resp.tablet_replica_infos.end()) {
371
0
        return ResultError(Status::InternalError("replicas not found, tablet_id={}", tablet_id));
372
0
    }
373
0
    auto replicas = it->second;
374
0
    auto local_host = BackendOptions::get_localhost();
375
0
    bool include_local_host = false;
376
0
    DBUG_EXECUTE_IF("get_peer_replicas_address.enable_local_host", { include_local_host = true; });
377
0
    auto ret_view =
378
0
            replicas | std::views::filter([&local_host, include_local_host](const auto& replica) {
379
0
                return local_host.find(replica.host) == std::string::npos || include_local_host;
380
0
            }) |
381
0
            std::views::transform([](auto& replica) {
382
0
                return std::make_pair(std::move(replica.host), replica.brpc_port);
383
0
            });
384
0
    return std::vector(ret_view.begin(), ret_view.end());
385
0
}
386
387
Result<CaptureRowsetResult> BaseTablet::_remote_capture_rowsets(
388
0
        const Version& version_range) const {
389
0
    auto start_tm_us = MonotonicMicros();
390
0
    Defer defer {
391
0
            [&]() { g_remote_fetch_tablet_rowsets_latency << MonotonicMicros() - start_tm_us; }};
392
0
#ifndef BE_TEST
393
0
    auto maybe_be_addresses = get_peer_replicas_addresses(tablet_id());
394
#else
395
    Result<std::vector<std::pair<std::string, int32_t>>> maybe_be_addresses;
396
    TEST_SYNC_POINT_CALLBACK("get_peer_replicas_addresses", &maybe_be_addresses);
397
#endif
398
0
    DBUG_EXECUTE_IF("Tablet::_remote_get_rowsets_meta.inject_replica_address_fail",
399
0
                    { maybe_be_addresses = ResultError(Status::InternalError("inject failure")); });
400
0
    if (!maybe_be_addresses) {
401
0
        return ResultError(std::move(maybe_be_addresses.error()));
402
0
    }
403
0
    auto be_addresses = std::move(maybe_be_addresses.value());
404
0
    if (be_addresses.empty()) {
405
0
        LOG(WARNING) << "no peers replica for tablet=" << tablet_id();
406
0
        return ResultError(Status::InternalError("no replicas for tablet={}", tablet_id()));
407
0
    }
408
409
0
    auto cntl = std::make_shared<GetRowsetsCntl>();
410
0
    cntl->tablet_id = tablet_id();
411
0
    cntl->req_addrs = std::move(be_addresses);
412
0
    cntl->version_range = version_range;
413
0
    bool is_mow = keys_type() == KeysType::UNIQUE_KEYS && enable_unique_key_merge_on_write();
414
0
    CaptureRowsetResult result;
415
0
    if (is_mow) {
416
0
        result.delete_bitmap =
417
0
                std::make_unique<DeleteBitmap>(_tablet_meta->delete_bitmap().snapshot());
418
0
        DeleteBitmapPB delete_bitmap_keys;
419
0
        auto keyset = result.delete_bitmap->delete_bitmap |
420
0
                      std::views::transform([](const auto& kv) { return kv.first; });
421
0
        for (const auto& key : keyset) {
422
0
            const auto& [rs_id, seg_id, version] = key;
423
0
            delete_bitmap_keys.mutable_rowset_ids()->Add(rs_id.to_string());
424
0
            delete_bitmap_keys.mutable_segment_ids()->Add(seg_id);
425
0
            delete_bitmap_keys.mutable_versions()->Add(version);
426
0
        }
427
0
        cntl->delete_bitmap_keys = std::move(delete_bitmap_keys);
428
0
    }
429
430
0
    RETURN_IF_ERROR_RESULT(cntl->start_req_bg());
431
0
    auto maybe_meta = cntl->wait_for_ret();
432
0
    if (!maybe_meta) {
433
0
        auto err = Status::InternalError(
434
0
                "tried to get rowsets from peer replicas and failed, "
435
0
                "reason={}",
436
0
                maybe_meta.error());
437
0
        return ResultError(std::move(err));
438
0
    }
439
440
0
    auto& remote_meta = maybe_meta.value();
441
0
    const auto& rs_metas = remote_meta.rowsets;
442
0
    for (const auto& rs_meta : rs_metas) {
443
0
        RowsetSharedPtr rs;
444
0
        auto st = RowsetFactory::create_rowset(_tablet_meta->tablet_schema(), {}, rs_meta, &rs);
445
0
        if (!st) {
446
0
            return ResultError(std::move(st));
447
0
        }
448
0
        result.rowsets.push_back(std::move(rs));
449
0
    }
450
0
    if (is_mow) {
451
        DCHECK_NE(result.delete_bitmap, nullptr);
452
0
        result.delete_bitmap->merge(*remote_meta.delete_bitmap);
453
0
    }
454
0
    return result;
455
0
}
456
457
} // namespace doris