/src/rocksdb/db/db_iter.h
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 | | // Copyright (c) 2011 The LevelDB Authors. All rights reserved. |
7 | | // Use of this source code is governed by a BSD-style license that can be |
8 | | // found in the LICENSE file. See the AUTHORS file for names of contributors. |
9 | | |
10 | | #pragma once |
11 | | #include <cstdint> |
12 | | #include <memory> |
13 | | #include <mutex> |
14 | | #include <string> |
15 | | #include <utility> |
16 | | #include <vector> |
17 | | |
18 | | #include "db/blob/blob_fetcher.h" |
19 | | #include "db/blob/blob_index.h" |
20 | | #include "db/db_impl/db_impl.h" |
21 | | #include "db/wide/read_path_blob_resolver.h" |
22 | | #include "db/wide/wide_columns_helper.h" |
23 | | #include "memory/arena.h" |
24 | | #include "options/cf_options.h" |
25 | | #include "rocksdb/db.h" |
26 | | #include "rocksdb/iterator.h" |
27 | | #include "rocksdb/wide_columns.h" |
28 | | #include "table/iterator_wrapper.h" |
29 | | #include "util/autovector.h" |
30 | | #include "util/dirty_tracked.h" |
31 | | |
32 | | namespace ROCKSDB_NAMESPACE { |
33 | | class BlobFileCache; |
34 | | class Version; |
35 | | namespace port { |
36 | | class RWMutex; |
37 | | } |
38 | | |
39 | | // This file declares the factory functions of DBIter, in its original form |
40 | | // or a wrapped form with class ArenaWrappedDBIter, which is defined here. |
41 | | // Class DBIter, which is declared and implemented inside db_iter.cc, is |
42 | | // an iterator that converts internal keys (yielded by an InternalIterator) |
43 | | // that were live at the specified sequence number into appropriate user |
44 | | // keys. |
45 | | // Each internal key consists of a user key, a sequence number, and a value |
46 | | // type. DBIter deals with multiple key versions, tombstones, merge operands, |
47 | | // etc, and exposes an Iterator. |
48 | | // For example, DBIter may wrap following InternalIterator: |
49 | | // user key: AAA value: v3 seqno: 100 type: Put |
50 | | // user key: AAA value: v2 seqno: 97 type: Put |
51 | | // user key: AAA value: v1 seqno: 95 type: Put |
52 | | // user key: BBB value: v1 seqno: 90 type: Put |
53 | | // user key: BBC value: N/A seqno: 98 type: Delete |
54 | | // user key: BBC value: v1 seqno: 95 type: Put |
55 | | // If the snapshot passed in is 102, then the DBIter is expected to |
56 | | // expose the following iterator: |
57 | | // key: AAA value: v3 |
58 | | // key: BBB value: v1 |
59 | | // If the snapshot passed in is 96, then it should expose: |
60 | | // key: AAA value: v1 |
61 | | // key: BBB value: v1 |
62 | | // key: BBC value: v1 |
63 | | // |
64 | | |
65 | | // Memtables and sstables that make the DB representation contain |
66 | | // (userkey,seq,type) => uservalue entries. DBIter |
67 | | // combines multiple entries for the same userkey found in the DB |
68 | | // representation into a single entry while accounting for sequence |
69 | | // numbers, deletion markers, overwrites, etc. |
70 | | class DBIter final : public Iterator { |
71 | | public: |
72 | | // Return a new DBIter that reads from `internal_iter` at the specified |
73 | | // `sequence` number. |
74 | | // |
75 | | // @param active_mem Pointer to the active memtable that `internal_iter` |
76 | | // is reading from. If not null, the memtable can be marked for flush |
77 | | // according to options mutable_cf_options.memtable_op_scan_flush_trigger |
78 | | // and mutable_cf_options.memtable_avg_op_scan_flush_trigger. |
79 | | // @param arena_mode If true, the DBIter will be allocated from the arena. |
80 | | static DBIter* NewIter(Env* env, const ReadOptions& read_options, |
81 | | const ImmutableOptions& ioptions, |
82 | | const MutableCFOptions& mutable_cf_options, |
83 | | const Comparator* user_key_comparator, |
84 | | InternalIterator* internal_iter, |
85 | | const Version* version, const SequenceNumber& sequence, |
86 | | ReadCallback* read_callback, |
87 | | ReadOnlyMemTable* active_mem, |
88 | | ColumnFamilyHandleImpl* cfh = nullptr, |
89 | | bool expose_blob_index = false, Arena* arena = nullptr, |
90 | | DBImpl* db_impl = nullptr, |
91 | 15.1k | ColumnFamilyData* cfd = nullptr) { |
92 | 15.1k | if (cfh != nullptr) { |
93 | 0 | db_impl = cfh->db(); |
94 | 0 | cfd = cfh->cfd(); |
95 | 0 | } |
96 | 15.1k | void* mem = arena ? arena->AllocateAligned(sizeof(DBIter)) |
97 | 15.1k | : operator new(sizeof(DBIter)); |
98 | 15.1k | DBIter* db_iter = new (mem) |
99 | 15.1k | DBIter(env, read_options, ioptions, mutable_cf_options, |
100 | 15.1k | user_key_comparator, internal_iter, version, sequence, arena, |
101 | 15.1k | read_callback, db_impl, cfd, expose_blob_index, active_mem); |
102 | 15.1k | return db_iter; |
103 | 15.1k | } |
104 | | |
105 | | // The following is grossly complicated. TODO: clean it up |
106 | | // Which direction is the iterator currently moving? |
107 | | // (1) When moving forward: |
108 | | // (1a) if current_entry_is_merged_ = false, the internal iterator is |
109 | | // positioned at the exact entry that yields this->key(), this->value() |
110 | | // (1b) if current_entry_is_merged_ = true, the internal iterator is |
111 | | // positioned immediately after the last entry that contributed to the |
112 | | // current this->value(). That entry may or may not have key equal to |
113 | | // this->key(). |
114 | | // (2) When moving backwards, the internal iterator is positioned |
115 | | // just before all entries whose user key == this->key(). |
116 | | enum Direction : uint8_t { kForward, kReverse }; |
117 | | |
118 | | // LocalStatistics contain Statistics counters that will be aggregated per |
119 | | // each iterator instance and then will be sent to the global statistics when |
120 | | // the iterator is destroyed. |
121 | | // |
122 | | // The purpose of this approach is to avoid perf regression happening |
123 | | // when multiple threads bump the atomic counters from a DBIter::Next(). |
124 | | struct LocalStatistics { |
125 | 15.1k | explicit LocalStatistics() { ResetCounters(); } |
126 | | |
127 | 30.3k | void ResetCounters() { |
128 | 30.3k | next_count_ = 0; |
129 | 30.3k | next_found_count_ = 0; |
130 | 30.3k | prev_count_ = 0; |
131 | 30.3k | prev_found_count_ = 0; |
132 | 30.3k | bytes_read_ = 0; |
133 | 30.3k | skip_count_ = 0; |
134 | 30.3k | } |
135 | | |
136 | 15.1k | void BumpGlobalStatistics(Statistics* global_statistics) { |
137 | 15.1k | RecordTick(global_statistics, NUMBER_DB_NEXT, next_count_); |
138 | 15.1k | RecordTick(global_statistics, NUMBER_DB_NEXT_FOUND, next_found_count_); |
139 | 15.1k | RecordTick(global_statistics, NUMBER_DB_PREV, prev_count_); |
140 | 15.1k | RecordTick(global_statistics, NUMBER_DB_PREV_FOUND, prev_found_count_); |
141 | 15.1k | RecordTick(global_statistics, ITER_BYTES_READ, bytes_read_); |
142 | 15.1k | RecordTick(global_statistics, NUMBER_ITER_SKIP, skip_count_); |
143 | 15.1k | PERF_COUNTER_ADD(iter_read_bytes, bytes_read_); |
144 | 15.1k | ResetCounters(); |
145 | 15.1k | } |
146 | | |
147 | | // Map to Tickers::NUMBER_DB_NEXT |
148 | | uint64_t next_count_; |
149 | | // Map to Tickers::NUMBER_DB_NEXT_FOUND |
150 | | uint64_t next_found_count_; |
151 | | // Map to Tickers::NUMBER_DB_PREV |
152 | | uint64_t prev_count_; |
153 | | // Map to Tickers::NUMBER_DB_PREV_FOUND |
154 | | uint64_t prev_found_count_; |
155 | | // Map to Tickers::ITER_BYTES_READ |
156 | | uint64_t bytes_read_; |
157 | | // Map to Tickers::NUMBER_ITER_SKIP |
158 | | uint64_t skip_count_; |
159 | | }; |
160 | | |
161 | | // No copying allowed |
162 | | DBIter(const DBIter&) = delete; |
163 | | void operator=(const DBIter&) = delete; |
164 | | |
165 | 15.1k | ~DBIter() override { |
166 | 15.1k | MarkMemtableForFlushForAvgTrigger(); |
167 | 15.1k | ThreadStatus::OperationType cur_op_type = |
168 | 15.1k | ThreadStatusUtil::GetThreadOperation(); |
169 | 15.1k | ThreadStatusUtil::SetThreadOperation( |
170 | 15.1k | ThreadStatus::OperationType::OP_UNKNOWN); |
171 | | // Release pinned data if any |
172 | 15.1k | if (pinned_iters_mgr_.PinningEnabled()) { |
173 | 3.14k | pinned_iters_mgr_.ReleasePinnedData(); |
174 | 3.14k | } |
175 | 15.1k | RecordTick(statistics_, NO_ITERATOR_DELETED); |
176 | 15.1k | ResetInternalKeysSkippedCounter(); |
177 | 15.1k | local_stats_.BumpGlobalStatistics(statistics_); |
178 | 15.1k | iter_.DeleteIter(arena_mode_); |
179 | 15.1k | ThreadStatusUtil::SetThreadOperation(cur_op_type); |
180 | 15.1k | } |
181 | 11.3k | void SetIter(InternalIterator* iter) { |
182 | 11.3k | assert(iter_.iter() == nullptr); |
183 | 11.3k | iter_.Set(iter); |
184 | 11.3k | iter_.iter()->SetPinnedItersMgr(&pinned_iters_mgr_); |
185 | 11.3k | } |
186 | | |
187 | 45.0k | bool Valid() const override { |
188 | | #ifdef ROCKSDB_ASSERT_STATUS_CHECKED |
189 | | if (valid_) { |
190 | | status_.PermitUncheckedError(); |
191 | | } |
192 | | #endif // ROCKSDB_ASSERT_STATUS_CHECKED |
193 | 45.0k | return valid_; |
194 | 45.0k | } |
195 | 34.5k | Slice key() const override { |
196 | 34.5k | assert(valid_); |
197 | 34.5k | if (timestamp_lb_) { |
198 | 0 | return saved_key_.GetInternalKey(); |
199 | 34.5k | } else { |
200 | 34.5k | const Slice ukey_and_ts = saved_key_.GetUserKey(); |
201 | 34.5k | return Slice(ukey_and_ts.data(), ukey_and_ts.size() - timestamp_size_); |
202 | 34.5k | } |
203 | 34.5k | } |
204 | 34.5k | Slice value() const override { |
205 | 34.5k | assert(valid_); |
206 | | |
207 | 34.5k | return value_columns_state_->value(); |
208 | 34.5k | } |
209 | | |
210 | 0 | const WideColumns& columns() const override { |
211 | 0 | assert(valid_); |
212 | 0 | assert(!value_columns_state_->HasLazyEntityColumns() || |
213 | 0 | value_columns_state_->HasMaterializedColumns()); |
214 | 0 | return value_columns_state_->wide_columns(); |
215 | 0 | } |
216 | | |
217 | 0 | Status status() const override { |
218 | 0 | if (status_.ok() && iter_.iter() != nullptr) { |
219 | 0 | return iter_.status(); |
220 | 0 | } else { |
221 | 0 | assert(!valid_); |
222 | 0 | return status_; |
223 | 0 | } |
224 | 0 | } |
225 | 0 | Slice timestamp() const override { |
226 | 0 | assert(valid_); |
227 | 0 | assert(timestamp_size_ > 0); |
228 | 0 | if (direction_ == kReverse) { |
229 | 0 | return saved_timestamp_; |
230 | 0 | } |
231 | 0 | const Slice ukey_and_ts = saved_key_.GetUserKey(); |
232 | 0 | assert(timestamp_size_ < ukey_and_ts.size()); |
233 | 0 | return ExtractTimestampFromUserKey(ukey_and_ts, timestamp_size_); |
234 | 0 | } |
235 | 0 | bool IsBlob() const { |
236 | 0 | assert(valid_); |
237 | 0 | return blob_state_->is_blob; |
238 | 0 | } |
239 | | |
240 | | Status GetProperty(std::string prop_name, std::string* prop) override; |
241 | | |
242 | | void Next() final override; |
243 | | void Prev() final override; |
244 | | // 'target' does not contain timestamp, even if user timestamp feature is |
245 | | // enabled. |
246 | | void Seek(const Slice& target) final override; |
247 | | void SeekForPrev(const Slice& target) final override; |
248 | | void SeekToFirst() final override; |
249 | | void SeekToLast() final override; |
250 | 0 | Env* env() const { return env_; } |
251 | 0 | void set_sequence(uint64_t s) { |
252 | 0 | sequence_ = s; |
253 | 0 | if (read_callback_) { |
254 | 0 | read_callback_->Refresh(s); |
255 | 0 | } |
256 | 0 | iter_.SetRangeDelReadSeqno(s); |
257 | 0 | } |
258 | 0 | void set_valid(bool v) { valid_ = v; } |
259 | 0 | void set_status(Status s) { status_ = std::move(s); } |
260 | | |
261 | | bool PrepareValue() override; |
262 | | |
263 | | void Prepare(const MultiScanArgs& scan_opts) override; |
264 | | Status ValidateScanOptions(const MultiScanArgs& multiscan_opts) const; |
265 | | Status SetScanOptionsForPrepare(const MultiScanArgs& scan_opts); |
266 | | void PrepareInternalChildren(); |
267 | | |
268 | | private: |
269 | | DBIter(Env* _env, const ReadOptions& read_options, |
270 | | const ImmutableOptions& ioptions, |
271 | | const MutableCFOptions& mutable_cf_options, const Comparator* cmp, |
272 | | InternalIterator* iter, const Version* version, SequenceNumber s, |
273 | | bool arena_mode, ReadCallback* read_callback, DBImpl* db_impl, |
274 | | ColumnFamilyData* cfd, bool expose_blob_index, |
275 | | ReadOnlyMemTable* active_mem); |
276 | | |
277 | | class BlobReader { |
278 | | public: |
279 | | BlobReader(const Version* version, const ReadOptions& read_options, |
280 | | BlobFileCache* blob_file_cache, bool allow_write_path_fallback) |
281 | 15.1k | : blob_fetcher_(version, ReadOptions(read_options), blob_file_cache, |
282 | 15.1k | allow_write_path_fallback) {} |
283 | | |
284 | 0 | const Slice& GetBlobValue() const { return blob_value_; } |
285 | | Status RetrieveAndSetBlobValue(const Slice& user_key, |
286 | | const Slice& blob_index); |
287 | 0 | void ResetBlobValue() { blob_value_.Reset(); } |
288 | | // The blob fetcher backing this reader, for resolving wide-column entity |
289 | | // blob references (the merge path). Valid for this BlobReader's lifetime. |
290 | 0 | const BlobFetcher& blob_fetcher() const { return blob_fetcher_; } |
291 | | |
292 | | private: |
293 | | PinnableSlice blob_value_; |
294 | | OwningVersionBlobFetcher blob_fetcher_; |
295 | | }; |
296 | | struct BlobState { |
297 | | BlobReader reader; |
298 | | Slice lazy_blob_index; |
299 | | bool is_blob = false; |
300 | | |
301 | | template <typename... Args> |
302 | 15.1k | explicit BlobState(Args&&... args) : reader(std::forward<Args>(args)...) {} |
303 | | |
304 | 0 | void Reset() { |
305 | 0 | reader.ResetBlobValue(); |
306 | 0 | lazy_blob_index.clear(); |
307 | 0 | is_blob = false; |
308 | 0 | } |
309 | | }; |
310 | | |
311 | | // Groups the current iterator result together with the backing storage and |
312 | | // lazy entity-resolution metadata it depends on. Resetting this object drops |
313 | | // all aliases into the saved entity buffer and resolver cache at once. |
314 | | class ValueColumnsState { |
315 | | public: |
316 | | ValueColumnsState(const Version* version, const ReadOptions& read_options, |
317 | | ColumnFamilyData* cfd) |
318 | 15.1k | : entity_blob_resolver_( |
319 | 15.1k | version, read_options, cfd ? cfd->blob_file_cache() : nullptr, |
320 | 15.1k | cfd != nullptr && cfd->blob_partition_manager() != nullptr) {} |
321 | | |
322 | 0 | Slice& value() { return value_; } |
323 | 34.5k | const Slice& value() const { return value_; } |
324 | 0 | WideColumns& wide_columns() { return wide_columns_; } |
325 | 0 | const WideColumns& wide_columns() const { return wide_columns_; } |
326 | 0 | bool HasMaterializedColumns() const { return !wide_columns_.empty(); } |
327 | 0 | bool HasLazyEntityColumns() const { return !lazy_entity_columns_.empty(); } |
328 | | |
329 | 0 | std::string& saved_value() { return saved_value_; } |
330 | 0 | const std::string& saved_value() const { return saved_value_; } |
331 | | |
332 | 0 | std::vector<WideColumn>& lazy_entity_columns() { |
333 | 0 | return lazy_entity_columns_; |
334 | 0 | } |
335 | 0 | const std::vector<WideColumn>& lazy_entity_columns() const { |
336 | 0 | return lazy_entity_columns_; |
337 | 0 | } |
338 | | |
339 | 0 | std::vector<std::pair<size_t, BlobIndex>>& lazy_blob_columns() { |
340 | 0 | return lazy_blob_columns_; |
341 | 0 | } |
342 | 0 | const std::vector<std::pair<size_t, BlobIndex>>& lazy_blob_columns() const { |
343 | 0 | return lazy_blob_columns_; |
344 | 0 | } |
345 | | |
346 | 0 | ReadPathBlobResolver& entity_blob_resolver() { |
347 | 0 | return entity_blob_resolver_; |
348 | 0 | } |
349 | 0 | const ReadPathBlobResolver& entity_blob_resolver() const { |
350 | 0 | return entity_blob_resolver_; |
351 | 0 | } |
352 | | |
353 | 0 | std::mutex& lazy_entity_columns_mutex() const { |
354 | 0 | return lazy_entity_columns_mutex_; |
355 | 0 | } |
356 | | |
357 | | // DBIter calls this after ResetValueAndColumns() before repopulating the |
358 | | // current entry from a serialized wide-column entity. |
359 | 0 | void AssertReadyForEntity() const { |
360 | 0 | assert(value_.empty()); |
361 | 0 | assert(wide_columns_.empty()); |
362 | 0 | assert(lazy_entity_columns_.empty()); |
363 | 0 | assert(lazy_blob_columns_.empty()); |
364 | 0 | } |
365 | | |
366 | 43.9k | inline void ClearSavedValue() { |
367 | 43.9k | if (saved_value_.capacity() > 1048576) { |
368 | 0 | std::string empty; |
369 | 0 | swap(empty, saved_value_); |
370 | 43.9k | } else { |
371 | 43.9k | saved_value_.clear(); |
372 | 43.9k | } |
373 | 43.9k | } |
374 | | |
375 | | // Preserve the serialized entity bytes when DBIter needs stable backing |
376 | | // storage for lazy column slices across iterator movement. |
377 | 0 | void SaveEntitySliceIfNeeded(const Slice& slice) { |
378 | 0 | if (slice.data() != saved_value_.data() || |
379 | 0 | slice.size() != saved_value_.size()) { |
380 | 0 | saved_value_.assign(slice.data(), slice.size()); |
381 | 0 | } |
382 | 0 | } |
383 | | |
384 | | // Clears the previous lazy entity metadata and returns the saved entity |
385 | | // buffer as input for Deserialize(). |
386 | 0 | Slice PrepareForLazyEntityDeserialize() { |
387 | 0 | ClearLazyEntity(); |
388 | 0 | return Slice(saved_value_); |
389 | 0 | } |
390 | | |
391 | | // Publishes the deserialized lazy entity metadata to the blob resolver. |
392 | 0 | void BindLazyEntity(const Slice& user_key) { |
393 | 0 | entity_blob_resolver_.Reset(user_key, &lazy_entity_columns_, |
394 | 0 | &lazy_blob_columns_); |
395 | 0 | } |
396 | | |
397 | | // Drops the lazy entity metadata and any resolver aliases into it. |
398 | 37.5k | void ClearLazyEntity() { |
399 | 37.5k | lazy_entity_columns_.clear(); |
400 | 37.5k | lazy_blob_columns_.clear(); |
401 | 37.5k | entity_blob_resolver_.Reset(Slice(), nullptr, nullptr); |
402 | 37.5k | } |
403 | | |
404 | | // Clears materialized wide columns on error before DBIter invalidates |
405 | | // itself. |
406 | 0 | void ClearWideColumns() { wide_columns_.clear(); } |
407 | | |
408 | | // Fast path for inline entities whose default column is already |
409 | | // materialized in wide_columns_. |
410 | 0 | void MaybeSetValueFromMaterializedDefaultColumn() { |
411 | 0 | if (WideColumnsHelper::HasDefaultColumn(wide_columns_)) { |
412 | 0 | value_ = WideColumnsHelper::GetDefaultColumn(wide_columns_); |
413 | 0 | } |
414 | 0 | } |
415 | | |
416 | 39.4k | void SetFromPlain(const Slice& slice) { |
417 | 39.4k | assert(value_.empty()); |
418 | 39.4k | assert(wide_columns_.empty()); |
419 | 39.4k | assert(lazy_entity_columns_.empty()); |
420 | 39.4k | assert(lazy_blob_columns_.empty()); |
421 | | |
422 | 39.4k | value_ = slice; |
423 | 39.4k | wide_columns_.emplace_back(kDefaultWideColumnName, slice); |
424 | 39.4k | } |
425 | | |
426 | 37.5k | void Reset() { |
427 | 37.5k | value_.clear(); |
428 | 37.5k | wide_columns_.clear(); |
429 | 37.5k | ClearLazyEntity(); |
430 | 37.5k | } |
431 | | |
432 | | private: |
433 | | std::string saved_value_; |
434 | | // Value of the default column. |
435 | | Slice value_; |
436 | | // All columns (i.e. name-value pairs). |
437 | | WideColumns wide_columns_; |
438 | | // Lazy resolution state for V2 entities with blob columns. |
439 | | ReadPathBlobResolver entity_blob_resolver_; |
440 | | std::vector<WideColumn> lazy_entity_columns_; |
441 | | std::vector<std::pair<size_t, BlobIndex>> lazy_blob_columns_; |
442 | | mutable std::mutex lazy_entity_columns_mutex_; |
443 | | }; |
444 | | |
445 | | // For all methods in this block: |
446 | | // PRE: iter_->Valid() && status_.ok() |
447 | | // Return false if there was an error, and status() is non-ok, valid_ = false; |
448 | | // in this case callers would usually stop what they were doing and return. |
449 | | bool ReverseToForward(); |
450 | | bool ReverseToBackward(); |
451 | | // Set saved_key_ to the seek key to target, with proper sequence number set. |
452 | | // It might get adjusted if the seek key is smaller than iterator lower bound. |
453 | | // target does not have timestamp. |
454 | | void SetSavedKeyToSeekTarget(const Slice& target); |
455 | | // Set saved_key_ to the seek key to target, with proper sequence number set. |
456 | | // It might get adjusted if the seek key is larger than iterator upper bound. |
457 | | // target does not have timestamp. |
458 | | void SetSavedKeyToSeekForPrevTarget(const Slice& target); |
459 | | bool FindValueForCurrentKey(bool& found_visible); |
460 | | bool FindValueForCurrentKeyUsingSeek(); |
461 | | bool FindUserKeyBeforeSavedKey(); |
462 | | // If `skipping_saved_key` is true, the function will keep iterating until it |
463 | | // finds a user key that is larger than `saved_key_`. |
464 | | // When prefix_ is set, the iterator stops when all keys for the prefix are |
465 | | // exhausted and the iterator is set to invalid. |
466 | | bool FindNextUserEntry(bool skipping_saved_key); |
467 | | // Internal implementation of FindNextUserEntry(). |
468 | | bool FindNextUserEntryInternal(bool skipping_saved_key); |
469 | | bool ParseKey(ParsedInternalKey* key); |
470 | | bool MergeValuesNewToOld(); |
471 | | |
472 | | // When prefix_ is set, we need to set the iterator to invalid if no more |
473 | | // entry can be found within the prefix. |
474 | | void PrevInternal(); |
475 | | bool TooManyInternalKeysSkipped(bool increment = true); |
476 | | bool IsVisible(SequenceNumber sequence, const Slice& ts, |
477 | | bool* more_recent = nullptr); |
478 | | |
479 | | // Temporarily pin the blocks that we encounter until ReleaseTempPinnedData() |
480 | | // is called |
481 | 3.26k | void TempPinData() { |
482 | 3.26k | if (!pin_thru_lifetime_) { |
483 | 3.26k | pinned_iters_mgr_.StartPinning(); |
484 | 3.26k | } |
485 | 3.26k | } |
486 | | |
487 | | // Release blocks pinned by TempPinData() |
488 | 52.2k | void ReleaseTempPinnedData() { |
489 | 52.2k | if (!pin_thru_lifetime_ && pinned_iters_mgr_.PinningEnabled()) { |
490 | 114 | pinned_iters_mgr_.ReleasePinnedData(); |
491 | 114 | } |
492 | 52.2k | } |
493 | | |
494 | 43.9k | inline void ClearSavedValue() { |
495 | 43.9k | value_columns_state_.mut()->ClearSavedValue(); |
496 | 43.9k | } |
497 | | |
498 | 26.5k | inline void ResetInternalKeysSkippedCounter() { |
499 | 26.5k | local_stats_.skip_count_ += num_internal_keys_skipped_; |
500 | 26.5k | if (valid_) { |
501 | 1.85k | local_stats_.skip_count_--; |
502 | 1.85k | } |
503 | 26.5k | num_internal_keys_skipped_ = 0; |
504 | 26.5k | } |
505 | | |
506 | 7.47k | bool expect_total_order_inner_iter() { |
507 | 7.47k | assert(expect_total_order_inner_iter_ || prefix_extractor_ != nullptr); |
508 | 7.47k | return expect_total_order_inner_iter_; |
509 | 7.47k | } |
510 | | |
511 | | // If lower bound of timestamp is given by ReadOptions.iter_start_ts, we need |
512 | | // to return versions of the same key. We cannot just skip if the key value |
513 | | // is the same but timestamps are different but fall in timestamp range. |
514 | 46.5k | inline int CompareKeyForSkip(const Slice& a, const Slice& b) { |
515 | 46.5k | return timestamp_lb_ != nullptr |
516 | 46.5k | ? user_comparator_.Compare(a, b) |
517 | 46.5k | : user_comparator_.CompareWithoutTimestamp(a, b); |
518 | 46.5k | } |
519 | | |
520 | 39.4k | void SetValueAndColumnsFromPlain(const Slice& slice) { |
521 | 39.4k | value_columns_state_.mut()->SetFromPlain(slice); |
522 | 39.4k | } |
523 | | |
524 | | bool SetValueAndColumnsFromBlobImpl(const Slice& user_key, |
525 | | const Slice& blob_index); |
526 | | bool SetValueAndColumnsFromBlob(const Slice& user_key, |
527 | | const Slice& blob_index); |
528 | | |
529 | | bool SetValueAndColumnsFromEntity(Slice slice); |
530 | | bool MaterializeLazyEntityColumns() const; |
531 | | |
532 | | bool SetValueAndColumnsFromMergeResult(const Status& merge_status, |
533 | | ValueType result_type); |
534 | | |
535 | 48.9k | void ResetValueAndColumns() { value_columns_state_.Reset(); } |
536 | | |
537 | 48.9k | void ResetBlobData() { blob_state_.Reset(); } |
538 | | |
539 | | // The following methods perform the actual merge operation for the |
540 | | // no/plain/blob/wide-column base value cases. |
541 | | // If user-defined timestamp is enabled, `user_key` includes timestamp. |
542 | | bool MergeWithNoBaseValue(const Slice& user_key); |
543 | | bool MergeWithPlainBaseValue(const Slice& value, const Slice& user_key); |
544 | | bool MergeWithBlobBaseValue(const Slice& blob_index, const Slice& user_key); |
545 | | bool MergeWithWideColumnBaseValue(const Slice& entity, const Slice& user_key); |
546 | | |
547 | 49.9k | bool PrepareValueInternal() { |
548 | | // Capture this before PrepareValue(): PrepareValue() updates the wrapper |
549 | | // state to "prepared" on success. We still call PrepareValue() |
550 | | // unconditionally to preserve its contract/error handling, but only need |
551 | | // to re-parse ikey_ when this call may have actually materialized the |
552 | | // underlying iterator value/key. |
553 | 49.9k | const bool value_was_prepared = iter_.IsValuePrepared(); |
554 | 49.9k | if (!iter_.PrepareValue()) { |
555 | 0 | assert(!iter_.status().ok()); |
556 | 0 | valid_ = false; |
557 | 0 | return false; |
558 | 0 | } |
559 | | // ikey_ could change as BlockBasedTableIterator does Block cache |
560 | | // lookup and index_iter_ could point to different block resulting |
561 | | // in ikey_ pointing to wrong key. So ikey_ needs to be updated in |
562 | | // case of Seek/Next calls to point to right key again. |
563 | 49.9k | if (!value_was_prepared && !ParseKey(&ikey_)) { |
564 | 0 | return false; |
565 | 0 | } |
566 | 49.9k | return true; |
567 | 49.9k | } |
568 | | |
569 | | // Record a deletion into the current contiguous tombstone run. |
570 | | // In forward iteration, first_key is set only for the first tombstone |
571 | | // (always_update_first_key=false). In reverse, keys arrive in decreasing |
572 | | // order so first_key is updated every time (always_update_first_key=true). |
573 | | void TrackContiguousTombstone(const Slice& user_key, |
574 | | bool always_update_first_key); |
575 | | |
576 | | // If a contiguous tombstone run is pending, insert a range tombstone |
577 | | // (if threshold is met) and reset tracking state. When |
578 | | // check_prefix_match is true, the insertion is skipped (but tracking is |
579 | | // still reset) if end_key is outside the seek prefix. |
580 | | void FlushPendingTombstoneRun(const Slice& end_key, |
581 | | bool check_prefix_match = false); |
582 | | |
583 | | // If enough contiguous tombstones have been tracked, insert a range |
584 | | // tombstone [first_key, end_key) into the mutable memtable. |
585 | | // end_key is the exclusive upper bound -- typically the next live key. |
586 | | void MaybeInsertRangeTombstone(const Slice& end_key); |
587 | 13.4k | void ResetContiguousTombstoneTracking() { |
588 | 13.4k | contiguous_tombstone_count_ = 0; |
589 | 13.4k | range_tomb_first_key_.Clear(); |
590 | 13.4k | range_tomb_end_key_.Clear(); |
591 | 13.4k | } |
592 | | |
593 | | // Returns true if there is no prefix constraint (prefix_ not set) or |
594 | | // if `key` is in the prefix extractor's domain and its prefix matches. |
595 | | // Out-of-domain keys return false when a prefix is set. |
596 | 55.5k | bool PrefixCheck(const Slice& key) const { |
597 | 55.5k | return !prefix_.has_value() || (prefix_extractor_->InDomain(key) && |
598 | 0 | prefix_extractor_->Transform(key).compare( |
599 | 0 | prefix_->GetUserKey()) == 0); |
600 | 55.5k | } |
601 | | |
602 | | // Returns true if a prefix should be extracted from the seek target and |
603 | | // used for prefix boundary tracking. True when prefix_same_as_start is |
604 | | // set, or when range tombstone conversion is enabled during a legacy |
605 | | // prefix seek. If target is out of domain then false is returned. |
606 | 8.38k | bool ShouldSetPrefix(const Slice& target) const { |
607 | 8.38k | return (prefix_same_as_start_ || |
608 | 8.38k | (min_tombstones_for_range_conversion_ > 0 && |
609 | 0 | !expect_total_order_inner_iter_)) && |
610 | 0 | prefix_extractor_->InDomain(target); |
611 | 8.38k | } |
612 | | |
613 | 11.3k | void ResetSeekState() { |
614 | 11.3k | ReleaseTempPinnedData(); |
615 | 11.3k | ResetBlobData(); |
616 | 11.3k | ResetValueAndColumns(); |
617 | 11.3k | ResetInternalKeysSkippedCounter(); |
618 | 11.3k | ResetContiguousTombstoneTracking(); |
619 | 11.3k | prefix_.reset(); |
620 | 11.3k | } |
621 | | |
622 | 26.5k | void MarkMemtableForFlushForAvgTrigger() { |
623 | 26.5k | if (avg_op_scan_flush_trigger_ && |
624 | 0 | mem_hidden_op_scanned_since_seek_ >= memtable_op_scan_flush_trigger_ && |
625 | 0 | mem_hidden_op_scanned_since_seek_ >= |
626 | 0 | static_cast<uint64_t>(iter_step_since_seek_) * |
627 | 0 | avg_op_scan_flush_trigger_) { |
628 | 0 | assert(memtable_op_scan_flush_trigger_ > 0); |
629 | 0 | active_mem_->MarkForFlush(); |
630 | 0 | avg_op_scan_flush_trigger_ = 0; |
631 | 0 | memtable_op_scan_flush_trigger_ = 0; |
632 | 0 | } |
633 | 26.5k | iter_step_since_seek_ = 1; |
634 | 26.5k | mem_hidden_op_scanned_since_seek_ = 0; |
635 | 26.5k | } |
636 | | |
637 | 14.7k | void MarkMemtableForFlushForPerOpTrigger(uint64_t& mem_hidden_op_scanned) { |
638 | 14.7k | if (memtable_op_scan_flush_trigger_ && |
639 | 0 | ikey_.sequence >= memtable_seqno_lb_) { |
640 | 0 | if (++mem_hidden_op_scanned >= memtable_op_scan_flush_trigger_) { |
641 | 0 | active_mem_->MarkForFlush(); |
642 | | // Turn off the flush trigger checks. |
643 | 0 | memtable_op_scan_flush_trigger_ = 0; |
644 | 0 | avg_op_scan_flush_trigger_ = 0; |
645 | 0 | } |
646 | 0 | if (avg_op_scan_flush_trigger_) { |
647 | 0 | ++mem_hidden_op_scanned_since_seek_; |
648 | 0 | } |
649 | 0 | } |
650 | 14.7k | } |
651 | | |
652 | | const SliceTransform* prefix_extractor_; |
653 | | Env* const env_; |
654 | | SystemClock* clock_; |
655 | | Logger* logger_; |
656 | | UserComparatorWrapper user_comparator_; |
657 | | const MergeOperator* const merge_operator_; |
658 | | IteratorWrapper iter_; |
659 | | // TODO: blob_state_'s BlobReader (whole-value blob reads + the merge entity |
660 | | // path) and value_columns_state_'s ReadPathBlobResolver (lazy per-column |
661 | | // entity resolution) each own a Version-backed blob fetcher with the same |
662 | | // {version, read_options, blob_file_cache, allow_write_path_fallback}, so an |
663 | | // iterator holds two ReadOptions copies. Collapse to a single shared fetcher. |
664 | | // Deferred because the resolver must keep owning its fetcher for the planned |
665 | | // lazy entity read path, where the result outlives the read call and owns the |
666 | | // SuperVersion pin; unifying is best done together with that work. |
667 | | DirtyTracked<BlobState> blob_state_; |
668 | | ReadCallback* read_callback_; |
669 | | // Max visible sequence number. It is normally the snapshot seq unless we have |
670 | | // uncommitted data in db as in WriteUnCommitted. |
671 | | SequenceNumber sequence_; |
672 | | |
673 | | IterKey saved_key_; |
674 | | // Reusable internal key data structure. This is only used inside one function |
675 | | // and should not be used across functions. Reusing this object can reduce |
676 | | // overhead of calling construction of the function if creating it each time. |
677 | | ParsedInternalKey ikey_; |
678 | | |
679 | | // The approximate write time for the entry. It is deduced from the entry's |
680 | | // sequence number if the seqno to time mapping is available. For a |
681 | | // kTypeValuePreferredSeqno entry, this is the write time specified by the |
682 | | // user. |
683 | | uint64_t saved_write_unix_time_; |
684 | | DirtyTracked<ValueColumnsState> value_columns_state_; |
685 | | Slice pinned_value_; |
686 | | Statistics* statistics_; |
687 | | uint64_t max_skip_; |
688 | | uint64_t max_skippable_internal_keys_; |
689 | | uint64_t num_internal_keys_skipped_; |
690 | | const Slice* iterate_lower_bound_; |
691 | | const Slice* iterate_upper_bound_; |
692 | | |
693 | | // The prefix of the seek key. Set during Seek/SeekForPrev when either |
694 | | // prefix_same_as_start_ is true or the iterator uses prefix filtering |
695 | | // (!expect_total_order_inner_iter_ && InDomain(target)). Used in Next() |
696 | | // and Prev() to: |
697 | | // - invalidate the iterator when prefix_same_as_start_ is true and keys |
698 | | // in this prefix have been exhausted. |
699 | | // - bound range tombstone tracking to the seek prefix when |
700 | | // min_tombstones_for_range_conversion_ > 0. |
701 | | // Set via SetUserKey(), read via GetUserKey(). |
702 | | std::optional<IterKey> prefix_; |
703 | | |
704 | | Status status_; |
705 | | |
706 | | // List of operands for merge operator. |
707 | | MergeContext merge_context_; |
708 | | LocalStatistics local_stats_; |
709 | | PinnedIteratorsManager pinned_iters_mgr_; |
710 | | DBImpl* trace_db_; |
711 | | uint32_t trace_cf_id_; |
712 | | bool has_trace_state_; |
713 | | port::RWMutex* ingest_sst_lock_; |
714 | | const Slice* const timestamp_ub_; |
715 | | const Slice* const timestamp_lb_; |
716 | | const size_t timestamp_size_; |
717 | | std::string saved_timestamp_; |
718 | | std::optional<MultiScanArgs> scan_opts_; |
719 | | size_t scan_index_{0}; |
720 | | ReadOnlyMemTable* const active_mem_; |
721 | | SequenceNumber memtable_seqno_lb_; |
722 | | uint32_t memtable_op_scan_flush_trigger_; |
723 | | uint32_t avg_op_scan_flush_trigger_; |
724 | | uint32_t iter_step_since_seek_; |
725 | | uint32_t mem_hidden_op_scanned_since_seek_; |
726 | | uint32_t contiguous_tombstone_count_; |
727 | | Direction direction_; |
728 | | bool valid_; |
729 | | bool current_entry_is_merged_; |
730 | | // True if we know that the current entry's seqnum is 0. |
731 | | // This information is used as that the next entry will be for another |
732 | | // user key. |
733 | | bool is_key_seqnum_zero_; |
734 | | const bool prefix_same_as_start_; |
735 | | // Means that we will pin all data blocks we read as long the Iterator |
736 | | // is not deleted, will be true if ReadOptions::pin_data is true |
737 | | const bool pin_thru_lifetime_; |
738 | | // Expect the inner iterator to maintain a total order. |
739 | | // prefix_extractor_ must be non-NULL if the value is false. |
740 | | const bool expect_total_order_inner_iter_; |
741 | | const uint32_t min_tombstones_for_range_conversion_; |
742 | | // Whether the iterator is allowed to expose blob references. Set to true when |
743 | | // the stacked BlobDB implementation is used, false otherwise. |
744 | | bool expose_blob_index_; |
745 | | bool allow_unprepared_value_; |
746 | | bool arena_mode_; |
747 | | |
748 | | IterKey range_tomb_first_key_; |
749 | | IterKey range_tomb_end_key_; |
750 | | }; |
751 | | } // namespace ROCKSDB_NAMESPACE |