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 |