Coverage Report

Created: 2026-08-08 02:22

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
be/src/io/cache/peer_file_cache_reader.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
#include "io/cache/peer_file_cache_reader.h"
18
19
#include <brpc/controller.h>
20
#include <butil/iobuf.h>
21
#include <bvar/latency_recorder.h>
22
#include <bvar/reducer.h>
23
#include <fmt/format.h>
24
#include <gen_cpp/internal_service.pb.h>
25
#include <glog/logging.h>
26
27
#include <algorithm>
28
#include <utility>
29
30
#include "common/compiler_util.h" // IWYU pragma: keep
31
#include "common/metrics/doris_metrics.h"
32
#include "runtime/exec_env.h"
33
#include "runtime/runtime_profile.h"
34
#include "runtime/thread_context.h"
35
#include "runtime/workload_group/workload_group.h"
36
#include "runtime/workload_management/io_throttle.h"
37
#include "runtime/workload_management/resource_context.h"
38
#include "util/brpc_client_cache.h"
39
#include "util/bvar_helper.h"
40
#include "util/debug_points.h"
41
#include "util/defer_op.h"
42
#include "util/network_util.h"
43
44
namespace doris::io {
45
46
namespace {
47
48
struct ExpectedPeerFetch {
49
    std::vector<FileBlock::Range> expected_ranges;
50
    std::vector<FileBlock::Range> pending_ranges;
51
    std::vector<size_t> expected_block_indexes;
52
    size_t expected_bytes = 0;
53
};
54
55
91
size_t clip_requested_range(const FileBlock::Range& range, size_t file_size) {
56
91
    if (range.left >= file_size) {
57
2
        return 0;
58
2
    }
59
89
    return std::min(file_size - range.left, range.size());
60
91
}
61
62
ExpectedPeerFetch build_expected_peer_fetch(const std::vector<FileBlockSPtr>& blocks,
63
40
                                            size_t file_size, PFetchPeerDataRequest* req) {
64
40
    ExpectedPeerFetch expected;
65
131
    for (size_t block_idx = 0; block_idx < blocks.size(); ++block_idx) {
66
91
        const auto& blk = blocks[block_idx];
67
91
        auto* cb = req->add_cache_req();
68
91
        cb->set_block_offset(static_cast<int64_t>(blk->range().left));
69
91
        cb->set_block_size(static_cast<int64_t>(blk->range().size()));
70
91
        const size_t clipped_size = clip_requested_range(blk->range(), file_size);
71
91
        if (clipped_size == 0) {
72
2
            continue;
73
2
        }
74
89
        expected.expected_block_indexes.push_back(block_idx);
75
89
        expected.expected_ranges.emplace_back(blk->range().left,
76
89
                                              blk->range().left + clipped_size - 1);
77
89
        expected.pending_ranges.emplace_back(expected.expected_ranges.back());
78
89
        expected.expected_bytes += clipped_size;
79
89
    }
80
40
    return expected;
81
40
}
82
83
18
Status cut_attachment_payload(butil::IOBuf* attachment, size_t size, butil::IOBuf* out) {
84
18
    const size_t cut_size = attachment->cutn(out, size);
85
18
    if (cut_size != size) {
86
1
        return Status::InternalError<false>(
87
1
                "peer cache read incomplete attachment: need={}, got={}", size, cut_size);
88
1
    }
89
17
    return Status::OK();
90
18
}
91
92
bool subtract_pending_range(std::vector<FileBlock::Range>& pending_ranges,
93
57
                            const FileBlock::Range& response_range) {
94
58
    for (size_t idx = 0; idx < pending_ranges.size(); ++idx) {
95
57
        auto pending_range = pending_ranges[idx];
96
57
        if (response_range.left < pending_range.left ||
97
57
            response_range.right > pending_range.right) {
98
1
            continue;
99
1
        }
100
101
56
        if (response_range.left == pending_range.left &&
102
56
            response_range.right == pending_range.right) {
103
55
            pending_ranges.erase(pending_ranges.begin() + idx);
104
55
        } else if (response_range.left == pending_range.left) {
105
1
            pending_ranges[idx].left = response_range.right + 1;
106
1
        } else if (response_range.right == pending_range.right) {
107
0
            pending_ranges[idx].right = response_range.left - 1;
108
0
        } else {
109
0
            auto right_remain = FileBlock::Range(response_range.right + 1, pending_range.right);
110
0
            pending_ranges[idx].right = response_range.left - 1;
111
0
            pending_ranges.insert(pending_ranges.begin() + idx + 1, right_remain);
112
0
        }
113
56
        return true;
114
57
    }
115
1
    return false;
116
57
}
117
118
int find_expected_range_idx(const std::vector<FileBlock::Range>& expected_ranges,
119
58
                            const FileBlock::Range& response_range) {
120
98
    for (size_t idx = 0; idx < expected_ranges.size(); ++idx) {
121
97
        const auto& expected_range = expected_ranges[idx];
122
97
        if (response_range.left >= expected_range.left &&
123
97
            response_range.right <= expected_range.right) {
124
57
            return static_cast<int>(idx);
125
57
        }
126
97
    }
127
1
    return -1;
128
58
}
129
130
} // namespace
131
132
// read from peer
133
134
bvar::Adder<uint64_t> peer_cache_reader_failed_counter("peer_cache_reader", "failed_counter");
135
bvar::Adder<uint64_t> peer_cache_reader_succ_counter("peer_cache_reader", "succ_counter");
136
bvar::LatencyRecorder peer_bytes_per_read("peer_cache_reader", "bytes_per_read"); // also QPS
137
bvar::Adder<uint64_t> peer_cache_reader_total("peer_cache_reader", "total_num");
138
bvar::Adder<uint64_t> peer_cache_being_read("peer_cache_reader", "file_being_read");
139
bvar::Adder<uint64_t> peer_cache_reader_read_counter("peer_cache_reader", "read_at");
140
bvar::LatencyRecorder peer_cache_reader_latency("peer_cache_reader", "peer_latency");
141
bvar::PerSecond<bvar::Adder<uint64_t>> peer_get_request_qps("peer_cache_reader", "peer_get_request",
142
                                                            &peer_cache_reader_read_counter);
143
bvar::Adder<uint64_t> peer_bytes_read_total("peer_cache_reader", "bytes_read");
144
bvar::PerSecond<bvar::Adder<uint64_t>> peer_read_througthput("peer_cache_reader",
145
                                                             "peer_read_throughput",
146
                                                             &peer_bytes_read_total);
147
148
PeerFileCacheReader::PeerFileCacheReader(const io::Path& file_path, bool is_doris_table,
149
                                         std::string host, int port)
150
42
        : _path(file_path), _is_doris_table(is_doris_table), _host(host), _port(port) {
151
42
    peer_cache_reader_total << 1;
152
42
    peer_cache_being_read << 1;
153
42
}
154
155
42
PeerFileCacheReader::~PeerFileCacheReader() {
156
42
    peer_cache_being_read << -1;
157
42
}
158
159
Status PeerFileCacheReader::fetch_blocks(const std::vector<FileBlockSPtr>& blocks,
160
                                         PeerFetchResult* result, size_t file_size,
161
                                         const IOContext* ctx, bool request_fill, int64_t tablet_id,
162
42
                                         std::string resource_id) {
163
42
    (void)ctx;
164
42
    if (result == nullptr) {
165
0
        return Status::InvalidArgument("peer cache fetch requires non-null result");
166
0
    }
167
42
    result->clear();
168
42
    VLOG_DEBUG << "enter PeerFileCacheReader::fetch_blocks";
169
42
    if (blocks.empty()) {
170
1
        return Status::OK();
171
1
    }
172
41
    if (!_is_doris_table) {
173
1
        return Status::NotSupported<false>("peer cache fetch only supports doris table segments");
174
1
    }
175
176
40
    PFetchPeerDataRequest req;
177
40
    req.set_type(PFetchPeerDataRequest_Type_PEER_FILE_CACHE_BLOCK);
178
40
    req.set_path(_path.native());
179
40
    req.set_file_size(static_cast<int64_t>(file_size));
180
40
    auto* rowset_meta_pb = req.mutable_rowset_meta();
181
40
    rowset_meta_pb->Clear();
182
    // RowsetMetaPB still has deprecated proto2 required rowset_id. Set a dummy value so
183
    // the RPC can be serialized; current peer read/fill paths only read tablet_id/resource_id.
184
40
    rowset_meta_pb->set_rowset_id(0);
185
40
    rowset_meta_pb->set_resource_id(resource_id);
186
40
    rowset_meta_pb->set_tablet_id(tablet_id);
187
40
    if (request_fill) {
188
        // Ask the peer server to pull missing blocks from remote storage before serving them.
189
        // Only set for cross-CG reads targeting the designated fill compute group
190
        // (peer_cache_fill_compute_group_id). Server still gates with enable_peer_server_cache_fill.
191
2
        req.set_request_cache_fill(true);
192
2
    }
193
    // Always advertise attachment support. Older peers can still reply in protobuf mode.
194
40
    req.set_support_attachment(true);
195
40
    auto expected = build_expected_peer_fetch(blocks, file_size, &req);
196
40
    if (expected.expected_bytes == 0) {
197
1
        return Status::OK();
198
1
    }
199
200
39
    std::string realhost = _host;
201
39
    int port = _port;
202
203
39
    auto dns_cache = ExecEnv::GetInstance()->dns_cache();
204
39
    if (dns_cache == nullptr) {
205
39
        LOG(WARNING) << "DNS cache is not initialized, skipping hostname resolve";
206
39
    } else if (!is_valid_ip(realhost)) {
207
0
        Status status = dns_cache->get(_host, &realhost);
208
0
        if (!status.ok()) {
209
0
            peer_cache_reader_failed_counter << 1;
210
0
            LOG(WARNING) << "failed to get ip from host " << _host << ": " << status.to_string();
211
0
            return Status::InternalError<false>("failed to get ip from host {}", _host);
212
0
        }
213
0
    }
214
39
    std::string brpc_addr = get_host_port(realhost, port);
215
39
    Status st = Status::OK();
216
39
    std::shared_ptr<PBackendService_Stub> brpc_stub =
217
39
            ExecEnv::GetInstance()->brpc_internal_client_cache()->get_new_client_no_cache(
218
39
                    brpc_addr);
219
39
    if (!brpc_stub) {
220
0
        peer_cache_reader_failed_counter << 1;
221
0
        LOG(WARNING) << "failed to get brpc stub " << brpc_addr;
222
0
        st = Status::RpcError<false>("Address {} is wrong", brpc_addr);
223
0
        return st;
224
0
    }
225
226
39
    size_t filled = 0;
227
39
    size_t* bytes_read = &filled;
228
39
    LIMIT_REMOTE_SCAN_IO(bytes_read);
229
39
    int64_t begin_ts = std::chrono::duration_cast<std::chrono::microseconds>(
230
39
                               std::chrono::system_clock::now().time_since_epoch())
231
39
                               .count();
232
39
    Defer defer_latency {[&]() {
233
39
        int64_t end_ts = std::chrono::duration_cast<std::chrono::microseconds>(
234
39
                                 std::chrono::system_clock::now().time_since_epoch())
235
39
                                 .count();
236
39
        peer_cache_reader_latency << (end_ts - begin_ts);
237
39
    }};
238
239
39
    brpc::Controller cntl;
240
    // Use a longer timeout when fill is requested: server may spend up to
241
    // peer_server_cache_fill_timeout_ms (default 6000ms) pulling from S3 before responding.
242
39
    cntl.set_timeout_ms(request_fill ? 7000 : 5000);
243
39
    PFetchPeerDataResponse resp;
244
39
    peer_cache_reader_read_counter << 1;
245
39
    brpc_stub->fetch_peer_data(&cntl, &req, &resp, nullptr);
246
39
    if (cntl.Failed()) {
247
1
        return Status::RpcError<false>(cntl.ErrorText());
248
1
    }
249
38
    if (resp.has_status()) {
250
38
        Status st2 = Status::create<false>(resp.status());
251
38
        LOG_EVERY_N(WARNING, 1000) << "peer cache read failed, status=" << st2.msg();
252
38
        if (!st2.ok()) return st2;
253
38
    }
254
255
    // Metadata stays in resp.datas(); payload may come from protobuf bytes or BRPC attachment.
256
28
    const bool use_attachment = resp.has_data_in_attachment() && resp.data_in_attachment();
257
28
    butil::IOBuf remaining_attachment(cntl.response_attachment());
258
28
    result->chunks.reserve(resp.datas_size());
259
60
    for (const auto& data : resp.datas()) {
260
60
        if (data.block_offset() < 0 || data.block_size() < 0) {
261
1
            peer_cache_reader_failed_counter << 1;
262
1
            result->clear();
263
1
            return Status::InternalError<false>(
264
1
                    "peer cache read invalid block metadata: offset={}, size={}",
265
1
                    data.block_offset(), data.block_size());
266
1
        }
267
59
        const size_t block_off = static_cast<size_t>(data.block_offset());
268
59
        const size_t payload_size =
269
59
                use_attachment ? static_cast<size_t>(data.block_size()) : data.data().size();
270
59
        if (payload_size == 0) {
271
1
            continue;
272
1
        }
273
58
        const auto response_range = FileBlock::Range(block_off, block_off + payload_size - 1);
274
58
        const int expected_idx = find_expected_range_idx(expected.expected_ranges, response_range);
275
58
        if (expected_idx < 0) {
276
1
            peer_cache_reader_failed_counter << 1;
277
1
            result->clear();
278
1
            return Status::InternalError<false>(
279
1
                    "peer cache read block out of requested ranges: off={}, size={}", block_off,
280
1
                    payload_size);
281
1
        }
282
        // Attachment payload is a single byte stream. Consume it in resp.datas() order so the
283
        // peer can split or reorder requested ranges without forcing a fallback.
284
57
        if (!subtract_pending_range(expected.pending_ranges, response_range)) {
285
1
            peer_cache_reader_failed_counter << 1;
286
1
            result->clear();
287
1
            return Status::InternalError<false>("peer cache read unexpected block range: [{}, {}]",
288
1
                                                response_range.left, response_range.right);
289
1
        }
290
291
56
        PeerFetchChunk chunk;
292
56
        chunk.block_index = expected.expected_block_indexes[expected_idx];
293
56
        chunk.block_offset = response_range.left;
294
56
        VLOG_DEBUG << "peer cache read data=" << data.block_offset() << " size=" << payload_size
295
0
                   << " block_idx=" << chunk.block_index;
296
56
        if (use_attachment) {
297
18
            auto cut_st =
298
18
                    cut_attachment_payload(&remaining_attachment, payload_size, &chunk.payload);
299
18
            if (!cut_st.ok()) {
300
1
                peer_cache_reader_failed_counter << 1;
301
1
                result->clear();
302
1
                return cut_st;
303
1
            }
304
38
        } else if (chunk.payload.append(data.data().data(), payload_size) != 0) {
305
0
            peer_cache_reader_failed_counter << 1;
306
0
            result->clear();
307
0
            return Status::InternalError<false>(
308
0
                    "failed to append protobuf payload into iobuf: size={}", payload_size);
309
0
        }
310
55
        filled += payload_size;
311
55
        result->chunks.emplace_back(std::move(chunk));
312
55
    }
313
24
    VLOG_DEBUG << "peer cache read filled=" << filled;
314
    // Sparse reads are complete only when all requested block ranges are covered exactly.
315
24
    if (!expected.pending_ranges.empty() || filled != expected.expected_bytes ||
316
24
        (use_attachment && !remaining_attachment.empty())) {
317
2
        peer_cache_reader_failed_counter << 1;
318
2
        result->clear();
319
2
        return Status::InternalError<false>(
320
2
                "peer cache read incomplete: need={}, got={}, attachment_left={}",
321
2
                expected.expected_bytes, filled, remaining_attachment.size());
322
2
    }
323
22
    peer_bytes_read_total << filled;
324
22
    peer_bytes_per_read << filled;
325
22
    peer_cache_reader_succ_counter << 1;
326
22
    result->bytes_read = filled;
327
22
    return Status::OK();
328
24
}
329
330
} // namespace doris::io