/src/rocksdb/trace_replay/trace_replay.cc
Line | Count | Source |
1 | | // Copyright (c) 2011-present, Facebook, Inc. All rights reserved. |
2 | | // This source code is licensed under both the GPLv2 (found in the |
3 | | // COPYING file in the root directory) and Apache 2.0 License |
4 | | // (found in the LICENSE.Apache file in the root directory). |
5 | | |
6 | | #include "trace_replay/trace_replay.h" |
7 | | |
8 | | #include <chrono> |
9 | | #include <sstream> |
10 | | #include <thread> |
11 | | |
12 | | #include "db/db_impl/db_impl.h" |
13 | | #include "rocksdb/env.h" |
14 | | #include "rocksdb/iterator.h" |
15 | | #include "rocksdb/options.h" |
16 | | #include "rocksdb/slice.h" |
17 | | #include "rocksdb/system_clock.h" |
18 | | #include "rocksdb/trace_reader_writer.h" |
19 | | #include "rocksdb/write_batch.h" |
20 | | #include "util/coding.h" |
21 | | #include "util/string_util.h" |
22 | | |
23 | | namespace ROCKSDB_NAMESPACE { |
24 | | |
25 | | const std::string kTraceMagic = "feedcafedeadbeef"; |
26 | | |
27 | | namespace { |
28 | 0 | void DecodeCFAndKey(std::string& buffer, uint32_t* cf_id, Slice* key) { |
29 | 0 | Slice buf(buffer); |
30 | 0 | GetFixed32(&buf, cf_id); |
31 | 0 | GetLengthPrefixedSlice(&buf, key); |
32 | 0 | } |
33 | | } // namespace |
34 | | |
35 | 0 | Status TracerHelper::ParseVersionStr(std::string& v_string, int* v_num) { |
36 | 0 | if (v_string.find_first_of('.') == std::string::npos || |
37 | 0 | v_string.find_first_of('.') != v_string.find_last_of('.')) { |
38 | 0 | return Status::Corruption( |
39 | 0 | "Corrupted trace file. Incorrect version format."); |
40 | 0 | } |
41 | 0 | int tmp_num = 0; |
42 | 0 | for (int i = 0; i < static_cast<int>(v_string.size()); i++) { |
43 | 0 | if (v_string[i] == '.') { |
44 | 0 | continue; |
45 | 0 | } else if (isdigit(v_string[i])) { |
46 | 0 | tmp_num = tmp_num * 10 + (v_string[i] - '0'); |
47 | 0 | } else { |
48 | 0 | return Status::Corruption( |
49 | 0 | "Corrupted trace file. Incorrect version format"); |
50 | 0 | } |
51 | 0 | } |
52 | 0 | *v_num = tmp_num; |
53 | 0 | return Status::OK(); |
54 | 0 | } |
55 | | |
56 | | Status TracerHelper::ParseTraceHeader(const Trace& header, int* trace_version, |
57 | 0 | int* db_version) { |
58 | 0 | std::vector<std::string> s_vec; |
59 | 0 | int begin = 0, end; |
60 | 0 | for (int i = 0; i < 3; i++) { |
61 | 0 | assert(header.payload.find('\t', begin) != std::string::npos); |
62 | 0 | end = static_cast<int>(header.payload.find('\t', begin)); |
63 | 0 | s_vec.push_back(header.payload.substr(begin, end - begin)); |
64 | 0 | begin = end + 1; |
65 | 0 | } |
66 | |
|
67 | 0 | std::string t_v_str, db_v_str; |
68 | 0 | assert(s_vec.size() == 3); |
69 | 0 | assert(s_vec[1].find("Trace Version: ") != std::string::npos); |
70 | 0 | t_v_str = s_vec[1].substr(15); |
71 | 0 | assert(s_vec[2].find("RocksDB Version: ") != std::string::npos); |
72 | 0 | db_v_str = s_vec[2].substr(17); |
73 | |
|
74 | 0 | Status s; |
75 | 0 | s = ParseVersionStr(t_v_str, trace_version); |
76 | 0 | if (s != Status::OK()) { |
77 | 0 | return s; |
78 | 0 | } |
79 | 0 | s = ParseVersionStr(db_v_str, db_version); |
80 | 0 | return s; |
81 | 0 | } |
82 | | |
83 | 0 | void TracerHelper::EncodeTrace(const Trace& trace, std::string* encoded_trace) { |
84 | 0 | assert(encoded_trace); |
85 | 0 | PutFixed64(encoded_trace, trace.ts); |
86 | 0 | encoded_trace->push_back(trace.type); |
87 | 0 | PutFixed32(encoded_trace, static_cast<uint32_t>(trace.payload.size())); |
88 | 0 | encoded_trace->append(trace.payload); |
89 | 0 | } |
90 | | |
91 | | Status TracerHelper::DecodeTrace(const std::string& encoded_trace, |
92 | 0 | Trace* trace) { |
93 | 0 | assert(trace != nullptr); |
94 | 0 | Slice enc_slice = Slice(encoded_trace); |
95 | 0 | if (!GetFixed64(&enc_slice, &trace->ts)) { |
96 | 0 | return Status::Incomplete("Decode trace string failed"); |
97 | 0 | } |
98 | 0 | if (enc_slice.size() < kTraceTypeSize + kTracePayloadLengthSize) { |
99 | 0 | return Status::Incomplete("Decode trace string failed"); |
100 | 0 | } |
101 | 0 | trace->type = static_cast<TraceType>(enc_slice[0]); |
102 | 0 | enc_slice.remove_prefix(kTraceTypeSize + kTracePayloadLengthSize); |
103 | 0 | trace->payload = enc_slice.ToString(); |
104 | 0 | return Status::OK(); |
105 | 0 | } |
106 | | |
107 | | Status TracerHelper::DecodeHeader(const std::string& encoded_trace, |
108 | 0 | Trace* header) { |
109 | 0 | Status s = TracerHelper::DecodeTrace(encoded_trace, header); |
110 | |
|
111 | 0 | if (header->type != kTraceBegin) { |
112 | 0 | return Status::Corruption("Corrupted trace file. Incorrect header."); |
113 | 0 | } |
114 | 0 | if (header->payload.substr(0, kTraceMagic.length()) != kTraceMagic) { |
115 | 0 | return Status::Corruption("Corrupted trace file. Incorrect magic."); |
116 | 0 | } |
117 | | |
118 | 0 | return s; |
119 | 0 | } |
120 | | |
121 | | bool TracerHelper::SetPayloadMap(uint64_t& payload_map, |
122 | 0 | const TracePayloadType payload_type) { |
123 | 0 | uint64_t old_state = payload_map; |
124 | 0 | uint64_t tmp = 1; |
125 | 0 | payload_map |= (tmp << payload_type); |
126 | 0 | return old_state != payload_map; |
127 | 0 | } |
128 | | |
129 | | Status TracerHelper::DecodeTraceRecord(Trace* trace, int trace_file_version, |
130 | 0 | std::unique_ptr<TraceRecord>* record) { |
131 | 0 | assert(trace != nullptr); |
132 | |
|
133 | 0 | if (record != nullptr) { |
134 | 0 | record->reset(nullptr); |
135 | 0 | } |
136 | |
|
137 | 0 | switch (trace->type) { |
138 | | // Write |
139 | 0 | case kTraceWrite: { |
140 | 0 | PinnableSlice rep; |
141 | 0 | if (trace_file_version < 2) { |
142 | 0 | rep.PinSelf(trace->payload); |
143 | 0 | } else { |
144 | 0 | Slice buf(trace->payload); |
145 | 0 | GetFixed64(&buf, &trace->payload_map); |
146 | 0 | int64_t payload_map = static_cast<int64_t>(trace->payload_map); |
147 | 0 | Slice write_batch_data; |
148 | 0 | while (payload_map) { |
149 | | // Find the rightmost set bit. |
150 | 0 | uint32_t set_pos = |
151 | 0 | static_cast<uint32_t>(log2(payload_map & -payload_map)); |
152 | 0 | switch (set_pos) { |
153 | 0 | case TracePayloadType::kWriteBatchData: { |
154 | 0 | GetLengthPrefixedSlice(&buf, &write_batch_data); |
155 | 0 | break; |
156 | 0 | } |
157 | 0 | default: { |
158 | 0 | assert(false); |
159 | 0 | } |
160 | 0 | } |
161 | | // unset the rightmost bit. |
162 | 0 | payload_map &= (payload_map - 1); |
163 | 0 | } |
164 | 0 | rep.PinSelf(write_batch_data); |
165 | 0 | } |
166 | | |
167 | 0 | if (record != nullptr) { |
168 | 0 | record->reset(new WriteQueryTraceRecord(std::move(rep), trace->ts)); |
169 | 0 | } |
170 | |
|
171 | 0 | return Status::OK(); |
172 | 0 | } |
173 | | // Get |
174 | 0 | case kTraceGet: { |
175 | 0 | uint32_t cf_id = 0; |
176 | 0 | Slice get_key; |
177 | |
|
178 | 0 | if (trace_file_version < 2) { |
179 | 0 | DecodeCFAndKey(trace->payload, &cf_id, &get_key); |
180 | 0 | } else { |
181 | 0 | Slice buf(trace->payload); |
182 | 0 | GetFixed64(&buf, &trace->payload_map); |
183 | 0 | int64_t payload_map = static_cast<int64_t>(trace->payload_map); |
184 | 0 | while (payload_map) { |
185 | | // Find the rightmost set bit. |
186 | 0 | uint32_t set_pos = |
187 | 0 | static_cast<uint32_t>(log2(payload_map & -payload_map)); |
188 | 0 | switch (set_pos) { |
189 | 0 | case TracePayloadType::kGetCFID: { |
190 | 0 | GetFixed32(&buf, &cf_id); |
191 | 0 | break; |
192 | 0 | } |
193 | 0 | case TracePayloadType::kGetKey: { |
194 | 0 | GetLengthPrefixedSlice(&buf, &get_key); |
195 | 0 | break; |
196 | 0 | } |
197 | 0 | default: { |
198 | 0 | assert(false); |
199 | 0 | } |
200 | 0 | } |
201 | | // unset the rightmost bit. |
202 | 0 | payload_map &= (payload_map - 1); |
203 | 0 | } |
204 | 0 | } |
205 | | |
206 | 0 | if (record != nullptr) { |
207 | 0 | PinnableSlice ps; |
208 | 0 | ps.PinSelf(get_key); |
209 | 0 | record->reset(new GetQueryTraceRecord(cf_id, std::move(ps), trace->ts)); |
210 | 0 | } |
211 | |
|
212 | 0 | return Status::OK(); |
213 | 0 | } |
214 | | // Iterator Seek and SeekForPrev |
215 | 0 | case kTraceIteratorSeek: |
216 | 0 | case kTraceIteratorSeekForPrev: { |
217 | 0 | uint32_t cf_id = 0; |
218 | 0 | Slice iter_key; |
219 | 0 | Slice lower_bound; |
220 | 0 | Slice upper_bound; |
221 | |
|
222 | 0 | if (trace_file_version < 2) { |
223 | 0 | DecodeCFAndKey(trace->payload, &cf_id, &iter_key); |
224 | 0 | } else { |
225 | 0 | Slice buf(trace->payload); |
226 | 0 | GetFixed64(&buf, &trace->payload_map); |
227 | 0 | int64_t payload_map = static_cast<int64_t>(trace->payload_map); |
228 | 0 | while (payload_map) { |
229 | | // Find the rightmost set bit. |
230 | 0 | uint32_t set_pos = |
231 | 0 | static_cast<uint32_t>(log2(payload_map & -payload_map)); |
232 | 0 | switch (set_pos) { |
233 | 0 | case TracePayloadType::kIterCFID: { |
234 | 0 | GetFixed32(&buf, &cf_id); |
235 | 0 | break; |
236 | 0 | } |
237 | 0 | case TracePayloadType::kIterKey: { |
238 | 0 | GetLengthPrefixedSlice(&buf, &iter_key); |
239 | 0 | break; |
240 | 0 | } |
241 | 0 | case TracePayloadType::kIterLowerBound: { |
242 | 0 | GetLengthPrefixedSlice(&buf, &lower_bound); |
243 | 0 | break; |
244 | 0 | } |
245 | 0 | case TracePayloadType::kIterUpperBound: { |
246 | 0 | GetLengthPrefixedSlice(&buf, &upper_bound); |
247 | 0 | break; |
248 | 0 | } |
249 | 0 | default: { |
250 | 0 | assert(false); |
251 | 0 | } |
252 | 0 | } |
253 | | // unset the rightmost bit. |
254 | 0 | payload_map &= (payload_map - 1); |
255 | 0 | } |
256 | 0 | } |
257 | | |
258 | 0 | if (record != nullptr) { |
259 | 0 | PinnableSlice ps_key; |
260 | 0 | ps_key.PinSelf(iter_key); |
261 | 0 | PinnableSlice ps_lower; |
262 | 0 | ps_lower.PinSelf(lower_bound); |
263 | 0 | PinnableSlice ps_upper; |
264 | 0 | ps_upper.PinSelf(upper_bound); |
265 | 0 | record->reset(new IteratorSeekQueryTraceRecord( |
266 | 0 | static_cast<IteratorSeekQueryTraceRecord::SeekType>(trace->type), |
267 | 0 | cf_id, std::move(ps_key), std::move(ps_lower), std::move(ps_upper), |
268 | 0 | trace->ts)); |
269 | 0 | } |
270 | |
|
271 | 0 | return Status::OK(); |
272 | 0 | } |
273 | | // MultiGet |
274 | 0 | case kTraceMultiGet: { |
275 | 0 | if (trace_file_version < 2) { |
276 | 0 | return Status::Corruption("MultiGet is not supported."); |
277 | 0 | } |
278 | | |
279 | 0 | uint32_t multiget_size = 0; |
280 | 0 | std::vector<uint32_t> cf_ids; |
281 | 0 | std::vector<PinnableSlice> multiget_keys; |
282 | |
|
283 | 0 | Slice cfids_payload; |
284 | 0 | Slice keys_payload; |
285 | 0 | Slice buf(trace->payload); |
286 | 0 | GetFixed64(&buf, &trace->payload_map); |
287 | 0 | int64_t payload_map = static_cast<int64_t>(trace->payload_map); |
288 | 0 | while (payload_map) { |
289 | | // Find the rightmost set bit. |
290 | 0 | uint32_t set_pos = |
291 | 0 | static_cast<uint32_t>(log2(payload_map & -payload_map)); |
292 | 0 | switch (set_pos) { |
293 | 0 | case TracePayloadType::kMultiGetSize: { |
294 | 0 | GetFixed32(&buf, &multiget_size); |
295 | 0 | break; |
296 | 0 | } |
297 | 0 | case TracePayloadType::kMultiGetCFIDs: { |
298 | 0 | GetLengthPrefixedSlice(&buf, &cfids_payload); |
299 | 0 | break; |
300 | 0 | } |
301 | 0 | case TracePayloadType::kMultiGetKeys: { |
302 | 0 | GetLengthPrefixedSlice(&buf, &keys_payload); |
303 | 0 | break; |
304 | 0 | } |
305 | 0 | default: { |
306 | 0 | assert(false); |
307 | 0 | } |
308 | 0 | } |
309 | | // unset the rightmost bit. |
310 | 0 | payload_map &= (payload_map - 1); |
311 | 0 | } |
312 | 0 | if (multiget_size == 0) { |
313 | 0 | return Status::InvalidArgument("Empty MultiGet cf_ids or keys."); |
314 | 0 | } |
315 | | |
316 | | // Decode the cfids_payload and keys_payload |
317 | 0 | cf_ids.reserve(multiget_size); |
318 | 0 | multiget_keys.reserve(multiget_size); |
319 | 0 | for (uint32_t i = 0; i < multiget_size; i++) { |
320 | 0 | uint32_t tmp_cfid = 0; |
321 | 0 | Slice tmp_key; |
322 | 0 | GetFixed32(&cfids_payload, &tmp_cfid); |
323 | 0 | GetLengthPrefixedSlice(&keys_payload, &tmp_key); |
324 | 0 | cf_ids.push_back(tmp_cfid); |
325 | 0 | Slice s(tmp_key); |
326 | 0 | PinnableSlice ps; |
327 | 0 | ps.PinSelf(s); |
328 | 0 | multiget_keys.push_back(std::move(ps)); |
329 | 0 | } |
330 | |
|
331 | 0 | if (record != nullptr) { |
332 | 0 | record->reset(new MultiGetQueryTraceRecord( |
333 | 0 | std::move(cf_ids), std::move(multiget_keys), trace->ts)); |
334 | 0 | } |
335 | |
|
336 | 0 | return Status::OK(); |
337 | 0 | } |
338 | 0 | default: |
339 | 0 | return Status::NotSupported("Unsupported trace type."); |
340 | 0 | } |
341 | 0 | } |
342 | | |
343 | | Tracer::Tracer(SystemClock* clock, const TraceOptions& trace_options, |
344 | | std::unique_ptr<TraceWriter>&& trace_writer) |
345 | 0 | : clock_(clock), |
346 | 0 | trace_options_(trace_options), |
347 | 0 | trace_writer_(std::move(trace_writer)), |
348 | 0 | trace_request_count_(0), |
349 | 0 | trace_write_status_(Status::OK()) { |
350 | | // TODO: What if this fails? |
351 | 0 | WriteHeader().PermitUncheckedError(); |
352 | 0 | } |
353 | | |
354 | 0 | Tracer::~Tracer() { trace_writer_.reset(); } |
355 | | |
356 | 0 | Status Tracer::Write(WriteBatch* write_batch) { |
357 | 0 | TraceType trace_type = kTraceWrite; |
358 | 0 | if (ShouldSkipTrace(trace_type)) { |
359 | 0 | return Status::OK(); |
360 | 0 | } |
361 | 0 | Trace trace; |
362 | 0 | trace.ts = clock_->NowMicros(); |
363 | 0 | trace.type = trace_type; |
364 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, |
365 | 0 | TracePayloadType::kWriteBatchData); |
366 | 0 | PutFixed64(&trace.payload, trace.payload_map); |
367 | 0 | PutLengthPrefixedSlice(&trace.payload, Slice(write_batch->Data())); |
368 | 0 | return WriteTrace(trace); |
369 | 0 | } |
370 | | |
371 | 0 | Status Tracer::Get(ColumnFamilyHandle* column_family, const Slice& key) { |
372 | 0 | TraceType trace_type = kTraceGet; |
373 | 0 | if (ShouldSkipTrace(trace_type)) { |
374 | 0 | return Status::OK(); |
375 | 0 | } |
376 | 0 | Trace trace; |
377 | 0 | trace.ts = clock_->NowMicros(); |
378 | 0 | trace.type = trace_type; |
379 | | // Set the payloadmap of the struct member that will be encoded in the |
380 | | // payload. |
381 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, TracePayloadType::kGetCFID); |
382 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, TracePayloadType::kGetKey); |
383 | | // Encode the Get struct members into payload. Make sure add them in order. |
384 | 0 | PutFixed64(&trace.payload, trace.payload_map); |
385 | 0 | PutFixed32(&trace.payload, column_family->GetID()); |
386 | 0 | PutLengthPrefixedSlice(&trace.payload, key); |
387 | 0 | return WriteTrace(trace); |
388 | 0 | } |
389 | | |
390 | | Status Tracer::IteratorSeek(const uint32_t& cf_id, const Slice& key, |
391 | 0 | const Slice& lower_bound, const Slice upper_bound) { |
392 | 0 | TraceType trace_type = kTraceIteratorSeek; |
393 | 0 | if (ShouldSkipTrace(trace_type)) { |
394 | 0 | return Status::OK(); |
395 | 0 | } |
396 | 0 | Trace trace; |
397 | 0 | trace.ts = clock_->NowMicros(); |
398 | 0 | trace.type = trace_type; |
399 | | // Set the payloadmap of the struct member that will be encoded in the |
400 | | // payload. |
401 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, TracePayloadType::kIterCFID); |
402 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, TracePayloadType::kIterKey); |
403 | 0 | if (lower_bound.size() > 0) { |
404 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, |
405 | 0 | TracePayloadType::kIterLowerBound); |
406 | 0 | } |
407 | 0 | if (upper_bound.size() > 0) { |
408 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, |
409 | 0 | TracePayloadType::kIterUpperBound); |
410 | 0 | } |
411 | | // Encode the Iterator struct members into payload. Make sure add them in |
412 | | // order. |
413 | 0 | PutFixed64(&trace.payload, trace.payload_map); |
414 | 0 | PutFixed32(&trace.payload, cf_id); |
415 | 0 | PutLengthPrefixedSlice(&trace.payload, key); |
416 | 0 | if (lower_bound.size() > 0) { |
417 | 0 | PutLengthPrefixedSlice(&trace.payload, lower_bound); |
418 | 0 | } |
419 | 0 | if (upper_bound.size() > 0) { |
420 | 0 | PutLengthPrefixedSlice(&trace.payload, upper_bound); |
421 | 0 | } |
422 | 0 | return WriteTrace(trace); |
423 | 0 | } |
424 | | |
425 | | Status Tracer::IteratorSeekForPrev(const uint32_t& cf_id, const Slice& key, |
426 | | const Slice& lower_bound, |
427 | 0 | const Slice upper_bound) { |
428 | 0 | TraceType trace_type = kTraceIteratorSeekForPrev; |
429 | 0 | if (ShouldSkipTrace(trace_type)) { |
430 | 0 | return Status::OK(); |
431 | 0 | } |
432 | 0 | Trace trace; |
433 | 0 | trace.ts = clock_->NowMicros(); |
434 | 0 | trace.type = trace_type; |
435 | | // Set the payloadmap of the struct member that will be encoded in the |
436 | | // payload. |
437 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, TracePayloadType::kIterCFID); |
438 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, TracePayloadType::kIterKey); |
439 | 0 | if (lower_bound.size() > 0) { |
440 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, |
441 | 0 | TracePayloadType::kIterLowerBound); |
442 | 0 | } |
443 | 0 | if (upper_bound.size() > 0) { |
444 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, |
445 | 0 | TracePayloadType::kIterUpperBound); |
446 | 0 | } |
447 | | // Encode the Iterator struct members into payload. Make sure add them in |
448 | | // order. |
449 | 0 | PutFixed64(&trace.payload, trace.payload_map); |
450 | 0 | PutFixed32(&trace.payload, cf_id); |
451 | 0 | PutLengthPrefixedSlice(&trace.payload, key); |
452 | 0 | if (lower_bound.size() > 0) { |
453 | 0 | PutLengthPrefixedSlice(&trace.payload, lower_bound); |
454 | 0 | } |
455 | 0 | if (upper_bound.size() > 0) { |
456 | 0 | PutLengthPrefixedSlice(&trace.payload, upper_bound); |
457 | 0 | } |
458 | 0 | return WriteTrace(trace); |
459 | 0 | } |
460 | | |
461 | | Status Tracer::MultiGet(const size_t num_keys, |
462 | | ColumnFamilyHandle** column_families, |
463 | 0 | const Slice* keys) { |
464 | 0 | if (num_keys == 0) { |
465 | 0 | return Status::OK(); |
466 | 0 | } |
467 | 0 | std::vector<ColumnFamilyHandle*> v_column_families; |
468 | 0 | std::vector<Slice> v_keys; |
469 | 0 | v_column_families.resize(num_keys); |
470 | 0 | v_keys.resize(num_keys); |
471 | 0 | for (size_t i = 0; i < num_keys; i++) { |
472 | 0 | v_column_families[i] = column_families[i]; |
473 | 0 | v_keys[i] = keys[i]; |
474 | 0 | } |
475 | 0 | return MultiGet(v_column_families, v_keys); |
476 | 0 | } |
477 | | |
478 | | Status Tracer::MultiGet(const size_t num_keys, |
479 | 0 | ColumnFamilyHandle* column_family, const Slice* keys) { |
480 | 0 | if (num_keys == 0) { |
481 | 0 | return Status::OK(); |
482 | 0 | } |
483 | 0 | std::vector<ColumnFamilyHandle*> column_families; |
484 | 0 | std::vector<Slice> v_keys; |
485 | 0 | column_families.resize(num_keys); |
486 | 0 | v_keys.resize(num_keys); |
487 | 0 | for (size_t i = 0; i < num_keys; i++) { |
488 | 0 | column_families[i] = column_family; |
489 | 0 | v_keys[i] = keys[i]; |
490 | 0 | } |
491 | 0 | return MultiGet(column_families, v_keys); |
492 | 0 | } |
493 | | |
494 | | Status Tracer::MultiGet(const std::vector<ColumnFamilyHandle*>& column_families, |
495 | 0 | const std::vector<Slice>& keys) { |
496 | 0 | if (column_families.size() != keys.size()) { |
497 | 0 | return Status::Corruption("the CFs size and keys size does not match!"); |
498 | 0 | } |
499 | 0 | TraceType trace_type = kTraceMultiGet; |
500 | 0 | if (ShouldSkipTrace(trace_type)) { |
501 | 0 | return Status::OK(); |
502 | 0 | } |
503 | 0 | uint32_t multiget_size = static_cast<uint32_t>(keys.size()); |
504 | 0 | Trace trace; |
505 | 0 | trace.ts = clock_->NowMicros(); |
506 | 0 | trace.type = trace_type; |
507 | | // Set the payloadmap of the struct member that will be encoded in the |
508 | | // payload. |
509 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, |
510 | 0 | TracePayloadType::kMultiGetSize); |
511 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, |
512 | 0 | TracePayloadType::kMultiGetCFIDs); |
513 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, |
514 | 0 | TracePayloadType::kMultiGetKeys); |
515 | | // Encode the CFIDs inorder |
516 | 0 | std::string cfids_payload; |
517 | 0 | std::string keys_payload; |
518 | 0 | for (uint32_t i = 0; i < multiget_size; i++) { |
519 | 0 | assert(i < column_families.size()); |
520 | 0 | assert(i < keys.size()); |
521 | 0 | PutFixed32(&cfids_payload, column_families[i]->GetID()); |
522 | 0 | PutLengthPrefixedSlice(&keys_payload, keys[i]); |
523 | 0 | } |
524 | | // Encode the Get struct members into payload. Make sure add them in order. |
525 | 0 | PutFixed64(&trace.payload, trace.payload_map); |
526 | 0 | PutFixed32(&trace.payload, multiget_size); |
527 | 0 | PutLengthPrefixedSlice(&trace.payload, cfids_payload); |
528 | 0 | PutLengthPrefixedSlice(&trace.payload, keys_payload); |
529 | 0 | return WriteTrace(trace); |
530 | 0 | } |
531 | | |
532 | 0 | bool Tracer::ShouldSkipTrace(const TraceType& trace_type) { |
533 | 0 | if (IsTraceFileOverMax()) { |
534 | 0 | return true; |
535 | 0 | } |
536 | | |
537 | 0 | TraceFilterType filter_mask = kTraceFilterNone; |
538 | 0 | switch (trace_type) { |
539 | 0 | case kTraceNone: |
540 | 0 | case kTraceBegin: |
541 | 0 | case kTraceEnd: |
542 | 0 | filter_mask = kTraceFilterNone; |
543 | 0 | break; |
544 | 0 | case kTraceWrite: |
545 | 0 | filter_mask = kTraceFilterWrite; |
546 | 0 | break; |
547 | 0 | case kTraceGet: |
548 | 0 | filter_mask = kTraceFilterGet; |
549 | 0 | break; |
550 | 0 | case kTraceIteratorSeek: |
551 | 0 | filter_mask = kTraceFilterIteratorSeek; |
552 | 0 | break; |
553 | 0 | case kTraceIteratorSeekForPrev: |
554 | 0 | filter_mask = kTraceFilterIteratorSeekForPrev; |
555 | 0 | break; |
556 | 0 | case kBlockTraceIndexBlock: |
557 | 0 | case kBlockTraceFilterBlock: |
558 | 0 | case kBlockTraceDataBlock: |
559 | 0 | case kBlockTraceUncompressionDictBlock: |
560 | 0 | case kBlockTraceRangeDeletionBlock: |
561 | 0 | case kIOTracer: |
562 | 0 | filter_mask = kTraceFilterNone; |
563 | 0 | break; |
564 | 0 | case kTraceMultiGet: |
565 | 0 | filter_mask = kTraceFilterMultiGet; |
566 | 0 | break; |
567 | 0 | case kTraceMax: |
568 | 0 | assert(false); |
569 | 0 | filter_mask = kTraceFilterNone; |
570 | 0 | break; |
571 | 0 | } |
572 | 0 | if (filter_mask != kTraceFilterNone && trace_options_.filter & filter_mask) { |
573 | 0 | return true; |
574 | 0 | } |
575 | | |
576 | 0 | ++trace_request_count_; |
577 | 0 | if (trace_request_count_ < trace_options_.sampling_frequency) { |
578 | 0 | return true; |
579 | 0 | } |
580 | 0 | trace_request_count_ = 0; |
581 | 0 | return false; |
582 | 0 | } |
583 | | |
584 | 0 | bool Tracer::IsTraceFileOverMax() { |
585 | 0 | uint64_t trace_file_size = trace_writer_->GetFileSize(); |
586 | 0 | return (trace_file_size > trace_options_.max_trace_file_size); |
587 | 0 | } |
588 | | |
589 | 0 | Status Tracer::WriteHeader() { |
590 | 0 | std::ostringstream s; |
591 | 0 | s << kTraceMagic << "\t" << "Trace Version: " << kTraceFileMajorVersion << "." |
592 | 0 | << kTraceFileMinorVersion << "\t" << "RocksDB Version: " << kMajorVersion |
593 | 0 | << "." << kMinorVersion << "\t" << "Format: Timestamp OpType Payload\n"; |
594 | 0 | std::string header(s.str()); |
595 | |
|
596 | 0 | Trace trace; |
597 | 0 | trace.ts = clock_->NowMicros(); |
598 | 0 | trace.type = kTraceBegin; |
599 | 0 | trace.payload = header; |
600 | 0 | return WriteTrace(trace); |
601 | 0 | } |
602 | | |
603 | 0 | Status Tracer::WriteFooter() { |
604 | 0 | Trace trace; |
605 | 0 | trace.ts = clock_->NowMicros(); |
606 | 0 | trace.type = kTraceEnd; |
607 | 0 | TracerHelper::SetPayloadMap(trace.payload_map, |
608 | 0 | TracePayloadType::kEmptyPayload); |
609 | 0 | trace.payload = ""; |
610 | 0 | return WriteTrace(trace); |
611 | 0 | } |
612 | | |
613 | 0 | Status Tracer::WriteTrace(const Trace& trace) { |
614 | 0 | if (!trace_write_status_.ok()) { |
615 | 0 | return Status::Incomplete("Tracing has seen error: %s", |
616 | 0 | trace_write_status_.ToString()); |
617 | 0 | } |
618 | 0 | assert(trace_write_status_.ok()); |
619 | 0 | std::string encoded_trace; |
620 | 0 | TracerHelper::EncodeTrace(trace, &encoded_trace); |
621 | 0 | Status s = trace_writer_->Write(Slice(encoded_trace)); |
622 | 0 | if (!s.ok()) { |
623 | 0 | trace_write_status_ = s; |
624 | 0 | } |
625 | 0 | return s; |
626 | 0 | } |
627 | | |
628 | 0 | Status Tracer::Close() { return WriteFooter(); } |
629 | | |
630 | | } // namespace ROCKSDB_NAMESPACE |