/src/rocksdb/db/table_cache.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 | | // 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 | | #include "db/table_cache.h" |
11 | | |
12 | | #include <algorithm> |
13 | | |
14 | | #include "db/dbformat.h" |
15 | | #include "db/range_tombstone_fragmenter.h" |
16 | | #include "db/snapshot_impl.h" |
17 | | #include "db/version_edit.h" |
18 | | #include "file/file_util.h" |
19 | | #include "file/filename.h" |
20 | | #include "file/random_access_file_reader.h" |
21 | | #include "logging/logging.h" |
22 | | #include "monitoring/file_read_sample.h" |
23 | | #include "monitoring/perf_context_imp.h" |
24 | | #include "options/options_helper.h" |
25 | | #include "rocksdb/advanced_options.h" |
26 | | #include "rocksdb/statistics.h" |
27 | | #include "table/block_based/block_based_table_reader.h" |
28 | | #include "table/get_context.h" |
29 | | #include "table/internal_iterator.h" |
30 | | #include "table/iterator_wrapper.h" |
31 | | #include "table/multiget_context.h" |
32 | | #include "table/table_builder.h" |
33 | | #include "table/table_reader.h" |
34 | | #include "test_util/sync_point.h" |
35 | | #include "util/cast_util.h" |
36 | | #include "util/coding.h" |
37 | | #include "util/stop_watch.h" |
38 | | |
39 | | // Generate the regular and coroutine versions of some methods by |
40 | | // including table_cache_sync_and_async.h twice |
41 | | // Macros in the header will expand differently based on whether |
42 | | // WITH_COROUTINES or WITHOUT_COROUTINES is defined |
43 | | // clang-format off |
44 | | #define WITHOUT_COROUTINES |
45 | | #include "db/table_cache_sync_and_async.h" |
46 | | #undef WITHOUT_COROUTINES |
47 | | #define WITH_COROUTINES |
48 | | #include "db/table_cache_sync_and_async.h" |
49 | | #undef WITH_COROUTINES |
50 | | // clang-format on |
51 | | |
52 | | namespace ROCKSDB_NAMESPACE { |
53 | | |
54 | | namespace { |
55 | | |
56 | 299k | static Slice GetSliceForFileNumber(const uint64_t* file_number) { |
57 | 299k | return Slice(reinterpret_cast<const char*>(file_number), |
58 | 299k | sizeof(*file_number)); |
59 | 299k | } |
60 | | |
61 | 0 | void AppendVarint64(IterKey* key, uint64_t v) { |
62 | 0 | char buf[10]; |
63 | 0 | auto ptr = EncodeVarint64(buf, v); |
64 | 0 | key->TrimAppend(key->Size(), buf, ptr - buf); |
65 | 0 | } |
66 | | |
67 | | } // anonymous namespace |
68 | | |
69 | | const int kLoadConcurency = 128; |
70 | | |
71 | | TableCache::TableCache(const ImmutableOptions& ioptions, |
72 | | const FileOptions* file_options, Cache* const cache, |
73 | | BlockCacheTracer* const block_cache_tracer, |
74 | | const std::shared_ptr<IOTracer>& io_tracer, |
75 | | const std::string& db_session_id, bool fast_sst_open) |
76 | 106k | : ioptions_(ioptions), |
77 | 106k | file_options_(*file_options), |
78 | 106k | cache_(cache), |
79 | 106k | immortal_tables_(false), |
80 | 106k | should_pin_table_handles_(cache_.get()->GetCapacity() >= |
81 | 106k | kInfiniteCapacity), |
82 | 106k | fast_sst_open_(fast_sst_open), |
83 | 106k | block_cache_tracer_(block_cache_tracer), |
84 | 106k | loader_mutex_(kLoadConcurency), |
85 | 106k | io_tracer_(io_tracer), |
86 | 106k | db_session_id_(db_session_id) { |
87 | 106k | if (ioptions_.row_cache) { |
88 | | // If the same cache is shared by multiple instances, we need to |
89 | | // disambiguate its entries. |
90 | 0 | PutVarint64(&row_cache_id_, ioptions_.row_cache->NewId()); |
91 | 0 | } |
92 | 106k | } |
93 | | |
94 | 106k | TableCache::~TableCache() = default; |
95 | | |
96 | | Status TableCache::GetTableReader( |
97 | | const ReadOptions& ro, const FileOptions& file_options, |
98 | | const InternalKeyComparator& internal_comparator, |
99 | | const FileMetaData& file_meta, bool sequential_mode, |
100 | | HistogramImpl* file_read_hist, std::unique_ptr<TableReader>* table_reader, |
101 | | const MutableCFOptions& mutable_cf_options, bool skip_filters, int level, |
102 | | bool prefetch_index_and_filter_in_cache, |
103 | | size_t max_file_size_for_l0_meta_pin, Temperature file_temperature, |
104 | 123k | std::string* file_open_metadata, bool avoid_shared_metadata_cache) { |
105 | 123k | std::string fname = TableFileName( |
106 | 123k | ioptions_.cf_paths, file_meta.fd.GetNumber(), file_meta.fd.GetPathId()); |
107 | 123k | std::unique_ptr<FSRandomAccessFile> file; |
108 | 123k | FileOptions fopts = file_options; |
109 | 123k | fopts.temperature = file_temperature; |
110 | 123k | fopts.file_checksum = file_meta.file_checksum; |
111 | 123k | fopts.file_checksum_func_name = file_meta.file_checksum_func_name; |
112 | | // Pass file open metadata for fast SST open. Use a local copy since |
113 | | // fopts.file_metadata is a non-owning pointer and file_meta is const. |
114 | | // Only pass metadata when fast_sst_open is enabled; otherwise ignore |
115 | | // previously persisted metadata (e.g. stale filesystem credentials). |
116 | 123k | std::string file_open_metadata_copy; |
117 | 123k | if (fast_sst_open_.load(std::memory_order_relaxed) && |
118 | 0 | !file_meta.file_open_metadata.empty()) { |
119 | 0 | file_open_metadata_copy = file_meta.file_open_metadata; |
120 | 0 | fopts.file_metadata = &file_open_metadata_copy; |
121 | 0 | RecordTick(ioptions_.stats, FILE_OPEN_METADATA_PASSED); |
122 | 0 | } |
123 | 123k | Status s = PrepareIOFromReadOptions(ro, ioptions_.clock, fopts.io_options); |
124 | 123k | TEST_SYNC_POINT_CALLBACK("TableCache::GetTableReader:BeforeOpenFile", |
125 | 123k | const_cast<Status*>(&s)); |
126 | 123k | if (s.ok()) { |
127 | 123k | s = ioptions_.fs->NewRandomAccessFile(fname, fopts, &file, nullptr); |
128 | 123k | } |
129 | 123k | if (s.ok()) { |
130 | 123k | RecordTick(ioptions_.stats, NO_FILE_OPENS); |
131 | 123k | } else if (s.IsPathNotFound()) { |
132 | 0 | fname = Rocks2LevelTableFileName(fname); |
133 | | // If this file is also not found, we want to use the error message |
134 | | // that contains the table file name which is less confusing. |
135 | 0 | Status temp_s = |
136 | 0 | PrepareIOFromReadOptions(ro, ioptions_.clock, fopts.io_options); |
137 | 0 | if (temp_s.ok()) { |
138 | 0 | temp_s = ioptions_.fs->NewRandomAccessFile(fname, fopts, &file, nullptr); |
139 | 0 | } |
140 | 0 | if (temp_s.ok()) { |
141 | 0 | RecordTick(ioptions_.stats, NO_FILE_OPENS); |
142 | 0 | s = temp_s; |
143 | 0 | } |
144 | 0 | } |
145 | | |
146 | 123k | if (s.ok()) { |
147 | | // Retrieve file open metadata before wrapping the file |
148 | 123k | if (file_open_metadata != nullptr) { |
149 | 0 | IOStatus io_s = file->GetFileOpenMetadata(file_open_metadata); |
150 | 0 | if (io_s.ok() && !file_open_metadata->empty() && |
151 | 0 | file_open_metadata->size() <= |
152 | 0 | FSRandomAccessFile::kMaxFileOpenMetadataSize) { |
153 | 0 | RecordTick(ioptions_.stats, FILE_OPEN_METADATA_RETRIEVED); |
154 | 0 | } else { |
155 | 0 | if (io_s.ok() && file_open_metadata->size() > |
156 | 0 | FSRandomAccessFile::kMaxFileOpenMetadataSize) { |
157 | 0 | ROCKS_LOG_WARN(ioptions_.logger, |
158 | 0 | "File open metadata for %s too large (%zu bytes), " |
159 | 0 | "ignoring", |
160 | 0 | fname.c_str(), file_open_metadata->size()); |
161 | 0 | } |
162 | 0 | file_open_metadata->clear(); |
163 | 0 | } |
164 | 0 | } |
165 | 123k | if (!sequential_mode && ioptions_.advise_random_on_open) { |
166 | 123k | file->Hint(FSRandomAccessFile::kRandom); |
167 | 123k | } |
168 | 123k | if (ioptions_.default_temperature != Temperature::kUnknown && |
169 | 0 | file_temperature == Temperature::kUnknown) { |
170 | 0 | file_temperature = ioptions_.default_temperature; |
171 | 0 | } |
172 | 123k | StopWatch sw(ioptions_.clock, ioptions_.stats, TABLE_OPEN_IO_MICROS); |
173 | 123k | std::unique_ptr<RandomAccessFileReader> file_reader( |
174 | 123k | new RandomAccessFileReader(std::move(file), fname, ioptions_.clock, |
175 | 123k | io_tracer_, ioptions_.stats, SST_READ_MICROS, |
176 | 123k | file_read_hist, ioptions_.rate_limiter.get(), |
177 | 123k | ioptions_.listeners, file_temperature, |
178 | 123k | level == ioptions_.num_levels - 1)); |
179 | 123k | UniqueId64x2 expected_unique_id; |
180 | 123k | if (ioptions_.verify_sst_unique_id_in_manifest) { |
181 | 123k | expected_unique_id = file_meta.unique_id; |
182 | 18.4E | } else { |
183 | 18.4E | expected_unique_id = kNullUniqueId64x2; // null ID == no verification |
184 | 18.4E | } |
185 | 123k | TableReaderOptions table_reader_options( |
186 | 123k | ioptions_, mutable_cf_options.prefix_extractor, |
187 | 123k | mutable_cf_options.compression_manager.get(), file_options, |
188 | 123k | internal_comparator, mutable_cf_options.block_protection_bytes_per_key, |
189 | 123k | skip_filters, immortal_tables_, false /* force_direct_prefetch */, |
190 | 123k | level, block_cache_tracer_, max_file_size_for_l0_meta_pin, |
191 | 123k | db_session_id_, file_meta.fd.GetNumber(), expected_unique_id, |
192 | 123k | file_meta.fd.largest_seqno, file_meta.tail_size, |
193 | 123k | file_meta.user_defined_timestamps_persisted, |
194 | 123k | avoid_shared_metadata_cache); |
195 | | // Route same-file ("embedded") blob reads through the CFD's BlobSource for |
196 | | // caching + stats. nullptr in non-DB contexts (e.g. repair). |
197 | 123k | table_reader_options.blob_source = blob_source_; |
198 | 123k | s = mutable_cf_options.table_factory->NewTableReader( |
199 | 123k | ro, table_reader_options, std::move(file_reader), |
200 | 123k | file_meta.fd.GetFileSize(), table_reader, |
201 | 123k | prefetch_index_and_filter_in_cache); |
202 | 123k | TEST_SYNC_POINT("TableCache::GetTableReader:0"); |
203 | 123k | } |
204 | 123k | return s; |
205 | 123k | } |
206 | | |
207 | 0 | Cache::Handle* TableCache::Lookup(Cache* cache, uint64_t file_number) { |
208 | | // NOTE: sharing same Cache with BlobFileCache |
209 | 0 | Slice key = GetSliceForFileNumber(&file_number); |
210 | 0 | return cache->Lookup(key); |
211 | 0 | } |
212 | | |
213 | | // TODO: consider making handle RAII. |
214 | | Status TableCache::FindTable( |
215 | | const ReadOptions& ro, const FileOptions& file_options, |
216 | | const InternalKeyComparator& internal_comparator, |
217 | | const FileMetaData& file_meta, TypedHandle** handle, |
218 | | const MutableCFOptions& mutable_cf_options, TableReader** out_table_reader, |
219 | | const bool no_io, HistogramImpl* file_read_hist, bool skip_filters, |
220 | | int level, bool prefetch_index_and_filter_in_cache, |
221 | | size_t max_file_size_for_l0_meta_pin, Temperature file_temperature, |
222 | | bool pin_table_handle, std::string* file_open_metadata, |
223 | | std::unique_ptr<TableReader>* fresh_table_reader_owner, |
224 | 353k | const TableCacheOpenOptions& open_options) { |
225 | 353k | assert(out_table_reader != nullptr && *out_table_reader == nullptr); |
226 | 353k | assert(handle != nullptr && *handle == nullptr); |
227 | | // open_ephemeral_table_reader requests a fresh reader; |
228 | | // fresh_table_reader_owner is where we return it. The two must agree. |
229 | 353k | assert(open_options.open_ephemeral_table_reader == |
230 | 353k | (fresh_table_reader_owner != nullptr)); |
231 | 353k | PERF_TIMER_GUARD_WITH_CLOCK(find_table_nanos, ioptions_.clock); |
232 | | |
233 | | // Bypass path: open a fresh TableReader, skipping the pinned-reader fast path |
234 | | // and the shared cache. The caller takes ownership via |
235 | | // fresh_table_reader_owner. no_io is not allowed here since opening a new |
236 | | // reader always needs I/O. |
237 | 353k | if (fresh_table_reader_owner != nullptr) { |
238 | 0 | assert(!no_io); |
239 | 0 | if (no_io) { |
240 | | // Defensive in release builds: the caller violated the contract. |
241 | 0 | return Status::Incomplete( |
242 | 0 | "fresh TableReader requested but no_io is set; cannot open file"); |
243 | 0 | } |
244 | 0 | std::unique_ptr<TableReader> table_reader; |
245 | 0 | const bool effective_skip_filters = |
246 | 0 | skip_filters || open_options.skip_filters; |
247 | 0 | TEST_SYNC_POINT_CALLBACK("TableCache::FindTable:FreshTableReader", |
248 | 0 | const_cast<TableCacheOpenOptions*>(&open_options)); |
249 | 0 | Status s = GetTableReader( |
250 | 0 | ro, file_options, internal_comparator, file_meta, |
251 | 0 | false /* sequential mode */, file_read_hist, &table_reader, |
252 | 0 | mutable_cf_options, effective_skip_filters, level, |
253 | 0 | prefetch_index_and_filter_in_cache && |
254 | 0 | !open_options.avoid_shared_metadata_cache, |
255 | 0 | max_file_size_for_l0_meta_pin, file_temperature, file_open_metadata, |
256 | 0 | open_options.avoid_shared_metadata_cache); |
257 | 0 | if (!s.ok()) { |
258 | 0 | assert(table_reader == nullptr); |
259 | 0 | RecordTick(ioptions_.stats, NO_FILE_ERRORS); |
260 | 0 | IGNORE_STATUS_IF_ERROR(s); |
261 | 0 | return s; |
262 | 0 | } |
263 | 0 | *out_table_reader = table_reader.get(); |
264 | 0 | *fresh_table_reader_owner = std::move(table_reader); |
265 | 0 | *handle = nullptr; |
266 | 0 | return s; |
267 | 0 | } |
268 | | |
269 | | // Fast path: if table reader is already pinned, return it directly without a |
270 | | // cache lookup. |
271 | 353k | auto pinned_reader = file_meta.fd.pinned_reader.Get(); |
272 | 353k | if (pinned_reader != nullptr) { |
273 | 199k | *handle = nullptr; |
274 | 199k | *out_table_reader = pinned_reader; |
275 | 199k | return Status::OK(); |
276 | 199k | } |
277 | | |
278 | 154k | uint64_t number = file_meta.fd.GetNumber(); |
279 | | // NOTE: sharing same Cache with BlobFileCache |
280 | 154k | Slice key = GetSliceForFileNumber(&number); |
281 | 154k | *handle = cache_.Lookup(key); |
282 | 154k | TEST_SYNC_POINT_CALLBACK("TableCache::FindTable:0", |
283 | 154k | const_cast<bool*>(&no_io)); |
284 | | |
285 | 154k | Status s = Status::OK(); |
286 | 154k | if (*handle == nullptr) { |
287 | 123k | if (no_io) { |
288 | 0 | s = Status::Incomplete("Table not found in table_cache, no_io is set"); |
289 | 0 | return s; |
290 | 0 | } |
291 | 123k | MutexLock load_lock(&loader_mutex_.Get(key)); |
292 | | |
293 | | // Check if another thread has already pinned the table reader |
294 | 123k | pinned_reader = file_meta.fd.pinned_reader.Get(); |
295 | 123k | if (pinned_reader != nullptr) { |
296 | 0 | *handle = nullptr; |
297 | 0 | *out_table_reader = pinned_reader; |
298 | 0 | return s; |
299 | 0 | } |
300 | | |
301 | | // We check the cache again under loading mutex |
302 | 123k | *handle = cache_.Lookup(key); |
303 | 123k | if (*handle == nullptr) { |
304 | 123k | std::unique_ptr<TableReader> table_reader; |
305 | 123k | s = GetTableReader(ro, file_options, internal_comparator, file_meta, |
306 | 123k | false /* sequential mode */, file_read_hist, |
307 | 123k | &table_reader, mutable_cf_options, skip_filters, level, |
308 | 123k | prefetch_index_and_filter_in_cache, |
309 | 123k | max_file_size_for_l0_meta_pin, file_temperature, |
310 | 123k | file_open_metadata); |
311 | 123k | if (!s.ok()) { |
312 | 0 | assert(table_reader == nullptr); |
313 | 0 | RecordTick(ioptions_.stats, NO_FILE_ERRORS); |
314 | | // We do not cache error results so that if the error is transient, |
315 | | // or somebody repairs the file, we recover automatically. |
316 | 0 | IGNORE_STATUS_IF_ERROR(s); |
317 | 123k | } else { |
318 | 123k | s = cache_.Insert(key, table_reader.get(), 1, handle); |
319 | 123k | if (s.ok()) { |
320 | | // Release ownership of table reader. |
321 | 122k | (void)table_reader.release(); |
322 | 122k | } |
323 | 123k | } |
324 | 123k | } |
325 | | |
326 | 123k | if (s.ok()) { |
327 | 122k | *out_table_reader = cache_.Value(*handle); |
328 | 122k | if (pin_table_handle) { |
329 | 99.0k | file_meta.fd.pinned_reader.Pin(*handle, *out_table_reader); |
330 | 99.0k | *handle = nullptr; |
331 | 99.0k | } |
332 | 122k | } |
333 | 123k | } else { |
334 | 30.7k | *out_table_reader = cache_.Value(*handle); |
335 | 30.7k | if (pin_table_handle) { |
336 | | // handle is in cache but not pinned. This should happen fairly rarely, |
337 | | // and once the reader is pinned, we will no longer need to go through |
338 | | // these mutexes again. |
339 | 30.7k | MutexLock load_lock(&loader_mutex_.Get(key)); |
340 | 30.7k | if (file_meta.fd.pinned_reader.Get() != nullptr) { |
341 | | // Another thread has pinned the handle; release our lookup ref. |
342 | 0 | cache_.Release(*handle); |
343 | 30.7k | } else { |
344 | 30.7k | file_meta.fd.pinned_reader.Pin(*handle, *out_table_reader); |
345 | 30.7k | } |
346 | 30.7k | *handle = nullptr; |
347 | 30.7k | } |
348 | 30.7k | } |
349 | | |
350 | 154k | return s; |
351 | 154k | } |
352 | | |
353 | | InternalIterator* TableCache::NewIterator( |
354 | | const ReadOptions& options, const FileOptions& file_options, |
355 | | const InternalKeyComparator& icomparator, const FileMetaData& file_meta, |
356 | | RangeDelAggregator* range_del_agg, |
357 | | const MutableCFOptions& mutable_cf_options, TableReader** table_reader_ptr, |
358 | | HistogramImpl* file_read_hist, TableReaderCaller caller, Arena* arena, |
359 | | bool skip_filters, int level, size_t max_file_size_for_l0_meta_pin, |
360 | | const InternalKey* smallest_compaction_key, |
361 | | const InternalKey* largest_compaction_key, bool allow_unprepared_value, |
362 | | const SequenceNumber* read_seqno, |
363 | | std::unique_ptr<TruncatedRangeDelIterator>* range_del_iter, |
364 | | bool maybe_pin_table_handle, std::string* file_open_metadata, |
365 | 68.5k | const TableCacheOpenOptions& open_options) { |
366 | 68.5k | PERF_TIMER_GUARD(new_table_iterator_nanos); |
367 | | |
368 | 68.5k | Status s; |
369 | 68.5k | TableReader* table_reader = nullptr; |
370 | 68.5k | TypedHandle* handle = nullptr; |
371 | 68.5k | assert(!open_options.open_ephemeral_table_reader || |
372 | 68.5k | table_reader_ptr == nullptr); |
373 | | // Holds ownership of a freshly-opened TableReader when the caller asked us |
374 | | // to bypass the shared cache. When non-empty, the iterator we hand back must |
375 | | // arrange to free it on destruction. |
376 | 68.5k | std::unique_ptr<TableReader> ephemeral_reader; |
377 | 68.5k | if (table_reader_ptr != nullptr) { |
378 | 1.71k | *table_reader_ptr = nullptr; |
379 | 1.71k | } |
380 | 68.5k | const bool effective_skip_filters = skip_filters || open_options.skip_filters; |
381 | 68.5k | bool for_compaction = caller == TableReaderCaller::kCompaction; |
382 | 68.5k | TEST_SYNC_POINT_CALLBACK("TableCache::NewIterator::BeforeFindTable", |
383 | 68.5k | const_cast<FileDescriptor*>(&file_meta.fd)); |
384 | 68.5k | s = FindTable( |
385 | 68.5k | options, file_options, icomparator, file_meta, &handle, |
386 | 68.5k | mutable_cf_options, &table_reader, |
387 | 68.5k | options.read_tier == kBlockCacheTier /* no_io */, file_read_hist, |
388 | 68.5k | effective_skip_filters, level, |
389 | | /*prefetch_index_and_filter_in_cache=*/ |
390 | 68.5k | !open_options.avoid_shared_metadata_cache, max_file_size_for_l0_meta_pin, |
391 | 68.5k | file_meta.temperature, |
392 | 68.5k | maybe_pin_table_handle && should_pin_table_handles_, file_open_metadata, |
393 | 68.5k | open_options.open_ephemeral_table_reader ? &ephemeral_reader : nullptr, |
394 | 68.5k | open_options); |
395 | 68.5k | InternalIterator* result = nullptr; |
396 | 68.5k | if (s.ok()) { |
397 | 68.5k | if (HasTableFilter(options) && |
398 | 0 | !(*options.table_filter)(*table_reader->GetTableProperties())) { |
399 | 0 | result = NewEmptyInternalIterator<Slice>(arena); |
400 | 68.5k | } else { |
401 | 68.5k | result = table_reader->NewIterator( |
402 | 68.5k | options, mutable_cf_options.prefix_extractor.get(), arena, |
403 | 68.5k | effective_skip_filters, caller, |
404 | 68.5k | file_options.compaction_readahead_size, allow_unprepared_value); |
405 | 68.5k | } |
406 | 68.5k | if (handle != nullptr) { |
407 | 23.2k | cache_.RegisterReleaseAsCleanup(handle, *result); |
408 | 23.2k | handle = nullptr; // prevent from releasing below |
409 | 23.2k | } |
410 | | // Don't hand ephemeral_reader to result's cleanup yet: range-del |
411 | | // processing below can set s to non-OK, in which case result is replaced |
412 | | // by an error iterator at function end. Transfer ownership only once we |
413 | | // know s stays OK; otherwise ephemeral_reader's destructor frees it. |
414 | | |
415 | 68.5k | if (for_compaction) { |
416 | 20.6k | table_reader->SetupForCompaction(); |
417 | 20.6k | } |
418 | 68.5k | if (table_reader_ptr != nullptr) { |
419 | 1.71k | *table_reader_ptr = table_reader; |
420 | 1.71k | } |
421 | 68.5k | } |
422 | 68.5k | if (s.ok() && !options.ignore_range_deletions) { |
423 | 68.5k | if (range_del_iter != nullptr) { |
424 | 41.1k | auto new_range_del_iter = |
425 | 41.1k | read_seqno ? table_reader->NewRangeTombstoneIterator( |
426 | 25.9k | *read_seqno, options.timestamp) |
427 | 41.1k | : table_reader->NewRangeTombstoneIterator(options); |
428 | 41.1k | if (new_range_del_iter == nullptr || new_range_del_iter->empty()) { |
429 | 38.3k | delete new_range_del_iter; |
430 | 38.3k | *range_del_iter = nullptr; |
431 | 38.3k | } else { |
432 | 2.87k | *range_del_iter = std::make_unique<TruncatedRangeDelIterator>( |
433 | 2.87k | std::unique_ptr<FragmentedRangeTombstoneIterator>( |
434 | 2.87k | new_range_del_iter), |
435 | 2.87k | &icomparator, &file_meta.smallest, &file_meta.largest); |
436 | 2.87k | } |
437 | 41.1k | } |
438 | 68.5k | if (range_del_agg != nullptr) { |
439 | 24.7k | if (range_del_agg->AddFile(file_meta.fd.GetNumber())) { |
440 | 24.7k | std::unique_ptr<FragmentedRangeTombstoneIterator> new_range_del_iter( |
441 | 24.7k | static_cast<FragmentedRangeTombstoneIterator*>( |
442 | 24.7k | table_reader->NewRangeTombstoneIterator(options))); |
443 | 24.7k | if (new_range_del_iter != nullptr) { |
444 | 0 | s = new_range_del_iter->status(); |
445 | 0 | } |
446 | 24.7k | if (s.ok()) { |
447 | 24.7k | const InternalKey* smallest = &file_meta.smallest; |
448 | 24.7k | const InternalKey* largest = &file_meta.largest; |
449 | 24.7k | if (smallest_compaction_key != nullptr) { |
450 | 5.43k | smallest = smallest_compaction_key; |
451 | 5.43k | } |
452 | 24.7k | if (largest_compaction_key != nullptr) { |
453 | 5.43k | largest = largest_compaction_key; |
454 | 5.43k | } |
455 | 24.7k | range_del_agg->AddTombstones(std::move(new_range_del_iter), smallest, |
456 | 24.7k | largest); |
457 | 24.7k | } |
458 | 24.7k | } |
459 | 24.7k | } |
460 | 68.5k | } |
461 | | |
462 | 68.5k | if (handle != nullptr) { |
463 | 0 | cache_.Release(handle); |
464 | 0 | } |
465 | | // Range-del processing is done and s is final. Hand the ephemeral reader's |
466 | | // lifetime to the returned iterator; if s is non-OK, leave it for |
467 | | // ephemeral_reader's destructor. RegisterCleanup gets the raw pointer before |
468 | | // release(), so if its allocation throws the reader is still owned by |
469 | | // ephemeral_reader and freed on unwind. |
470 | 68.5k | if (s.ok() && ephemeral_reader && result != nullptr) { |
471 | 0 | TableReader* raw = ephemeral_reader.get(); |
472 | 0 | result->RegisterCleanup( |
473 | 0 | [](void* arg1, void* /*arg2*/) { |
474 | 0 | delete static_cast<TableReader*>(arg1); |
475 | 0 | }, |
476 | 0 | raw, nullptr); |
477 | | // release() returns raw, which the cleanup now owns; drop the unique_ptr's |
478 | | // ownership so the reader isn't freed twice. |
479 | 0 | [[maybe_unused]] TableReader* released = ephemeral_reader.release(); |
480 | 0 | assert(released == raw); |
481 | 0 | } |
482 | 68.5k | if (!s.ok()) { |
483 | | // Today result is always null here: it is only set when s was OK, and the |
484 | | // only later change to s, new_range_del_iter->status(), is always OK for a |
485 | | // FragmentedRangeTombstoneIterator. The assert documents that; the cleanup |
486 | | // still disposes of any stray result before the error iterator replaces it. |
487 | 0 | assert(result == nullptr); |
488 | 0 | if (result != nullptr) { |
489 | 0 | if (arena != nullptr) { |
490 | 0 | result->~InternalIterator(); |
491 | 0 | } else { |
492 | 0 | delete result; |
493 | 0 | } |
494 | 0 | result = nullptr; |
495 | 0 | } |
496 | 0 | result = NewErrorInternalIterator<Slice>(s, arena); |
497 | 0 | } |
498 | 68.5k | return result; |
499 | 68.5k | } |
500 | | |
501 | | Status TableCache::GetRangeTombstoneIterator( |
502 | | const ReadOptions& options, |
503 | | const InternalKeyComparator& internal_comparator, |
504 | | const FileMetaData& file_meta, const MutableCFOptions& mutable_cf_options, |
505 | 0 | std::unique_ptr<FragmentedRangeTombstoneIterator>* out_iter) { |
506 | 0 | assert(out_iter); |
507 | 0 | Status s; |
508 | 0 | TableReader* t = nullptr; |
509 | 0 | TypedHandle* handle = nullptr; |
510 | 0 | s = FindTable(options, file_options_, internal_comparator, file_meta, &handle, |
511 | 0 | mutable_cf_options, &t); |
512 | 0 | if (s.ok()) { |
513 | | // Note: NewRangeTombstoneIterator could return nullptr |
514 | 0 | out_iter->reset(t->NewRangeTombstoneIterator(options)); |
515 | 0 | } |
516 | 0 | if (handle) { |
517 | 0 | if (*out_iter) { |
518 | 0 | cache_.RegisterReleaseAsCleanup(handle, **out_iter); |
519 | 0 | } else { |
520 | 0 | cache_.Release(handle); |
521 | 0 | } |
522 | 0 | } |
523 | 0 | return s; |
524 | 0 | } |
525 | | |
526 | | uint64_t TableCache::CreateRowCacheKeyPrefix(const ReadOptions& options, |
527 | | const FileDescriptor& fd, |
528 | | const Slice& internal_key, |
529 | | GetContext* get_context, |
530 | 0 | IterKey& row_cache_key) { |
531 | 0 | uint64_t fd_number = fd.GetNumber(); |
532 | | // We use the user key as cache key instead of the internal key, |
533 | | // otherwise the whole cache would be invalidated every time the |
534 | | // sequence key increases. However, to support caching snapshot |
535 | | // reads, we append a sequence number (incremented by 1 to |
536 | | // distinguish from 0) other than internal_key seq no |
537 | | // to determine row cache entry visibility. |
538 | | // If the snapshot is larger than the largest seqno in the file, |
539 | | // all data should be exposed to the snapshot, so we treat it |
540 | | // the same as there is no snapshot. The exception is that if |
541 | | // a seq-checking callback is registered, some internal keys |
542 | | // may still be filtered out. |
543 | 0 | uint64_t cache_entry_seq_no = 0; |
544 | | |
545 | | // Maybe we can include the whole file ifsnapshot == fd.largest_seqno. |
546 | 0 | if (options.snapshot != nullptr && |
547 | 0 | (get_context->has_callback() || |
548 | 0 | static_cast_with_check<const SnapshotImpl>(options.snapshot) |
549 | 0 | ->GetSequenceNumber() <= fd.largest_seqno)) { |
550 | | // We should consider to use options.snapshot->GetSequenceNumber() |
551 | | // instead of GetInternalKeySeqno(k), which will make the code |
552 | | // easier to understand. |
553 | 0 | const MetadataReadBounds* metadata_read_bounds = |
554 | 0 | get_context->metadata_read_bounds(); |
555 | 0 | cache_entry_seq_no = 1 + (metadata_read_bounds != nullptr |
556 | 0 | ? metadata_read_bounds->read_snapshot_seq |
557 | 0 | : GetInternalKeySeqno(internal_key)); |
558 | 0 | } |
559 | | |
560 | | // Compute row cache key. |
561 | 0 | row_cache_key.TrimAppend(row_cache_key.Size(), row_cache_id_.data(), |
562 | 0 | row_cache_id_.size()); |
563 | 0 | AppendVarint64(&row_cache_key, fd_number); |
564 | 0 | AppendVarint64(&row_cache_key, cache_entry_seq_no); |
565 | | |
566 | | // Provide a sequence number for callback checking on cache hit. |
567 | | // As cache_entry_seq_no starts at 1, decrease it's value by 1 to get |
568 | | // a sequence number align with get context's logic. |
569 | 0 | return cache_entry_seq_no == 0 ? 0 : cache_entry_seq_no - 1; |
570 | 0 | } |
571 | | |
572 | | bool TableCache::GetFromRowCache(const Slice& user_key, IterKey& row_cache_key, |
573 | | size_t prefix_size, GetContext* get_context, |
574 | 0 | Status* read_status, SequenceNumber seq_no) { |
575 | 0 | bool found = false; |
576 | |
|
577 | 0 | row_cache_key.TrimAppend(prefix_size, user_key.data(), user_key.size()); |
578 | 0 | RowCacheInterface row_cache{ioptions_.row_cache.get()}; |
579 | 0 | if (auto row_handle = row_cache.Lookup(row_cache_key.GetUserKey())) { |
580 | | // Cleanable routine to release the cache entry |
581 | 0 | Cleanable value_pinner; |
582 | | // If it comes here value is located on the cache. |
583 | | // found_row_cache_entry points to the value on cache, |
584 | | // and value_pinner has cleanup procedure for the cached entry. |
585 | | // After replayGetContextLog() returns, get_context.pinnable_slice_ |
586 | | // will point to cache entry buffer (or a copy based on that) and |
587 | | // cleanup routine under value_pinner will be delegated to |
588 | | // get_context.pinnable_slice_. Cache entry is released when |
589 | | // get_context.pinnable_slice_ is reset. |
590 | 0 | row_cache.RegisterReleaseAsCleanup(row_handle, value_pinner); |
591 | | // If row cache hit, knowing cache key is the same to row_cache_key, |
592 | | // can use row_cache_key's seq no to construct InternalKey. |
593 | 0 | *read_status = replayGetContextLog(*row_cache.Value(row_handle), user_key, |
594 | 0 | get_context, &value_pinner, seq_no); |
595 | 0 | RecordTick(ioptions_.stats, ROW_CACHE_HIT); |
596 | 0 | found = true; |
597 | 0 | } else { |
598 | 0 | RecordTick(ioptions_.stats, ROW_CACHE_MISS); |
599 | 0 | } |
600 | 0 | return found; |
601 | 0 | } |
602 | | |
603 | | void TableCache::UpdateRangeTombstoneSeqnums( |
604 | | const ReadOptions& options, TableReader* t, |
605 | 0 | MultiGetContext::Range& table_range) { |
606 | 0 | std::unique_ptr<FragmentedRangeTombstoneIterator> latest_range_del_iter; |
607 | 0 | bool track_newer_versions = false; |
608 | 0 | if (!table_range.empty()) { |
609 | 0 | ReadCallback* callback = table_range.begin()->get_context->read_callback(); |
610 | 0 | if (callback != nullptr && callback->GetMetadataReadBounds() != nullptr) { |
611 | 0 | for (auto iter = table_range.begin(); iter != table_range.end(); ++iter) { |
612 | 0 | if (iter->get_context->NeedToTrackNewerVersions()) { |
613 | 0 | track_newer_versions = true; |
614 | 0 | break; |
615 | 0 | } |
616 | 0 | } |
617 | 0 | } |
618 | 0 | } |
619 | 0 | if (track_newer_versions) { |
620 | 0 | SequenceNumber latest_range_del_read_seq = 0; |
621 | 0 | for (auto iter = table_range.begin(); iter != table_range.end(); ++iter) { |
622 | 0 | if (iter->get_context->NeedToTrackNewerVersions()) { |
623 | 0 | latest_range_del_read_seq = std::max(latest_range_del_read_seq, |
624 | 0 | GetInternalKeySeqno(iter->ikey)); |
625 | 0 | } |
626 | 0 | } |
627 | 0 | if (latest_range_del_read_seq > 0) { |
628 | 0 | latest_range_del_iter.reset(t->NewRangeTombstoneIterator( |
629 | 0 | latest_range_del_read_seq, options.timestamp)); |
630 | 0 | } |
631 | 0 | } |
632 | |
|
633 | 0 | std::unique_ptr<FragmentedRangeTombstoneIterator> range_del_iter; |
634 | 0 | if (!options.ignore_range_deletions) { |
635 | 0 | range_del_iter.reset(t->NewRangeTombstoneIterator(options)); |
636 | 0 | } |
637 | 0 | if (range_del_iter != nullptr || latest_range_del_iter != nullptr) { |
638 | 0 | for (auto iter = table_range.begin(); iter != table_range.end(); ++iter) { |
639 | 0 | if (latest_range_del_iter != nullptr && |
640 | 0 | iter->get_context->NeedToTrackNewerVersions()) { |
641 | 0 | const SequenceNumber covering_seq = |
642 | 0 | latest_range_del_iter->MaxCoveringTombstoneSeqnum( |
643 | 0 | iter->ukey_with_ts, iter->get_context->read_callback()); |
644 | 0 | if (covering_seq != 0) { |
645 | 0 | iter->get_context->RecordNewerVersionIfNeeded(covering_seq, |
646 | 0 | kTypeRangeDeletion); |
647 | 0 | } |
648 | 0 | } |
649 | 0 | if (range_del_iter == nullptr) { |
650 | 0 | continue; |
651 | 0 | } |
652 | 0 | SequenceNumber* max_covering_tombstone_seq = |
653 | 0 | iter->get_context->max_covering_tombstone_seq(); |
654 | 0 | SequenceNumber seq = |
655 | 0 | range_del_iter->MaxCoveringTombstoneSeqnum(iter->ukey_with_ts); |
656 | 0 | if (seq > *max_covering_tombstone_seq) { |
657 | 0 | *max_covering_tombstone_seq = seq; |
658 | 0 | if (iter->get_context->NeedTimestamp()) { |
659 | 0 | iter->get_context->SetTimestampFromRangeTombstone( |
660 | 0 | range_del_iter->timestamp()); |
661 | 0 | } |
662 | 0 | } |
663 | 0 | } |
664 | 0 | } |
665 | 0 | } |
666 | | |
667 | | Status TableCache::MultiGetFilter( |
668 | | const ReadOptions& options, |
669 | | const InternalKeyComparator& internal_comparator, |
670 | | const FileMetaData& file_meta, const MutableCFOptions& mutable_cf_options, |
671 | | HistogramImpl* file_read_hist, int level, |
672 | 0 | MultiGetContext::Range* mget_range, TypedHandle** handle) { |
673 | 0 | assert(*handle == nullptr); |
674 | 0 | IterKey row_cache_key; |
675 | 0 | std::string row_cache_entry_buffer; |
676 | | |
677 | | // Check if we need to use the row cache. If yes, then we cannot do the |
678 | | // filtering here, since the filtering needs to happen after the row cache |
679 | | // lookup. |
680 | 0 | KeyContext& first_key = *mget_range->begin(); |
681 | 0 | bool track_newer_versions = false; |
682 | 0 | ReadCallback* callback = first_key.get_context->read_callback(); |
683 | 0 | if (callback != nullptr && callback->GetMetadataReadBounds() != nullptr) { |
684 | 0 | for (auto iter = mget_range->begin(); iter != mget_range->end(); ++iter) { |
685 | 0 | if (iter->get_context->NeedToTrackNewerVersions()) { |
686 | 0 | track_newer_versions = true; |
687 | 0 | break; |
688 | 0 | } |
689 | 0 | } |
690 | 0 | } |
691 | 0 | if (ioptions_.row_cache && !first_key.get_context->NeedToReadSequence() && |
692 | 0 | !track_newer_versions) { |
693 | 0 | return Status::NotSupported(); |
694 | 0 | } |
695 | 0 | Status s; |
696 | 0 | TableReader* t = nullptr; |
697 | 0 | MultiGetContext::Range tombstone_range(*mget_range, mget_range->begin(), |
698 | 0 | mget_range->end()); |
699 | 0 | s = FindTable(options, file_options_, internal_comparator, file_meta, handle, |
700 | 0 | mutable_cf_options, &t, |
701 | 0 | options.read_tier == kBlockCacheTier /* no_io */, |
702 | 0 | file_read_hist, |
703 | 0 | /*skip_filters=*/false, level, |
704 | 0 | true /* prefetch_index_and_filter_in_cache */, |
705 | 0 | /*max_file_size_for_l0_meta_pin=*/0, file_meta.temperature, |
706 | 0 | should_pin_table_handles_); |
707 | 0 | if (s.ok()) { |
708 | 0 | s = t->MultiGetFilter(options, mutable_cf_options.prefix_extractor.get(), |
709 | 0 | mget_range); |
710 | 0 | } |
711 | 0 | if (s.ok() && (!options.ignore_range_deletions || track_newer_versions)) { |
712 | | // Update the range tombstone sequence numbers for the keys here |
713 | | // as TableCache::MultiGet may or may not be called, and even if it |
714 | | // is, it may be called with fewer keys in the rangedue to filtering. |
715 | 0 | UpdateRangeTombstoneSeqnums(options, t, tombstone_range); |
716 | 0 | } |
717 | 0 | if (mget_range->empty() && *handle) { |
718 | 0 | cache_.Release(*handle); |
719 | 0 | *handle = nullptr; |
720 | 0 | } |
721 | |
|
722 | 0 | return s; |
723 | 0 | } |
724 | | |
725 | | Status TableCache::GetTableProperties( |
726 | | const FileOptions& file_options, const ReadOptions& read_options, |
727 | | const InternalKeyComparator& internal_comparator, |
728 | | const FileMetaData& file_meta, |
729 | | std::shared_ptr<const TableProperties>* properties, |
730 | 152k | const MutableCFOptions& mutable_cf_options, bool no_io) { |
731 | 152k | TypedHandle* table_handle = nullptr; |
732 | 152k | TableReader* table = nullptr; |
733 | 152k | Status s = |
734 | 152k | FindTable(read_options, file_options, internal_comparator, file_meta, |
735 | 152k | &table_handle, mutable_cf_options, &table, no_io); |
736 | 152k | if (!s.ok()) { |
737 | 0 | return s; |
738 | 0 | } |
739 | 152k | *properties = table->GetTableProperties(); |
740 | 152k | if (table_handle) { |
741 | 0 | cache_.Release(table_handle); |
742 | 0 | } |
743 | 152k | return s; |
744 | 152k | } |
745 | | |
746 | | Status TableCache::ApproximateKeyAnchors( |
747 | | const ReadOptions& ro, const InternalKeyComparator& internal_comparator, |
748 | | const FileMetaData& file_meta, const MutableCFOptions& mutable_cf_options, |
749 | | |
750 | 0 | std::vector<TableReader::Anchor>& anchors) { |
751 | 0 | Status s; |
752 | 0 | TableReader* t = nullptr; |
753 | 0 | TypedHandle* handle = nullptr; |
754 | 0 | s = FindTable(ro, file_options_, internal_comparator, file_meta, &handle, |
755 | 0 | mutable_cf_options, &t); |
756 | 0 | if (s.ok() && t != nullptr) { |
757 | 0 | s = t->ApproximateKeyAnchors(ro, anchors); |
758 | 0 | } |
759 | 0 | if (handle != nullptr) { |
760 | 0 | cache_.Release(handle); |
761 | 0 | } |
762 | 0 | return s; |
763 | 0 | } |
764 | | |
765 | | size_t TableCache::GetMemoryUsageByTableReader( |
766 | | const FileOptions& file_options, const ReadOptions& read_options, |
767 | | const InternalKeyComparator& internal_comparator, |
768 | 0 | const FileMetaData& file_meta, const MutableCFOptions& mutable_cf_options) { |
769 | 0 | TypedHandle* table_handle = nullptr; |
770 | 0 | TableReader* table = nullptr; |
771 | 0 | Status s = |
772 | 0 | FindTable(read_options, file_options, internal_comparator, file_meta, |
773 | 0 | &table_handle, mutable_cf_options, &table, true /* no_io */); |
774 | 0 | if (!s.ok()) { |
775 | 0 | return 0; |
776 | 0 | } |
777 | 0 | auto ret = table->ApproximateMemoryUsage(); |
778 | 0 | if (table_handle) { |
779 | 0 | cache_.Release(table_handle); |
780 | 0 | } |
781 | 0 | return ret; |
782 | 0 | } |
783 | | |
784 | 22.3k | void TableCache::Evict(Cache* cache, uint64_t file_number) { |
785 | 22.3k | cache->Erase(GetSliceForFileNumber(&file_number)); |
786 | 22.3k | } |
787 | | |
788 | | uint64_t TableCache::ApproximateOffsetOf( |
789 | | const ReadOptions& read_options, const Slice& key, |
790 | | const FileMetaData& file_meta, TableReaderCaller caller, |
791 | | const InternalKeyComparator& internal_comparator, |
792 | 0 | const MutableCFOptions& mutable_cf_options) { |
793 | 0 | uint64_t result = 0; |
794 | 0 | TableReader* table_reader = nullptr; |
795 | 0 | TypedHandle* table_handle = nullptr; |
796 | 0 | Status s = |
797 | 0 | FindTable(read_options, file_options_, internal_comparator, file_meta, |
798 | 0 | &table_handle, mutable_cf_options, &table_reader); |
799 | |
|
800 | 0 | if (s.ok() && table_reader != nullptr) { |
801 | 0 | result = table_reader->ApproximateOffsetOf(read_options, key, caller); |
802 | 0 | } |
803 | 0 | if (table_handle != nullptr) { |
804 | 0 | cache_.Release(table_handle); |
805 | 0 | } |
806 | |
|
807 | 0 | return result; |
808 | 0 | } |
809 | | |
810 | | uint64_t TableCache::ApproximateSize( |
811 | | const ReadOptions& read_options, const Slice& start, const Slice& end, |
812 | | const FileMetaData& file_meta, TableReaderCaller caller, |
813 | | const InternalKeyComparator& internal_comparator, |
814 | 0 | const MutableCFOptions& mutable_cf_options) { |
815 | 0 | uint64_t result = 0; |
816 | 0 | TableReader* table_reader = nullptr; |
817 | 0 | TypedHandle* table_handle = nullptr; |
818 | 0 | Status s = |
819 | 0 | FindTable(read_options, file_options_, internal_comparator, file_meta, |
820 | 0 | &table_handle, mutable_cf_options, &table_reader); |
821 | |
|
822 | 0 | if (s.ok() && table_reader != nullptr) { |
823 | 0 | result = table_reader->ApproximateSize(read_options, start, end, caller); |
824 | 0 | } |
825 | 0 | if (table_handle != nullptr) { |
826 | 0 | cache_.Release(table_handle); |
827 | 0 | } |
828 | |
|
829 | 0 | return result; |
830 | 0 | } |
831 | | |
832 | | void TableCache::ReleaseObsolete(Cache* cache, uint64_t file_number, |
833 | | Cache::Handle* h, |
834 | 123k | uint32_t uncache_aggressiveness) { |
835 | 123k | CacheInterface typed_cache(cache); |
836 | 123k | TypedHandle* table_handle = reinterpret_cast<TypedHandle*>(h); |
837 | 123k | if (table_handle == nullptr) { |
838 | 0 | table_handle = typed_cache.Lookup(GetSliceForFileNumber(&file_number)); |
839 | 0 | } |
840 | 123k | if (table_handle != nullptr) { |
841 | 123k | TableReader* table_reader = typed_cache.Value(table_handle); |
842 | 123k | table_reader->MarkObsolete(uncache_aggressiveness); |
843 | | // Mark the entry Invisible so that if concurrent readers hold references, |
844 | | // the entry will be erased when the last reference is released. |
845 | 123k | cache->Erase(GetSliceForFileNumber(&file_number)); |
846 | 123k | typed_cache.Release(table_handle); |
847 | 123k | } |
848 | 123k | } |
849 | | |
850 | | } // namespace ROCKSDB_NAMESPACE |