version_builder.cc 28.8 KB
Newer Older
1
//  Copyright (c) 2011-present, Facebook, Inc.  All rights reserved.
Siying Dong's avatar
Siying Dong committed
2
3
4
//  This source code is licensed under both the GPLv2 (found in the
//  COPYING file in the root directory) and Apache 2.0 License
//  (found in the LICENSE.Apache file in the root directory).
5
6
7
8
9
10
11
12
//
// Copyright (c) 2011 The LevelDB Authors. All rights reserved.
// Use of this source code is governed by a BSD-style license that can be
// found in the LICENSE file. See the AUTHORS file for names of contributors.

#include "db/version_builder.h"

#include <algorithm>
13
#include <atomic>
14
#include <cinttypes>
15
#include <functional>
16
#include <map>
17
#include <memory>
18
#include <set>
19
#include <sstream>
20
#include <thread>
21
22
#include <unordered_map>
#include <unordered_set>
23
#include <utility>
24
25
#include <vector>

26
#include "db/blob/blob_file_meta.h"
27
#include "db/dbformat.h"
28
#include "db/internal_stats.h"
29
30
#include "db/table_cache.h"
#include "db/version_set.h"
Dmitri Smirnov's avatar
Dmitri Smirnov committed
31
#include "port/port.h"
32
#include "table/table_reader.h"
33
#include "util/string_util.h"
34

35
namespace ROCKSDB_NAMESPACE {
36
37

bool NewestFirstBySeqNo(FileMetaData* a, FileMetaData* b) {
38
39
  if (a->fd.largest_seqno != b->fd.largest_seqno) {
    return a->fd.largest_seqno > b->fd.largest_seqno;
40
  }
41
42
  if (a->fd.smallest_seqno != b->fd.smallest_seqno) {
    return a->fd.smallest_seqno > b->fd.smallest_seqno;
43
  }
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
  // Break ties by file number
  return a->fd.GetNumber() > b->fd.GetNumber();
}

namespace {
bool BySmallestKey(FileMetaData* a, FileMetaData* b,
                   const InternalKeyComparator* cmp) {
  int r = cmp->Compare(a->smallest, b->smallest);
  if (r != 0) {
    return (r < 0);
  }
  // Break ties by file number
  return (a->fd.GetNumber() < b->fd.GetNumber());
}
}  // namespace

class VersionBuilder::Rep {
 private:
  // Helper to sort files_ in v
  // kLevel0 -- NewestFirstBySeqNo
  // kLevelNon0 -- BySmallestKey
  struct FileComparator {
66
    enum SortMethod { kLevel0 = 0, kLevelNon0 = 1, } sort_method;
67
68
    const InternalKeyComparator* internal_comparator;

69
70
    FileComparator() : internal_comparator(nullptr) {}

71
72
73
74
75
76
77
78
79
80
81
82
83
    bool operator()(FileMetaData* f1, FileMetaData* f2) const {
      switch (sort_method) {
        case kLevel0:
          return NewestFirstBySeqNo(f1, f2);
        case kLevelNon0:
          return BySmallestKey(f1, f2, internal_comparator);
      }
      assert(false);
      return false;
    }
  };

  struct LevelState {
84
85
86
    std::unordered_set<uint64_t> deleted_files;
    // Map from file number to file meta data.
    std::unordered_map<uint64_t, FileMetaData*> added_files;
87
88
  };

89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
  class BlobFileMetaDataDelta {
   public:
    bool IsEmpty() const {
      return !shared_meta_ && !additional_garbage_count_ &&
             !additional_garbage_bytes_;
    }

    std::shared_ptr<SharedBlobFileMetaData> GetSharedMeta() const {
      return shared_meta_;
    }

    uint64_t GetAdditionalGarbageCount() const {
      return additional_garbage_count_;
    }

    uint64_t GetAdditionalGarbageBytes() const {
      return additional_garbage_bytes_;
    }

    void SetSharedMeta(std::shared_ptr<SharedBlobFileMetaData> shared_meta) {
      assert(!shared_meta_);
      assert(shared_meta);

      shared_meta_ = std::move(shared_meta);
    }

    void AddGarbage(uint64_t count, uint64_t bytes) {
      additional_garbage_count_ += count;
      additional_garbage_bytes_ += bytes;
    }

   private:
    std::shared_ptr<SharedBlobFileMetaData> shared_meta_;
    uint64_t additional_garbage_count_ = 0;
    uint64_t additional_garbage_bytes_ = 0;
  };

126
  const FileOptions& file_options_;
127
  const ImmutableCFOptions* const ioptions_;
128
129
  TableCache* table_cache_;
  VersionStorageInfo* base_vstorage_;
130
  VersionSet* version_set_;
131
  int num_levels_;
132
  LevelState* levels_;
133
  // Store sizes of levels larger than num_levels_. We do this instead of
134
135
136
  // storing them in levels_ to avoid regression in case there are no files
  // on invalid levels. The version is not consistent if in the end the files
  // on invalid levels don't cancel out.
137
  std::unordered_map<int, size_t> invalid_level_sizes_;
138
139
140
  // Whether there are invalid new files or invalid deletion on levels larger
  // than num_levels_.
  bool has_invalid_levels_;
141
142
  // Current levels of table files affected by additions/deletions.
  std::unordered_map<uint64_t, int> table_file_levels_;
143
144
145
  FileComparator level_zero_cmp_;
  FileComparator level_nonzero_cmp_;

146
147
  // Metadata delta for all blob files affected by the series of version edits.
  std::map<uint64_t, BlobFileMetaDataDelta> blob_file_meta_deltas_;
148

149
 public:
150
151
152
  Rep(const FileOptions& file_options, const ImmutableCFOptions* ioptions,
      TableCache* table_cache, VersionStorageInfo* base_vstorage,
      VersionSet* version_set)
153
      : file_options_(file_options),
154
        ioptions_(ioptions),
155
        table_cache_(table_cache),
156
        base_vstorage_(base_vstorage),
157
        version_set_(version_set),
158
159
        num_levels_(base_vstorage->num_levels()),
        has_invalid_levels_(false) {
160
161
    assert(ioptions_);

162
    levels_ = new LevelState[num_levels_];
163
164
165
166
167
168
169
    level_zero_cmp_.sort_method = FileComparator::kLevel0;
    level_nonzero_cmp_.sort_method = FileComparator::kLevelNon0;
    level_nonzero_cmp_.internal_comparator =
        base_vstorage_->InternalComparator();
  }

  ~Rep() {
170
    for (int level = 0; level < num_levels_; level++) {
171
172
      const auto& added = levels_[level].added_files;
      for (auto& pair : added) {
173
        UnrefFile(pair.second);
174
175
176
177
178
179
      }
    }

    delete[] levels_;
  }

180
181
182
183
184
185
186
187
188
189
190
191
  void UnrefFile(FileMetaData* f) {
    f->refs--;
    if (f->refs <= 0) {
      if (f->table_reader_handle) {
        assert(table_cache_ != nullptr);
        table_cache_->ReleaseHandle(f->table_reader_handle);
        f->table_reader_handle = nullptr;
      }
      delete f;
    }
  }

192
193
194
195
196
197
  bool IsBlobFileInVersion(uint64_t blob_file_number) const {
    auto delta_it = blob_file_meta_deltas_.find(blob_file_number);
    if (delta_it != blob_file_meta_deltas_.end()) {
      if (delta_it->second.GetSharedMeta()) {
        return true;
      }
198
199
200
201
202
203
204
205
    }

    assert(base_vstorage_);

    const auto& base_blob_files = base_vstorage_->GetBlobFiles();

    auto base_it = base_blob_files.find(blob_file_number);
    if (base_it != base_blob_files.end()) {
206
207
      assert(base_it->second);
      assert(base_it->second->GetSharedMeta());
208

209
      return true;
210
211
    }

212
    return false;
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
  }

  Status CheckConsistencyOfOldestBlobFileReference(
      const VersionStorageInfo* vstorage, uint64_t blob_file_number) const {
    assert(vstorage);

    // TODO: remove this check once we actually start recoding metadata for
    // blob files in the MANIFEST.
    if (vstorage->GetBlobFiles().empty()) {
      return Status::OK();
    }

    if (blob_file_number == kInvalidBlobFileNumber) {
      return Status::OK();
    }

229
    if (!IsBlobFileInVersion(blob_file_number)) {
230
231
232
233
234
235
236
237
238
239
      std::ostringstream oss;
      oss << "Blob file #" << blob_file_number
          << " is not part of this version";

      return Status::Corruption("VersionBuilder", oss.str());
    }

    return Status::OK();
  }

240
  Status CheckConsistency(VersionStorageInfo* vstorage) {
241
242
243
244
#ifdef NDEBUG
    if (!vstorage->force_consistency_checks()) {
      // Dont run consistency checks in release mode except if
      // explicitly asked to
245
      return Status::OK();
246
247
    }
#endif
248
249
250
    // Make sure the files are sorted correctly and that the oldest blob file
    // reference for each table file points to a valid blob file in this
    // version.
251
    for (int level = 0; level < num_levels_; level++) {
252
      auto& level_files = vstorage->LevelFiles(level);
253
254
255
256
257
258
259
260
261
262
263
264

      if (level_files.empty()) {
        continue;
      }

      assert(level_files[0]);
      Status s = CheckConsistencyOfOldestBlobFileReference(
          vstorage, level_files[0]->oldest_blob_file_number);
      if (!s.ok()) {
        return s;
      }

265
      for (size_t i = 1; i < level_files.size(); i++) {
266
267
268
269
270
271
272
        assert(level_files[i]);
        s = CheckConsistencyOfOldestBlobFileReference(
            vstorage, level_files[i]->oldest_blob_file_number);
        if (!s.ok()) {
          return s;
        }

273
274
        auto f1 = level_files[i - 1];
        auto f2 = level_files[i];
275
        if (level == 0) {
276
#ifndef NDEBUG
277
278
          auto pair = std::make_pair(&f1, &f2);
          TEST_SYNC_POINT_CALLBACK("VersionBuilder::CheckConsistency0", &pair);
279
#endif
280
          if (!level_zero_cmp_(f1, f2)) {
281
            return Status::Corruption("L0 files are not sorted properly");
282
283
          }

284
          if (f2->fd.smallest_seqno == f2->fd.largest_seqno) {
285
            // This is an external file that we ingested
286
287
            SequenceNumber external_file_seqno = f2->fd.smallest_seqno;
            if (!(external_file_seqno < f1->fd.largest_seqno ||
288
                  external_file_seqno == 0)) {
289
290
291
292
293
294
295
              return Status::Corruption(
                  "L0 file with seqno " +
                  NumberToString(f1->fd.smallest_seqno) + " " +
                  NumberToString(f1->fd.largest_seqno) +
                  " vs. file with global_seqno" +
                  NumberToString(external_file_seqno) + " with fileNumber " +
                  NumberToString(f1->fd.GetNumber()));
296
            }
297
          } else if (f1->fd.smallest_seqno <= f2->fd.smallest_seqno) {
298
299
300
301
302
303
304
            return Status::Corruption(
                "L0 files seqno " + NumberToString(f1->fd.smallest_seqno) +
                " " + NumberToString(f1->fd.largest_seqno) + " " +
                NumberToString(f1->fd.GetNumber()) + " vs. " +
                NumberToString(f2->fd.smallest_seqno) + " " +
                NumberToString(f2->fd.largest_seqno) + " " +
                NumberToString(f2->fd.GetNumber()));
305
          }
306
        } else {
307
308
309
310
#ifndef NDEBUG
          auto pair = std::make_pair(&f1, &f2);
          TEST_SYNC_POINT_CALLBACK("VersionBuilder::CheckConsistency1", &pair);
#endif
311
          if (!level_nonzero_cmp_(f1, f2)) {
312
313
            return Status::Corruption("L" + NumberToString(level) +
                                      " files are not sorted properly");
314
          }
315
316
317
318

          // Make sure there is no overlap in levels > 0
          if (vstorage->InternalComparator()->Compare(f1->largest,
                                                      f2->smallest) >= 0) {
319
320
321
322
            return Status::Corruption(
                "L" + NumberToString(level) + " have overlapping ranges " +
                (f1->largest).DebugString(true) + " vs. " +
                (f2->smallest).DebugString(true));
323
324
325
326
          }
        }
      }
    }
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343

    // Make sure that all blob files in the version have non-garbage data.
    const auto& blob_files = vstorage->GetBlobFiles();
    for (const auto& pair : blob_files) {
      const auto& blob_file_meta = pair.second;
      assert(blob_file_meta);

      if (blob_file_meta->GetGarbageBlobCount() >=
          blob_file_meta->GetTotalBlobCount()) {
        std::ostringstream oss;
        oss << "Blob file #" << blob_file_meta->GetBlobFileNumber()
            << " consists entirely of garbage";

        return Status::Corruption("VersionBuilder", oss.str());
      }
    }

344
345
346
347
    Status ret_s;
    TEST_SYNC_POINT_CALLBACK("VersionBuilder::CheckConsistencyBeforeReturn",
                             &ret_s);
    return ret_s;
348
349
  }

350
  bool CheckConsistencyForNumLevels() const {
351
352
353
354
    // Make sure there are no files on or beyond num_levels().
    if (has_invalid_levels_) {
      return false;
    }
355
356
357
358

    for (const auto& pair : invalid_level_sizes_) {
      const size_t level_size = pair.second;
      if (level_size != 0) {
359
360
361
        return false;
      }
    }
362

363
364
365
    return true;
  }

366
367
368
  Status ApplyBlobFileAddition(const BlobFileAddition& blob_file_addition) {
    const uint64_t blob_file_number = blob_file_addition.GetBlobFileNumber();

369
    if (IsBlobFileInVersion(blob_file_number)) {
370
371
372
373
374
375
      std::ostringstream oss;
      oss << "Blob file #" << blob_file_number << " already added";

      return Status::Corruption("VersionBuilder", oss.str());
    }

376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
    // Note: we use C++11 for now but in C++14, this could be done in a more
    // elegant way using generalized lambda capture.
    VersionSet* const vs = version_set_;
    const ImmutableCFOptions* const ioptions = ioptions_;

    auto deleter = [vs, ioptions](SharedBlobFileMetaData* shared_meta) {
      if (vs) {
        assert(ioptions);
        assert(!ioptions->cf_paths.empty());
        assert(shared_meta);

        vs->AddObsoleteBlobFile(shared_meta->GetBlobFileNumber(),
                                ioptions->cf_paths.front().path);
      }

      delete shared_meta;
    };

394
    auto shared_meta = SharedBlobFileMetaData::Create(
395
396
397
        blob_file_number, blob_file_addition.GetTotalBlobCount(),
        blob_file_addition.GetTotalBlobBytes(),
        blob_file_addition.GetChecksumMethod(),
398
        blob_file_addition.GetChecksumValue(), deleter);
399

400
401
    blob_file_meta_deltas_[blob_file_number].SetSharedMeta(
        std::move(shared_meta));
402
403
404
405
406
407
408

    return Status::OK();
  }

  Status ApplyBlobFileGarbage(const BlobFileGarbage& blob_file_garbage) {
    const uint64_t blob_file_number = blob_file_garbage.GetBlobFileNumber();

409
    if (!IsBlobFileInVersion(blob_file_number)) {
410
411
412
413
414
415
      std::ostringstream oss;
      oss << "Blob file #" << blob_file_number << " not found";

      return Status::Corruption("VersionBuilder", oss.str());
    }

416
417
418
    blob_file_meta_deltas_[blob_file_number].AddGarbage(
        blob_file_garbage.GetGarbageBlobCount(),
        blob_file_garbage.GetGarbageBlobBytes());
419
420
421
422

    return Status::OK();
  }

423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
  int GetCurrentLevelForTableFile(uint64_t file_number) const {
    auto it = table_file_levels_.find(file_number);
    if (it != table_file_levels_.end()) {
      return it->second;
    }

    assert(base_vstorage_);
    return base_vstorage_->GetFileLocation(file_number).GetLevel();
  }

  Status ApplyFileDeletion(int level, uint64_t file_number) {
    assert(level != VersionStorageInfo::FileLocation::Invalid().GetLevel());

    const int current_level = GetCurrentLevelForTableFile(file_number);

    if (level != current_level) {
      if (level >= num_levels_) {
        has_invalid_levels_ = true;
      }

      std::ostringstream oss;
      oss << "Cannot delete table file #" << file_number << " from level "
          << level << " since it is ";
      if (current_level ==
          VersionStorageInfo::FileLocation::Invalid().GetLevel()) {
        oss << "not in the LSM tree";
      } else {
        oss << "on level " << current_level;
      }

      return Status::Corruption("VersionBuilder", oss.str());
    }

    if (level >= num_levels_) {
      assert(invalid_level_sizes_[level] > 0);
      --invalid_level_sizes_[level];

      table_file_levels_[file_number] =
          VersionStorageInfo::FileLocation::Invalid().GetLevel();

      return Status::OK();
    }

    auto& level_state = levels_[level];

    auto& add_files = level_state.added_files;
    auto add_it = add_files.find(file_number);
    if (add_it != add_files.end()) {
      UnrefFile(add_it->second);
      add_files.erase(add_it);
    } else {
      auto& del_files = level_state.deleted_files;
      assert(del_files.find(file_number) == del_files.end());
      del_files.emplace(file_number);
    }

    table_file_levels_[file_number] =
        VersionStorageInfo::FileLocation::Invalid().GetLevel();

    return Status::OK();
  }

  Status ApplyFileAddition(int level, const FileMetaData& meta) {
    assert(level != VersionStorageInfo::FileLocation::Invalid().GetLevel());

    const uint64_t file_number = meta.fd.GetNumber();

    const int current_level = GetCurrentLevelForTableFile(file_number);

    if (current_level !=
        VersionStorageInfo::FileLocation::Invalid().GetLevel()) {
      if (level >= num_levels_) {
        has_invalid_levels_ = true;
      }

      std::ostringstream oss;
      oss << "Cannot add table file #" << file_number << " to level " << level
          << " since it is already in the LSM tree on level " << current_level;
      return Status::Corruption("VersionBuilder", oss.str());
    }

    if (level >= num_levels_) {
      ++invalid_level_sizes_[level];
      table_file_levels_[file_number] = level;

      return Status::OK();
    }

    auto& level_state = levels_[level];

    auto& del_files = level_state.deleted_files;
    auto del_it = del_files.find(file_number);
    if (del_it != del_files.end()) {
      del_files.erase(del_it);
    } else {
      FileMetaData* const f = new FileMetaData(meta);
      f->refs = 1;

      auto& add_files = level_state.added_files;
      assert(add_files.find(file_number) == add_files.end());
      add_files.emplace(file_number, f);
    }

    table_file_levels_[file_number] = level;

    return Status::OK();
  }

531
  // Apply all of the edits in *edit to the current state.
532
  Status Apply(VersionEdit* edit) {
533
534
535
536
537
    {
      const Status s = CheckConsistency(base_vstorage_);
      if (!s.ok()) {
        return s;
      }
538
    }
539
540

    // Delete files
541
542
543
    for (const auto& deleted_file : edit->GetDeletedFiles()) {
      const int level = deleted_file.first;
      const uint64_t file_number = deleted_file.second;
544

545
546
547
      const Status s = ApplyFileDeletion(level, file_number);
      if (!s.ok()) {
        return s;
548
      }
549
550
551
552
553
    }

    // Add new files
    for (const auto& new_file : edit->GetNewFiles()) {
      const int level = new_file.first;
554
555
556
557
558
      const FileMetaData& meta = new_file.second;

      const Status s = ApplyFileAddition(level, meta);
      if (!s.ok()) {
        return s;
559
      }
560
    }
561
562
563

    // Add new blob files
    for (const auto& blob_file_addition : edit->GetBlobFileAdditions()) {
564
      const Status s = ApplyBlobFileAddition(blob_file_addition);
565
566
567
568
569
570
571
      if (!s.ok()) {
        return s;
      }
    }

    // Increase the amount of garbage for blob files affected by GC
    for (const auto& blob_file_garbage : edit->GetBlobFileGarbages()) {
572
      const Status s = ApplyBlobFileGarbage(blob_file_garbage);
573
574
575
576
577
      if (!s.ok()) {
        return s;
      }
    }

578
    return Status::OK();
579
580
  }

581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
  static std::shared_ptr<BlobFileMetaData> CreateMetaDataForNewBlobFile(
      const BlobFileMetaDataDelta& delta) {
    auto shared_meta = delta.GetSharedMeta();
    assert(shared_meta);

    auto meta = BlobFileMetaData::Create(std::move(shared_meta),
                                         delta.GetAdditionalGarbageCount(),
                                         delta.GetAdditionalGarbageBytes());

    return meta;
  }

  static std::shared_ptr<BlobFileMetaData>
  GetOrCreateMetaDataForExistingBlobFile(
      const std::shared_ptr<BlobFileMetaData>& base_meta,
      const BlobFileMetaDataDelta& delta) {
    assert(base_meta);
    assert(!delta.GetSharedMeta());

    if (delta.IsEmpty()) {
      return base_meta;
    }

    auto shared_meta = base_meta->GetSharedMeta();
    assert(shared_meta);

    auto meta = BlobFileMetaData::Create(
        std::move(shared_meta),
        base_meta->GetGarbageBlobCount() + delta.GetAdditionalGarbageCount(),
        base_meta->GetGarbageBlobBytes() + delta.GetAdditionalGarbageBytes());

    return meta;
  }

615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
  void AddBlobFileIfNeeded(
      VersionStorageInfo* vstorage,
      const std::shared_ptr<BlobFileMetaData>& meta) const {
    assert(vstorage);
    assert(meta);

    if (meta->GetGarbageBlobCount() < meta->GetTotalBlobCount()) {
      vstorage->AddBlobFile(meta);
    }
  }

  // Merge the blob file metadata from the base version with the changes (edits)
  // applied, and save the result into *vstorage.
  void SaveBlobFilesTo(VersionStorageInfo* vstorage) const {
    assert(base_vstorage_);
    assert(vstorage);

    const auto& base_blob_files = base_vstorage_->GetBlobFiles();
    auto base_it = base_blob_files.begin();
    const auto base_it_end = base_blob_files.end();

636
637
    auto delta_it = blob_file_meta_deltas_.begin();
    const auto delta_it_end = blob_file_meta_deltas_.end();
638

639
    while (base_it != base_it_end && delta_it != delta_it_end) {
640
      const uint64_t base_blob_file_number = base_it->first;
641
      const uint64_t delta_blob_file_number = delta_it->first;
642

643
644
645
      if (base_blob_file_number < delta_blob_file_number) {
        const auto& base_meta = base_it->second;
        assert(base_meta);
646
647
648
649
650
651
        assert(base_meta->GetGarbageBlobCount() <
               base_meta->GetTotalBlobCount());

        vstorage->AddBlobFile(base_meta);

        ++base_it;
652
653
654
655
656
      } else if (delta_blob_file_number < base_blob_file_number) {
        // Note: blob file numbers are strictly increasing over time and
        // once blob files get marked obsolete, they never reappear. Thus,
        // this case is not possible.
        assert(false);
657

658
        ++delta_it;
659
      } else {
660
        assert(base_blob_file_number == delta_blob_file_number);
661

662
663
664
665
666
667
        const auto& base_meta = base_it->second;
        const auto& delta = delta_it->second;

        auto meta = GetOrCreateMetaDataForExistingBlobFile(base_meta, delta);

        AddBlobFileIfNeeded(vstorage, meta);
668
669

        ++base_it;
670
        ++delta_it;
671
672
673
674
675
676
677
678
679
680
681
682
      }
    }

    while (base_it != base_it_end) {
      const auto& base_meta = base_it->second;
      assert(base_meta);
      assert(base_meta->GetGarbageBlobCount() < base_meta->GetTotalBlobCount());

      vstorage->AddBlobFile(base_meta);
      ++base_it;
    }

683
684
685
686
687
688
    while (delta_it != delta_it_end) {
      const auto& delta = delta_it->second;

      auto meta = CreateMetaDataForNewBlobFile(delta);

      AddBlobFileIfNeeded(vstorage, meta);
689

690
      ++delta_it;
691
692
693
    }
  }

694
  // Save the current state in *v.
695
696
697
698
699
700
701
702
703
704
  Status SaveTo(VersionStorageInfo* vstorage) {
    Status s = CheckConsistency(base_vstorage_);
    if (!s.ok()) {
      return s;
    }

    s = CheckConsistency(vstorage);
    if (!s.ok()) {
      return s;
    }
705

706
    for (int level = 0; level < num_levels_; level++) {
707
708
709
710
      const auto& cmp = (level == 0) ? level_zero_cmp_ : level_nonzero_cmp_;
      // Merge the set of added files with the set of pre-existing files.
      // Drop any deleted files.  Store the result in *v.
      const auto& base_files = base_vstorage_->LevelFiles(level);
711
712
713
714
715
      const auto& unordered_added_files = levels_[level].added_files;
      vstorage->Reserve(level,
                        base_files.size() + unordered_added_files.size());

      // Sort added files for the level.
716
717
      std::vector<FileMetaData*> added_files;
      added_files.reserve(unordered_added_files.size());
718
719
720
721
      for (const auto& pair : unordered_added_files) {
        added_files.push_back(pair.second);
      }
      std::sort(added_files.begin(), added_files.end(), cmp);
722

723
#ifndef NDEBUG
724
      FileMetaData* prev_added_file = nullptr;
725
      for (const auto& added : added_files) {
726
        if (level > 0 && prev_added_file != nullptr) {
727
          assert(base_vstorage_->InternalComparator()->Compare(
728
                     prev_added_file->smallest, added->smallest) <= 0);
729
        }
730
731
        prev_added_file = added;
      }
732
733
#endif

734
735
736
737
738
739
740
741
742
743
      auto base_iter = base_files.begin();
      auto base_end = base_files.end();
      auto added_iter = added_files.begin();
      auto added_end = added_files.end();
      while (added_iter != added_end || base_iter != base_end) {
        if (base_iter == base_end ||
                (added_iter != added_end && cmp(*added_iter, *base_iter))) {
          MaybeAddFile(vstorage, level, *added_iter++);
        } else {
          MaybeAddFile(vstorage, level, *base_iter++);
744
745
746
747
        }
      }
    }

748
749
    SaveBlobFilesTo(vstorage);

750
751
    s = CheckConsistency(vstorage);
    return s;
752
753
  }

754
755
756
757
  Status LoadTableHandlers(InternalStats* internal_stats, int max_threads,
                           bool prefetch_index_and_filter_in_cache,
                           bool is_initial_load,
                           const SliceTransform* prefix_extractor) {
758
    assert(table_cache_ != nullptr);
759
760
761
762
763
764

    size_t table_cache_capacity = table_cache_->get_cache()->GetCapacity();
    bool always_load = (table_cache_capacity == TableCache::kInfiniteCapacity);
    size_t max_load = port::kMaxSizet;

    if (!always_load) {
765
      // If it is initial loading and not set to always loading all the
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
      // files, we only load up to kInitialLoadLimit files, to limit the
      // time reopening the DB.
      const size_t kInitialLoadLimit = 16;
      size_t load_limit;
      // If the table cache is not 1/4 full, we pin the table handle to
      // file metadata to avoid the cache read costs when reading the file.
      // The downside of pinning those files is that LRU won't be followed
      // for those files. This doesn't matter much because if number of files
      // of the DB excceeds table cache capacity, eventually no table reader
      // will be pinned and LRU will be followed.
      if (is_initial_load) {
        load_limit = std::min(kInitialLoadLimit, table_cache_capacity / 4);
      } else {
        load_limit = table_cache_capacity / 4;
      }

      size_t table_cache_usage = table_cache_->get_cache()->GetUsage();
      if (table_cache_usage >= load_limit) {
784
785
        // TODO (yanqin) find a suitable status code.
        return Status::OK();
786
787
788
789
790
      } else {
        max_load = load_limit - table_cache_usage;
      }
    }

791
792
    // <file metadata, level>
    std::vector<std::pair<FileMetaData*, int>> files_meta;
793
    std::vector<Status> statuses;
794
    for (int level = 0; level < num_levels_; level++) {
795
796
      for (auto& file_meta_pair : levels_[level].added_files) {
        auto* file_meta = file_meta_pair.second;
797
798
799
800
801
        // If the file has been opened before, just skip it.
        if (!file_meta->table_reader_handle) {
          files_meta.emplace_back(file_meta, level);
          statuses.emplace_back(Status::OK());
        }
802
803
804
805
806
807
        if (files_meta.size() >= max_load) {
          break;
        }
      }
      if (files_meta.size() >= max_load) {
        break;
808
809
810
811
      }
    }

    std::atomic<size_t> next_file_meta_idx(0);
812
    std::function<void()> load_handlers_func([&]() {
813
814
815
816
817
818
      while (true) {
        size_t file_idx = next_file_meta_idx.fetch_add(1);
        if (file_idx >= files_meta.size()) {
          break;
        }

819
820
        auto* file_meta = files_meta[file_idx].first;
        int level = files_meta[file_idx].second;
821
        statuses[file_idx] = table_cache_->FindTable(
822
            file_options_, *(base_vstorage_->InternalComparator()),
823
824
825
826
            file_meta->fd, &file_meta->table_reader_handle, prefix_extractor,
            false /*no_io */, true /* record_read_stats */,
            internal_stats->GetFileReadHist(level), false, level,
            prefetch_index_and_filter_in_cache);
827
828
829
830
        if (file_meta->table_reader_handle != nullptr) {
          // Load table_reader
          file_meta->fd.table_reader = table_cache_->GetTableReaderFromHandle(
              file_meta->table_reader_handle);
831
        }
832
      }
833
    });
834

835
836
837
838
839
840
841
    std::vector<port::Thread> threads;
    for (int i = 1; i < max_threads; i++) {
      threads.emplace_back(load_handlers_func);
    }
    load_handlers_func();
    for (auto& t : threads) {
      t.join();
842
    }
843
844
845
846
847
848
    for (const auto& s : statuses) {
      if (!s.ok()) {
        return s;
      }
    }
    return Status::OK();
849
850
851
852
  }

  void MaybeAddFile(VersionStorageInfo* vstorage, int level, FileMetaData* f) {
    if (levels_[level].deleted_files.count(f->fd.GetNumber()) > 0) {
853
      // f is to-be-deleted table file
854
      vstorage->RemoveCurrentStats(f);
855
    } else {
856
857
      assert(ioptions_);
      vstorage->AddFile(level, f, ioptions_->info_log);
858
859
860
861
    }
  }
};

862
VersionBuilder::VersionBuilder(const FileOptions& file_options,
863
                               const ImmutableCFOptions* ioptions,
864
                               TableCache* table_cache,
865
                               VersionStorageInfo* base_vstorage,
866
867
868
                               VersionSet* version_set)
    : rep_(new Rep(file_options, ioptions, table_cache, base_vstorage,
                   version_set)) {}
869

870
VersionBuilder::~VersionBuilder() = default;
871
872
873
874
875

bool VersionBuilder::CheckConsistencyForNumLevels() {
  return rep_->CheckConsistencyForNumLevels();
}

876
Status VersionBuilder::Apply(VersionEdit* edit) { return rep_->Apply(edit); }
877

878
879
Status VersionBuilder::SaveTo(VersionStorageInfo* vstorage) {
  return rep_->SaveTo(vstorage);
880
}
881

882
883
884
885
886
887
888
Status VersionBuilder::LoadTableHandlers(
    InternalStats* internal_stats, int max_threads,
    bool prefetch_index_and_filter_in_cache, bool is_initial_load,
    const SliceTransform* prefix_extractor) {
  return rep_->LoadTableHandlers(internal_stats, max_threads,
                                 prefetch_index_and_filter_in_cache,
                                 is_initial_load, prefix_extractor);
889
}
890

891
892
893
BaseReferencedVersionBuilder::BaseReferencedVersionBuilder(
    ColumnFamilyData* cfd)
    : version_builder_(new VersionBuilder(
894
895
896
          cfd->current()->version_set()->file_options(), cfd->ioptions(),
          cfd->table_cache(), cfd->current()->storage_info(),
          cfd->current()->version_set())),
897
898
899
900
901
902
903
      version_(cfd->current()) {
  version_->Ref();
}

BaseReferencedVersionBuilder::BaseReferencedVersionBuilder(
    ColumnFamilyData* cfd, Version* v)
    : version_builder_(new VersionBuilder(
904
905
          cfd->current()->version_set()->file_options(), cfd->ioptions(),
          cfd->table_cache(), v->storage_info(), v->version_set())),
906
907
908
909
910
911
912
913
      version_(v) {
  assert(version_ != cfd->current());
}

BaseReferencedVersionBuilder::~BaseReferencedVersionBuilder() {
  version_->Unref();
}

914
}  // namespace ROCKSDB_NAMESPACE