db_impl.cc 115.9 KB
Newer Older
1
//  Copyright (c) 2011-present, Facebook, Inc.  All rights reserved.
S
Siying Dong 已提交
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
//
J
jorlow@chromium.org 已提交
6 7 8 9 10
// 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/db_impl.h"

L
liuhuahang 已提交
11
#ifndef __STDC_FORMAT_MACROS
I
Igor Canadi 已提交
12
#define __STDC_FORMAT_MACROS
L
liuhuahang 已提交
13
#endif
14
#include <stdint.h>
D
David Bernard 已提交
15
#ifdef OS_SOLARIS
D
David Bernard 已提交
16
#include <alloca.h>
D
David Bernard 已提交
17
#endif
18

J
jorlow@chromium.org 已提交
19
#include <algorithm>
20
#include <cstdio>
21
#include <map>
J
jorlow@chromium.org 已提交
22
#include <set>
23
#include <stdexcept>
24
#include <string>
25
#include <unordered_map>
26
#include <unordered_set>
T
Tomislav Novak 已提交
27
#include <utility>
28
#include <vector>
29

J
jorlow@chromium.org 已提交
30
#include "db/builder.h"
I
Igor Canadi 已提交
31
#include "db/compaction_job.h"
32
#include "db/db_info_dumper.h"
33
#include "db/db_iter.h"
K
kailiu 已提交
34
#include "db/dbformat.h"
35
#include "db/error_handler.h"
36
#include "db/event_helpers.h"
37
#include "db/external_sst_file_ingestion_job.h"
38 39
#include "db/flush_job.h"
#include "db/forward_iterator.h"
I
Igor Canadi 已提交
40
#include "db/job_context.h"
J
jorlow@chromium.org 已提交
41 42
#include "db/log_reader.h"
#include "db/log_writer.h"
43
#include "db/malloc_stats.h"
44
#include "db/map_builder.h"
J
jorlow@chromium.org 已提交
45
#include "db/memtable.h"
K
kailiu 已提交
46
#include "db/memtable_list.h"
47
#include "db/merge_context.h"
48
#include "db/merge_helper.h"
A
Andrew Kryczka 已提交
49
#include "db/range_del_aggregator.h"
J
jorlow@chromium.org 已提交
50
#include "db/table_cache.h"
K
kailiu 已提交
51
#include "db/table_properties_collector.h"
52
#include "db/transaction_log_impl.h"
J
jorlow@chromium.org 已提交
53 54
#include "db/version_set.h"
#include "db/write_batch_internal.h"
A
agiardullo 已提交
55
#include "db/write_callback.h"
56 57
#include "memtable/hash_linklist_rep.h"
#include "memtable/hash_skiplist_rep.h"
58 59 60 61 62 63 64
#include "monitoring/iostats_context_imp.h"
#include "monitoring/perf_context_imp.h"
#include "monitoring/thread_status_updater.h"
#include "monitoring/thread_status_util.h"
#include "options/cf_options.h"
#include "options/options_helper.h"
#include "options/options_parser.h"
65
#include "port/port.h"
I
Igor Canadi 已提交
66
#include "rocksdb/cache.h"
67
#include "rocksdb/compaction_filter.h"
A
Aaron G 已提交
68
#include "rocksdb/convenience.h"
69 70 71 72 73
#include "rocksdb/db.h"
#include "rocksdb/env.h"
#include "rocksdb/merge_operator.h"
#include "rocksdb/statistics.h"
#include "rocksdb/status.h"
S
Siying Dong 已提交
74
#include "rocksdb/table.h"
75
#include "rocksdb/write_buffer_manager.h"
J
jorlow@chromium.org 已提交
76
#include "table/block.h"
77
#include "table/block_based_table_factory.h"
78
#include "table/merging_iterator.h"
K
kailiu 已提交
79
#include "table/table_builder.h"
J
jorlow@chromium.org 已提交
80
#include "table/two_level_iterator.h"
A
Aaron G 已提交
81
#include "tools/sst_dump_tool_imp.h"
82
#include "util/auto_roll_logger.h"
K
kailiu 已提交
83
#include "util/autovector.h"
84
#include "util/build_version.h"
J
jorlow@chromium.org 已提交
85
#include "util/coding.h"
I
Igor Canadi 已提交
86
#include "util/compression.h"
87
#include "util/crc32c.h"
88
#include "util/file_reader_writer.h"
89
#include "util/file_util.h"
90
#include "util/filename.h"
H
Haobo Xu 已提交
91
#include "util/log_buffer.h"
92
#include "util/logging.h"
J
jorlow@chromium.org 已提交
93
#include "util/mutexlock.h"
94
#include "util/sst_file_manager_impl.h"
95
#include "util/stop_watch.h"
96
#include "util/string_util.h"
97
#include "util/sync_point.h"
J
jorlow@chromium.org 已提交
98

99
namespace rocksdb {
100
const std::string kDefaultColumnFamilyName("default");
101
void DumpRocksDBBuildVersion(Logger* log);
102

A
Aaron Gao 已提交
103 104 105
CompressionType GetCompressionFlush(
    const ImmutableCFOptions& ioptions,
    const MutableCFOptions& mutable_cf_options) {
106 107 108
  // Compressing memtable flushes might not help unless the sequential load
  // optimization is used for leveled compaction. Otherwise the CPU and
  // latency overhead is not offset by saving much space.
109
  if (ioptions.compaction_style == kCompactionStyleUniversal) {
110 111
    if (mutable_cf_options.compaction_options_universal
            .compression_size_percent < 0) {
112 113 114 115 116 117 118
      return mutable_cf_options.compression;
    } else {
      return kNoCompression;
    }
  } else if (!ioptions.compression_per_level.empty()) {
    // For leveled compress when min_level_to_compress != 0.
    return ioptions.compression_per_level[0];
119
  } else {
A
Aaron Gao 已提交
120
    return mutable_cf_options.compression;
121 122
  }
}
I
Igor Canadi 已提交
123

S
Siying Dong 已提交
124
namespace {
125
void DumpSupportInfo(Logger* logger) {
126
  ROCKS_LOG_HEADER(logger, "Compression algorithms supported:");
127 128 129 130 131 132 133
  for (auto& compression : OptionsHelper::compression_type_string_map) {
    if (compression.second != kNoCompression &&
        compression.second != kDisableCompressionOption) {
      ROCKS_LOG_HEADER(logger, "\t%s supported: %d", compression.first.c_str(),
                       CompressionTypeSupported(compression.second));
    }
  }
134 135
  ROCKS_LOG_HEADER(logger, "Fast CRC32 supported: %s",
                   crc32c::IsFastCrc32Supported().c_str());
I
Igor Canadi 已提交
136
}
137 138

int64_t kDefaultLowPriThrottledRate = 2 * 1024 * 1024;
139
}  // namespace
140

141
DBImpl::DBImpl(const DBOptions& options, const std::string& dbname,
142
               const bool seq_per_batch, const bool batch_per_txn)
J
jorlow@chromium.org 已提交
143
    : env_(options.env),
H
heyongqiang 已提交
144
      dbname_(dbname),
145
      own_info_log_(options.info_log == nullptr),
146 147 148
      initial_db_options_(SanitizeOptions(dbname, options)),
      immutable_db_options_(initial_db_options_),
      mutable_db_options_(initial_db_options_),
149
      stats_(immutable_db_options_.statistics.get()),
150
      db_lock_(nullptr),
151 152
      mutex_(stats_, env_, DB_MUTEX_WAIT_MICROS,
             immutable_db_options_.use_adaptive_mutex),
I
Igor Canadi 已提交
153
      shutting_down_(false),
J
jorlow@chromium.org 已提交
154
      bg_cv_(&mutex_),
155
      logfile_number_(0),
156
      log_dir_synced_(false),
I
Igor Canadi 已提交
157
      log_empty_(true),
158
      default_cf_handle_(nullptr),
159
      log_sync_cv_(&mutex_),
I
Igor Canadi 已提交
160 161
      total_log_size_(0),
      max_total_in_memory_state_(0),
S
sdong 已提交
162
      is_snapshot_supported_(true),
163
      write_buffer_manager_(immutable_db_options_.write_buffer_manager.get()),
164
      write_thread_(immutable_db_options_),
165
      nonmem_write_thread_(immutable_db_options_),
166
      write_controller_(mutable_db_options_.delayed_write_rate),
167 168 169 170 171
      // Use delayed_write_rate as a base line to determine the initial
      // low pri write rate limit. It may be adjusted later.
      low_pri_write_rate_limiter_(NewGenericRateLimiter(std::min(
          static_cast<int64_t>(mutable_db_options_.delayed_write_rate / 8),
          kDefaultLowPriThrottledRate))),
S
sdong 已提交
172
      last_batch_group_size_(0),
173 174
      unscheduled_flushes_(0),
      unscheduled_compactions_(0),
175
      bg_bottom_compaction_scheduled_(0),
176
      bg_compaction_scheduled_(0),
177
      num_running_compactions_(0),
178
      bg_flush_scheduled_(0),
179
      num_running_flushes_(0),
180
      bg_purge_scheduled_(0),
181
      disable_delete_obsolete_files_(0),
182
      pending_purge_obsolete_files_(0),
183
      delete_obsolete_files_last_run_(env_->NowMicros()),
184
      last_stats_dump_time_microsec_(0),
185
      next_job_id_(1),
186
      has_unpersisted_data_(false),
S
Siying Dong 已提交
187
      unable_to_release_oldest_log_(false),
188
      env_options_(BuildDBOptions(immutable_db_options_, mutable_db_options_)),
189
      env_options_for_compaction_(env_->OptimizeForCompactionTableWrite(
190
          env_options_, immutable_db_options_)),
191
      num_running_ingest_file_(0),
I
Igor Canadi 已提交
192
#ifndef ROCKSDB_LITE
193
      wal_manager_(immutable_db_options_, env_options_, seq_per_batch),
I
Igor Canadi 已提交
194
#endif  // ROCKSDB_LITE
195
      event_logger_(immutable_db_options_.info_log.get()),
196
      bg_work_paused_(0),
197
      bg_compaction_paused_(0),
198
      refitting_level_(false),
199
      opened_successfully_(false),
200
      two_write_queues_(options.two_write_queues),
201
      manual_wal_flush_(options.manual_wal_flush),
202
      seq_per_batch_(seq_per_batch),
203
      batch_per_txn_(batch_per_txn),
204 205 206 207 208 209 210 211 212 213 214
      // last_sequencee_ is always maintained by the main queue that also writes
      // to the memtable. When two_write_queues_ is disabled last seq in
      // memtable is the same as last seq published to the readers. When it is
      // enabled but seq_per_batch_ is disabled, last seq in memtable still
      // indicates last published seq since wal-only writes that go to the 2nd
      // queue do not consume a sequence number. Otherwise writes performed by
      // the 2nd queue could change what is visible to the readers. In this
      // cases, last_seq_same_as_publish_seq_==false, the 2nd queue maintains a
      // separate variable to indicate the last published sequence.
      last_seq_same_as_publish_seq_(
          !(seq_per_batch && options.two_write_queues)),
215 216 217 218
      // Since seq_per_batch_ is currently set only by WritePreparedTxn which
      // requires a custom gc for compaction, we use that to set use_custom_gc_
      // as well.
      use_custom_gc_(seq_per_batch),
219
      shutdown_initiated_(false),
220
      own_sfm_(options.sst_file_manager == nullptr),
221
      preserve_deletes_(options.preserve_deletes),
222
      closed_(false),
223
      error_handler_(this, immutable_db_options_, &mutex_) {
224 225 226
  // !batch_per_trx_ implies seq_per_batch_ because it is only unset for
  // WriteUnprepared, which should use seq_per_batch_.
  assert(batch_per_txn_ || seq_per_batch_);
H
heyongqiang 已提交
227
  env_->GetAbsolutePath(dbname, &db_absolute_path_);
228

J
jorlow@chromium.org 已提交
229
  // Reserve ten files or so for other uses and give the rest to TableCache.
230
  // Give a large number for setting of "infinite" open files.
L
Leonidas Galanis 已提交
231
  const int table_cache_size = (mutable_db_options_.max_open_files == -1)
232
                                   ? TableCache::kInfiniteCapacity
L
Leonidas Galanis 已提交
233
                                   : mutable_db_options_.max_open_files - 10;
234 235
  table_cache_ = NewLRUCache(table_cache_size,
                             immutable_db_options_.table_cache_numshardbits);
236

237
  versions_.reset(new VersionSet(dbname_, &immutable_db_options_, env_options_,
238
                                 table_cache_.get(), write_buffer_manager_,
239
                                 &write_controller_));
240 241
  column_family_memtables_.reset(
      new ColumnFamilyMemTablesImpl(versions_->GetColumnFamilySet()));
242

243 244 245 246 247
  DumpRocksDBBuildVersion(immutable_db_options_.info_log.get());
  DumpDBFileSummary(immutable_db_options_, dbname_);
  immutable_db_options_.Dump(immutable_db_options_.info_log.get());
  mutable_db_options_.Dump(immutable_db_options_.info_log.get());
  DumpSupportInfo(immutable_db_options_.info_log.get());
248 249 250 251 252

  // always open the DB with 0 here, which means if preserve_deletes_==true
  // we won't drop any deletion markers until SetPreserveDeletesSequenceNumber()
  // is called by client and this seqnum is advanced.
  preserve_deletes_seqnum_.store(0);
J
jorlow@chromium.org 已提交
253 254
}

255 256 257 258 259 260 261 262 263 264
Status DBImpl::Resume() {
  ROCKS_LOG_INFO(immutable_db_options_.info_log, "Resuming DB");

  InstrumentedMutexLock db_mutex(&mutex_);

  if (!error_handler_.IsDBStopped() && !error_handler_.IsBGWorkStopped()) {
    // Nothing to do
    return Status::OK();
  }

265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299
  if (error_handler_.IsRecoveryInProgress()) {
    // Don't allow a mix of manual and automatic recovery
    return Status::Busy();
  }

  mutex_.Unlock();
  Status s = error_handler_.RecoverFromBGError(true);
  mutex_.Lock();
  return s;
}

// This function implements the guts of recovery from a background error. It
// is eventually called for both manual as well as automatic recovery. It does
// the following -
// 1. Wait for currently scheduled background flush/compaction to exit, in
//    order to inadvertently causing an error and thinking recovery failed
// 2. Flush memtables if there's any data for all the CFs. This may result
//    another error, which will be saved by error_handler_ and reported later
//    as the recovery status
// 3. Find and delete any obsolete files
// 4. Schedule compactions if needed for all the CFs. This is needed as the
//    flush in the prior step might have been a no-op for some CFs, which
//    means a new super version wouldn't have been installed
Status DBImpl::ResumeImpl() {
  mutex_.AssertHeld();
  WaitForBackgroundWork();

  Status bg_error = error_handler_.GetBGError();
  Status s;
  if (shutdown_initiated_) {
    // Returning shutdown status to SFM during auto recovery will cause it
    // to abort the recovery and allow the shutdown to progress
    s = Status::ShutdownInProgress();
  }
  if (s.ok() && bg_error.severity() > Status::Severity::kHardError) {
奏之章 已提交
300 301
    ROCKS_LOG_INFO(
        immutable_db_options_.info_log,
302
        "DB resume requested but failed due to Fatal/Unrecoverable error");
303 304 305 306 307 308 309 310 311 312 313 314
    s = bg_error;
  }

  // We cannot guarantee consistency of the WAL. So force flush Memtables of
  // all the column families
  if (s.ok()) {
    s = FlushAllCFs(FlushReason::kErrorRecovery);
    if (!s.ok()) {
      ROCKS_LOG_INFO(immutable_db_options_.info_log,
                     "DB resume requested but failed due to Flush failure [%s]",
                     s.ToString().c_str());
    }
315 316 317 318
  }

  JobContext job_context(0);
  FindObsoleteFiles(&job_context, true);
319 320 321
  if (s.ok()) {
    s = error_handler_.ClearBGError();
  }
322 323 324 325 326 327 328 329
  mutex_.Unlock();

  job_context.manifest_file_number = 1;
  if (job_context.HaveSomethingToDelete()) {
    PurgeObsoleteFiles(job_context);
  }
  job_context.Clean();

330 331 332
  if (s.ok()) {
    ROCKS_LOG_INFO(immutable_db_options_.info_log, "Successfully resumed DB");
  }
333
  mutex_.Lock();
334 335 336 337 338 339 340 341 342 343 344 345 346 347
  // Check for shutdown again before scheduling further compactions,
  // since we released and re-acquired the lock above
  if (shutdown_initiated_) {
    s = Status::ShutdownInProgress();
  }
  if (s.ok()) {
    for (auto cfd : *versions_->GetColumnFamilySet()) {
      SchedulePendingCompaction(cfd);
    }
    MaybeScheduleFlushOrCompaction();
  }

  // Wake up any waiters - in this case, it could be the shutdown thread
  bg_cv_.SignalAll();
348 349 350

  // No need to check BGError again. If something happened, event listener would be
  // notified and the operation causing it would have failed
351 352 353 354 355 356 357 358 359
  return s;
}

void DBImpl::WaitForBackgroundWork() {
  // Wait for background work to finish
  while (bg_bottom_compaction_scheduled_ || bg_compaction_scheduled_ ||
         bg_flush_scheduled_) {
    bg_cv_.Wait();
  }
360 361
}

362
// Will lock the mutex_,  will wait for completion if wait is true
363
void DBImpl::CancelAllBackgroundWork(bool wait) {
364
  InstrumentedMutexLock l(&mutex_);
365

366 367
  ROCKS_LOG_INFO(immutable_db_options_.info_log,
                 "Shutdown: canceling all background work");
368

369
  if (!shutting_down_.load(std::memory_order_acquire) &&
370
      has_unpersisted_data_.load(std::memory_order_relaxed) &&
Y
Yi Wu 已提交
371
      !mutable_db_options_.avoid_flush_during_shutdown) {
372
    for (auto cfd : *versions_->GetColumnFamilySet()) {
373
      if (!cfd->IsDropped() && cfd->initialized() && !cfd->mem()->IsEmpty()) {
374 375
        cfd->Ref();
        mutex_.Unlock();
376
        FlushMemTable(cfd, FlushOptions(), FlushReason::kShutDown);
377
        mutex_.Lock();
378
        cfd->Unref();
379 380
      }
    }
381
    versions_->GetColumnFamilySet()->FreeDeadColumnFamilies();
382
  }
383 384 385 386 387 388

  shutting_down_.store(true, std::memory_order_release);
  bg_cv_.SignalAll();
  if (!wait) {
    return;
  }
389
  WaitForBackgroundWork();
390 391
}

392
Status DBImpl::CloseHelper() {
393 394 395 396 397 398 399 400 401 402
  // Guarantee that there is no background error recovery in progress before
  // continuing with the shutdown
  mutex_.Lock();
  shutdown_initiated_ = true;
  error_handler_.CancelErrorRecovery();
  while (error_handler_.IsRecoveryInProgress()) {
    bg_cv_.Wait();
  }
  mutex_.Unlock();

403 404
  // CancelAllBackgroundWork called with false means we just set the shutdown
  // marker. After this we do a variant of the waiting and unschedule work
405 406
  // (to consider: moving all the waiting into CancelAllBackgroundWork(true))
  CancelAllBackgroundWork(false);
407 408
  int bottom_compactions_unscheduled =
      env_->UnSchedule(this, Env::Priority::BOTTOM);
409 410
  int compactions_unscheduled = env_->UnSchedule(this, Env::Priority::LOW);
  int flushes_unscheduled = env_->UnSchedule(this, Env::Priority::HIGH);
411
  Status ret;
412
  mutex_.Lock();
413
  bg_bottom_compaction_scheduled_ -= bottom_compactions_unscheduled;
414 415 416 417
  bg_compaction_scheduled_ -= compactions_unscheduled;
  bg_flush_scheduled_ -= flushes_unscheduled;

  // Wait for background work to finish
418
  while (bg_bottom_compaction_scheduled_ || bg_compaction_scheduled_ ||
419
         bg_flush_scheduled_ || bg_purge_scheduled_ ||
420 421
         pending_purge_obsolete_files_ ||
         error_handler_.IsRecoveryInProgress()) {
422
    TEST_SYNC_POINT("DBImpl::~DBImpl:WaitJob");
423 424
    bg_cv_.Wait();
  }
425
  TEST_SYNC_POINT_CALLBACK("DBImpl::CloseHelper:PendingPurgeFinished",
426
                           &files_grabbed_for_purge_);
427
  EraseThreadStatusDbInfo();
I
Igor Canadi 已提交
428 429
  flush_scheduler_.Clear();

430
  while (!flush_queue_.empty()) {
431 432 433 434 435 436
    const FlushRequest& flush_req = PopFirstFromFlushQueue();
    for (const auto& iter : flush_req) {
      ColumnFamilyData* cfd = iter.first;
      if (cfd->Unref()) {
        delete cfd;
      }
437 438 439 440 441 442 443 444 445
    }
  }
  while (!compaction_queue_.empty()) {
    auto cfd = PopFirstFromCompactionQueue();
    if (cfd->Unref()) {
      delete cfd;
    }
  }

I
Igor Canadi 已提交
446 447 448 449 450
  if (default_cf_handle_ != nullptr) {
    // we need to delete handle outside of lock because it does its own locking
    mutex_.Unlock();
    delete default_cf_handle_;
    mutex_.Lock();
451 452
  }

I
Igor Canadi 已提交
453 454 455 456 457 458 459 460 461 462
  // Clean up obsolete files due to SuperVersion release.
  // (1) Need to delete to obsolete files before closing because RepairDB()
  // scans all existing files in the file system and builds manifest file.
  // Keeping obsolete files confuses the repair process.
  // (2) Need to check if we Open()/Recover() the DB successfully before
  // deleting because if VersionSet recover fails (may be due to corrupted
  // manifest file), it is not able to identify live files correctly. As a
  // result, all "live" files can get deleted by accident. However, corrupted
  // manifest is recoverable by RepairDB().
  if (opened_successfully_) {
463
    JobContext job_context(next_job_id_.fetch_add(1));
I
Igor Canadi 已提交
464
    FindObsoleteFiles(&job_context, true);
465 466

    mutex_.Unlock();
I
Igor Canadi 已提交
467
    // manifest number starting from 2
I
Igor Canadi 已提交
468 469 470
    job_context.manifest_file_number = 1;
    if (job_context.HaveSomethingToDelete()) {
      PurgeObsoleteFiles(job_context);
471
    }
I
Igor Canadi 已提交
472
    job_context.Clean();
473
    mutex_.Lock();
474 475
  }

476 477 478
  for (auto l : logs_to_free_) {
    delete l;
  }
479
  for (auto& log : logs_) {
480 481 482
    uint64_t log_number = log.writer->get_log_number();
    Status s = log.ClearWriter();
    if (!s.ok()) {
483 484 485 486 487
      ROCKS_LOG_WARN(
          immutable_db_options_.info_log,
          "Unable to Sync WAL file %s with error -- %s",
          LogFileName(immutable_db_options_.wal_dir, log_number).c_str(),
          s.ToString().c_str());
488 489 490 491 492
      // Retain the first error
      if (ret.ok()) {
        ret = s;
      }
    }
493
  }
494
  logs_.clear();
495

496 497 498 499 500 501 502 503 504 505 506 507 508 509 510
  // Table cache may have table handles holding blocks from the block cache.
  // We need to release them before the block cache is destroyed. The block
  // cache may be destroyed inside versions_.reset(), when column family data
  // list is destroyed, so leaving handles in table cache after
  // versions_.reset() may cause issues.
  // Here we clean all unreferenced handles in table cache.
  // Now we assume all user queries have finished, so only version set itself
  // can possibly hold the blocks from block cache. After releasing unreferenced
  // handles here, only handles held by version set left and inside
  // versions_.reset(), we will release them. There, we need to make sure every
  // time a handle is released, we erase it from the cache too. By doing that,
  // we can guarantee that after versions_.reset(), table cache is empty
  // so the cache can be safely destroyed.
  table_cache_->EraseUnRefEntries();

I
Islam AbdelRahman 已提交
511 512 513 514
  for (auto& txn_entry : recovered_transactions_) {
    delete txn_entry.second;
  }

515
  // versions need to be destroyed before table_cache since it can hold
516 517
  // references to table_cache.
  versions_.reset();
518
  mutex_.Unlock();
I
Igor Canadi 已提交
519 520 521
  if (db_lock_ != nullptr) {
    env_->UnlockFile(db_lock_);
  }
522

523
  ROCKS_LOG_INFO(immutable_db_options_.info_log, "Shutdown complete");
524
  LogFlush(immutable_db_options_.info_log);
525

526 527 528 529 530 531 532 533 534 535 536
#ifndef ROCKSDB_LITE
  // If the sst_file_manager was allocated by us during DB::Open(), ccall
  // Close() on it before closing the info_log. Otherwise, background thread
  // in SstFileManagerImpl might try to log something
  if (immutable_db_options_.sst_file_manager && own_sfm_) {
    auto sfm = static_cast<SstFileManagerImpl*>(
        immutable_db_options_.sst_file_manager.get());
    sfm->Close();
  }
#endif // ROCKSDB_LITE

537
  if (immutable_db_options_.info_log && own_info_log_) {
538 539 540 541
    Status s = immutable_db_options_.info_log->Close();
    if (ret.ok()) {
      ret = s;
    }
542
  }
543
  return ret;
J
jorlow@chromium.org 已提交
544 545
}

546
Status DBImpl::CloseImpl() { return CloseHelper(); }
547 548 549 550 551 552 553

DBImpl::~DBImpl() {
  if (!closed_) {
    closed_ = true;
    CloseHelper();
  }
}
554

J
jorlow@chromium.org 已提交
555
void DBImpl::MaybeIgnoreError(Status* s) const {
556
  if (s->ok() || immutable_db_options_.paranoid_checks) {
J
jorlow@chromium.org 已提交
557 558
    // No change needed
  } else {
559 560
    ROCKS_LOG_WARN(immutable_db_options_.info_log, "Ignoring error %s",
                   s->ToString().c_str());
J
jorlow@chromium.org 已提交
561 562 563 564
    *s = Status::OK();
  }
}

565
const Status DBImpl::CreateArchivalDirectory() {
566 567 568
  if (immutable_db_options_.wal_ttl_seconds > 0 ||
      immutable_db_options_.wal_size_limit_mb > 0) {
    std::string archivalPath = ArchivalDirectory(immutable_db_options_.wal_dir);
569 570 571 572 573
    return env_->CreateDirIfMissing(archivalPath);
  }
  return Status::OK();
}

574
void DBImpl::PrintStatistics() {
575
  auto dbstats = immutable_db_options_.statistics.get();
576
  if (dbstats) {
577 578
    ROCKS_LOG_WARN(immutable_db_options_.info_log, "STATISTICS:\n %s",
                   dbstats->ToString().c_str());
579 580 581
  }
}

582
void DBImpl::MaybeDumpStats() {
583 584 585 586 587
  mutex_.Lock();
  unsigned int stats_dump_period_sec =
      mutable_db_options_.stats_dump_period_sec;
  mutex_.Unlock();
  if (stats_dump_period_sec == 0) return;
H
Haobo Xu 已提交
588 589 590

  const uint64_t now_micros = env_->NowMicros();

591
  if (last_stats_dump_time_microsec_ + stats_dump_period_sec * 1000000 <=
592
      now_micros) {
H
Haobo Xu 已提交
593 594 595 596 597
    // Multiple threads could race in here simultaneously.
    // However, the last one will update last_stats_dump_time_microsec_
    // atomically. We could see more than one dump during one dump
    // period in rare cases.
    last_stats_dump_time_microsec_ = now_micros;
598

599
#ifndef ROCKSDB_LITE
600 601 602 603 604 605 606
    const DBPropertyInfo* cf_property_info =
        GetPropertyInfo(DB::Properties::kCFStats);
    assert(cf_property_info != nullptr);
    const DBPropertyInfo* db_property_info =
        GetPropertyInfo(DB::Properties::kDBStats);
    assert(db_property_info != nullptr);

H
Haobo Xu 已提交
607
    std::string stats;
608
    {
609
      InstrumentedMutexLock l(&mutex_);
610 611
      default_cf_internal_stats_->GetStringProperty(
          *db_property_info, DB::Properties::kDBStats, &stats);
612
      for (auto cfd : *versions_->GetColumnFamilySet()) {
613 614 615 616 617
        if (cfd->initialized()) {
          cfd->internal_stats()->GetStringProperty(
              *cf_property_info, DB::Properties::kCFStatsNoFileHistogram,
              &stats);
        }
618 619
      }
      for (auto cfd : *versions_->GetColumnFamilySet()) {
620 621 622 623
        if (cfd->initialized()) {
          cfd->internal_stats()->GetStringProperty(
              *cf_property_info, DB::Properties::kCFFileHistogram, &stats);
        }
624 625
      }
    }
626 627 628
    ROCKS_LOG_WARN(immutable_db_options_.info_log,
                   "------- DUMPING STATS -------");
    ROCKS_LOG_WARN(immutable_db_options_.info_log, "%s", stats.c_str());
S
Siying Dong 已提交
629 630 631 632
    if (immutable_db_options_.dump_malloc_stats) {
      stats.clear();
      DumpMallocStats(&stats);
      if (!stats.empty()) {
633 634 635
        ROCKS_LOG_WARN(immutable_db_options_.info_log,
                       "------- Malloc STATS -------");
        ROCKS_LOG_WARN(immutable_db_options_.info_log, "%s", stats.c_str());
S
Siying Dong 已提交
636 637
      }
    }
638
#endif  // !ROCKSDB_LITE
639

640
    PrintStatistics();
641 642 643
  }
}

644 645 646 647 648 649 650 651 652 653
void DBImpl::ScheduleBgLogWriterClose(JobContext* job_context) {
  if (!job_context->logs_to_free.empty()) {
    for (auto l : job_context->logs_to_free) {
      AddToLogsToFreeQueue(l);
    }
    job_context->logs_to_free.clear();
    SchedulePurge();
  }
}

654 655 656 657 658 659 660 661 662 663
Directory* DBImpl::GetDataDir(ColumnFamilyData* cfd, size_t path_id) const {
  assert(cfd);
  Directory* ret_dir = cfd->GetDataDir(path_id);
  if (ret_dir == nullptr) {
    return directories_.GetDataDir(path_id);
  }
  return ret_dir;
}

Directory* DBImpl::Directories::GetDataDir(size_t path_id) const {
S
Siying Dong 已提交
664 665 666 667 668
  assert(path_id < data_dirs_.size());
  Directory* ret_dir = data_dirs_[path_id].get();
  if (ret_dir == nullptr) {
    // Should use db_dir_
    return db_dir_.get();
669
  }
S
Siying Dong 已提交
670
  return ret_dir;
671 672
}

673 674
Status DBImpl::SetOptions(
    ColumnFamilyHandle* column_family,
S
Siying Dong 已提交
675 676
    const std::unordered_map<std::string, std::string>& options_map) {
#ifdef ROCKSDB_LITE
677 678
  (void)column_family;
  (void)options_map;
S
Siying Dong 已提交
679 680 681 682 683 684 685 686
  return Status::NotSupported("Not supported in ROCKSDB LITE");
#else
  auto* cfd = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family)->cfd();
  if (options_map.empty()) {
    ROCKS_LOG_WARN(immutable_db_options_.info_log,
                   "SetOptions() on column family [%s], empty input",
                   cfd->GetName().c_str());
    return Status::InvalidArgument("empty input");
687 688
  }

S
Siying Dong 已提交
689 690 691 692
  MutableCFOptions new_options;
  Status s;
  Status persist_options_status;
  WriteThread::Writer w;
693
  SuperVersionContext sv_context(/* create_superversion */ true);
S
Siying Dong 已提交
694 695 696 697 698 699 700 701 702 703 704 705
  {
    InstrumentedMutexLock l(&mutex_);
    s = cfd->SetOptions(options_map);
    if (s.ok()) {
      new_options = *cfd->GetLatestMutableCFOptions();
      // Append new version to recompute compaction score.
      VersionEdit dummy_edit;
      versions_->LogAndApply(cfd, new_options, &dummy_edit, &mutex_,
                             directories_.GetDbDir());
      // Trigger possible flush/compactions. This has to be before we persist
      // options to file, otherwise there will be a deadlock with writer
      // thread.
706
      InstallSuperVersionAndScheduleWork(cfd, &sv_context, new_options);
707

Y
Yi Wu 已提交
708 709
      persist_options_status = WriteOptionsFile(
          false /*need_mutex_lock*/, true /*need_enter_write_thread*/);
710
      bg_cv_.SignalAll();
711 712
    }
  }
713
  sv_context.Clean();
714

715 716 717
  ROCKS_LOG_INFO(
      immutable_db_options_.info_log,
      "SetOptions() on column family [%s], inputs:", cfd->GetName().c_str());
S
Siying Dong 已提交
718 719 720
  for (const auto& o : options_map) {
    ROCKS_LOG_INFO(immutable_db_options_.info_log, "%s: %s\n", o.first.c_str(),
                   o.second.c_str());
721
  }
S
Siying Dong 已提交
722 723 724 725 726
  if (s.ok()) {
    ROCKS_LOG_INFO(immutable_db_options_.info_log,
                   "[%s] SetOptions() succeeded", cfd->GetName().c_str());
    new_options.Dump(immutable_db_options_.info_log.get());
    if (!persist_options_status.ok()) {
Y
Yi Wu 已提交
727
      s = persist_options_status;
728
    }
729
  } else {
S
Siying Dong 已提交
730 731
    ROCKS_LOG_WARN(immutable_db_options_.info_log, "[%s] SetOptions() failed",
                   cfd->GetName().c_str());
732
  }
S
Siying Dong 已提交
733 734 735
  LogFlush(immutable_db_options_.info_log);
  return s;
#endif  // ROCKSDB_LITE
736 737
}

S
Siying Dong 已提交
738 739 740
Status DBImpl::SetDBOptions(
    const std::unordered_map<std::string, std::string>& options_map) {
#ifdef ROCKSDB_LITE
741
  (void)options_map;
S
Siying Dong 已提交
742 743 744 745 746 747
  return Status::NotSupported("Not supported in ROCKSDB LITE");
#else
  if (options_map.empty()) {
    ROCKS_LOG_WARN(immutable_db_options_.info_log,
                   "SetDBOptions(), empty input.");
    return Status::InvalidArgument("empty input");
I
Igor Canadi 已提交
748
  }
749

S
Siying Dong 已提交
750 751 752
  MutableDBOptions new_options;
  Status s;
  Status persist_options_status;
753
  bool wal_changed = false;
S
Siying Dong 已提交
754 755 756 757 758 759 760 761 762 763 764 765 766
  WriteThread::Writer w;
  WriteContext write_context;
  {
    InstrumentedMutexLock l(&mutex_);
    s = GetMutableDBOptionsFromStrings(mutable_db_options_, options_map,
                                       &new_options);
    if (s.ok()) {
      if (new_options.max_background_compactions >
          mutable_db_options_.max_background_compactions) {
        env_->IncBackgroundThreadsIfNeeded(
            new_options.max_background_compactions, Env::Priority::LOW);
        MaybeScheduleFlushOrCompaction();
      }
767

768 769
      write_controller_.set_max_delayed_write_rate(
          new_options.delayed_write_rate);
L
Leonidas Galanis 已提交
770
      table_cache_.get()->SetCapacity(new_options.max_open_files == -1
771
                                          ? TableCache::kInfiniteCapacity
L
Leonidas Galanis 已提交
772
                                          : new_options.max_open_files - 10);
773 774 775 776 777
      wal_changed = mutable_db_options_.wal_bytes_per_sync !=
                    new_options.wal_bytes_per_sync;
      if (new_options.bytes_per_sync == 0) {
        new_options.bytes_per_sync = 1024 * 1024;
      }
778
      mutable_db_options_ = new_options;
779 780
      env_options_for_compaction_ = EnvOptions(
          BuildDBOptions(immutable_db_options_, mutable_db_options_));
781
      env_options_for_compaction_ = env_->OptimizeForCompactionTableWrite(
782
          env_options_for_compaction_, immutable_db_options_);
783
      versions_->ChangeEnvOptions(mutable_db_options_);
784 785 786 787
      env_options_for_compaction_ = env_->OptimizeForCompactionTableRead(
          env_options_for_compaction_, immutable_db_options_);
      env_options_for_compaction_.compaction_readahead_size =
          mutable_db_options_.compaction_readahead_size;
788
      write_thread_.EnterUnbatched(&w, &mutex_);
789 790
      if (total_log_size_ > GetMaxTotalWalSize() || wal_changed) {
        Status purge_wal_status = SwitchWAL(&write_context);
791 792 793 794 795
        if (!purge_wal_status.ok()) {
          ROCKS_LOG_WARN(immutable_db_options_.info_log,
                         "Unable to purge WAL files in SetDBOptions() -- %s",
                         purge_wal_status.ToString().c_str());
        }
796
      }
Y
Yi Wu 已提交
797 798
      persist_options_status = WriteOptionsFile(
          false /*need_mutex_lock*/, false /*need_enter_write_thread*/);
799
      write_thread_.ExitUnbatched(&w);
800 801
    }
  }
802
  ROCKS_LOG_INFO(immutable_db_options_.info_log, "SetDBOptions(), inputs:");
803
  for (const auto& o : options_map) {
804 805
    ROCKS_LOG_INFO(immutable_db_options_.info_log, "%s: %s\n", o.first.c_str(),
                   o.second.c_str());
806 807
  }
  if (s.ok()) {
808
    ROCKS_LOG_INFO(immutable_db_options_.info_log, "SetDBOptions() succeeded");
809 810 811 812 813 814 815
    new_options.Dump(immutable_db_options_.info_log.get());
    if (!persist_options_status.ok()) {
      if (immutable_db_options_.fail_if_options_file_error) {
        s = Status::IOError(
            "SetDBOptions() succeeded, but unable to persist options",
            persist_options_status.ToString());
      }
816 817 818
      ROCKS_LOG_WARN(immutable_db_options_.info_log,
                     "Unable to persist options in SetDBOptions() -- %s",
                     persist_options_status.ToString().c_str());
819 820
    }
  } else {
821
    ROCKS_LOG_WARN(immutable_db_options_.info_log, "SetDBOptions failed");
L
Lei Jin 已提交
822
  }
823
  LogFlush(immutable_db_options_.info_log);
824
  return s;
I
Igor Canadi 已提交
825
#endif  // ROCKSDB_LITE
826 827
}

828
// return the same level if it cannot be moved
A
Andrew Kryczka 已提交
829 830 831
int DBImpl::FindMinimumEmptyLevelFitting(
    ColumnFamilyData* cfd, const MutableCFOptions& /*mutable_cf_options*/,
    int level) {
832
  mutex_.AssertHeld();
S
sdong 已提交
833
  const auto* vstorage = cfd->current()->storage_info();
834
  int minimum_level = level;
835
  for (int i = level - 1; i > 0; --i) {
836
    // stop if level i is not empty
S
sdong 已提交
837
    if (vstorage->NumLevelFiles(i) > 0) break;
838
    // stop if level i is too small (cannot fit the level files)
839
    if (vstorage->MaxBytesForLevel(i) < vstorage->NumLevelBytes(level)) {
840 841
      break;
    }
842 843 844 845 846 847

    minimum_level = i;
  }
  return minimum_level;
}

848
Status DBImpl::FlushWAL(bool sync) {
849
  if (manual_wal_flush_) {
850 851 852 853 854 855 856
    // We need to lock log_write_mutex_ since logs_ might change concurrently
    InstrumentedMutexLock wl(&log_write_mutex_);
    log::Writer* cur_log_writer = logs_.back().writer;
    auto s = cur_log_writer->WriteBuffer();
    if (!s.ok()) {
      ROCKS_LOG_ERROR(immutable_db_options_.info_log, "WAL flush error %s",
                      s.ToString().c_str());
857 858 859 860 861
      // In case there is a fs error we should set it globally to prevent the
      // future writes
      WriteStatusCheck(s);
      // whether sync or not, we should abort the rest of function upon error
      return s;
862 863 864 865 866 867
    }
    if (!sync) {
      ROCKS_LOG_DEBUG(immutable_db_options_.info_log, "FlushWAL sync=false");
      return s;
    }
  }
868 869 870
  if (!sync) {
    return Status::OK();
  }
871 872 873 874 875
  // sync = true
  ROCKS_LOG_DEBUG(immutable_db_options_.info_log, "FlushWAL sync=true");
  return SyncWAL();
}

876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896
Status DBImpl::SyncWAL() {
  autovector<log::Writer*, 1> logs_to_sync;
  bool need_log_dir_sync;
  uint64_t current_log_number;

  {
    InstrumentedMutexLock l(&mutex_);
    assert(!logs_.empty());

    // This SyncWAL() call only cares about logs up to this number.
    current_log_number = logfile_number_;

    while (logs_.front().number <= current_log_number &&
           logs_.front().getting_synced) {
      log_sync_cv_.Wait();
    }
    // First check that logs are safe to sync in background.
    for (auto it = logs_.begin();
         it != logs_.end() && it->number <= current_log_number; ++it) {
      if (!it->writer->file()->writable_file()->IsSyncThreadSafe()) {
        return Status::NotSupported(
897 898 899 900
            "SyncWAL() is not supported for this implementation of WAL file",
            immutable_db_options_.allow_mmap_writes
                ? "try setting Options::allow_mmap_writes to false"
                : Slice());
901 902 903 904 905
      }
    }
    for (auto it = logs_.begin();
         it != logs_.end() && it->number <= current_log_number; ++it) {
      auto& log = *it;
S
Siying Dong 已提交
906 907 908 909
      assert(!log.getting_synced);
      log.getting_synced = true;
      logs_to_sync.push_back(log.writer);
    }
910

S
Siying Dong 已提交
911
    need_log_dir_sync = !log_dir_synced_;
912 913
  }

914
  TEST_SYNC_POINT("DBWALTest::SyncWALNotWaitWrite:1");
S
Siying Dong 已提交
915 916 917 918 919 920 921 922 923 924
  RecordTick(stats_, WAL_FILE_SYNCED);
  Status status;
  for (log::Writer* log : logs_to_sync) {
    status = log->file()->SyncWithoutFlush(immutable_db_options_.use_fsync);
    if (!status.ok()) {
      break;
    }
  }
  if (status.ok() && need_log_dir_sync) {
    status = directories_.GetWalDir()->Fsync();
925
  }
926
  TEST_SYNC_POINT("DBWALTest::SyncWALNotWaitWrite:2");
927

S
Siying Dong 已提交
928 929 930 931 932 933
  TEST_SYNC_POINT("DBImpl::SyncWAL:BeforeMarkLogsSynced:1");
  {
    InstrumentedMutexLock l(&mutex_);
    MarkLogsSynced(current_log_number, need_log_dir_sync, status);
  }
  TEST_SYNC_POINT("DBImpl::SyncWAL:BeforeMarkLogsSynced:2");
934

S
Siying Dong 已提交
935
  return status;
936 937
}

938 939
void DBImpl::MarkLogsSynced(uint64_t up_to, bool synced_dir,
                            const Status& status) {
S
Siying Dong 已提交
940
  mutex_.AssertHeld();
941
  if (synced_dir && logfile_number_ == up_to && status.ok()) {
S
Siying Dong 已提交
942 943 944 945 946 947 948
    log_dir_synced_ = true;
  }
  for (auto it = logs_.begin(); it != logs_.end() && it->number <= up_to;) {
    auto& log = *it;
    assert(log.getting_synced);
    if (status.ok() && logs_.size() > 1) {
      logs_to_free_.push_back(log.ReleaseWriter());
949 950
      // To modify logs_ both mutex_ and log_write_mutex_ must be held
      InstrumentedMutexLock l(&log_write_mutex_);
S
Siying Dong 已提交
951 952 953 954 955 956 957 958 959
      it = logs_.erase(it);
    } else {
      log.getting_synced = false;
      ++it;
    }
  }
  assert(!status.ok() || logs_.empty() || logs_[0].number > up_to ||
         (logs_.size() == 1 && !logs_[0].getting_synced));
  log_sync_cv_.SignalAll();
960 961
}

S
Siying Dong 已提交
962 963
SequenceNumber DBImpl::GetLatestSequenceNumber() const {
  return versions_->LastSequence();
964 965
}

966 967
void DBImpl::SetLastPublishedSequence(SequenceNumber seq) {
  versions_->SetLastPublishedSequence(seq);
968 969
}

970 971 972 973 974 975 976 977 978
bool DBImpl::SetPreserveDeletesSequenceNumber(SequenceNumber seqnum) {
  if (seqnum > preserve_deletes_seqnum_.load()) {
    preserve_deletes_seqnum_.store(seqnum);
    return true;
  } else {
    return false;
  }
}

S
Siying Dong 已提交
979 980 981 982 983 984 985 986 987
InternalIterator* DBImpl::NewInternalIterator(
    Arena* arena, RangeDelAggregator* range_del_agg,
    ColumnFamilyHandle* column_family) {
  ColumnFamilyData* cfd;
  if (column_family == nullptr) {
    cfd = default_cf_handle_->cfd();
  } else {
    auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
    cfd = cfh->cfd();
988
  }
989

S
Siying Dong 已提交
990 991 992 993 994 995 996
  mutex_.Lock();
  SuperVersion* super_version = cfd->GetSuperVersion()->Ref();
  mutex_.Unlock();
  ReadOptions roptions;
  return NewInternalIterator(roptions, cfd, super_version, arena,
                             range_del_agg);
}
I
Islam AbdelRahman 已提交
997

S
Siying Dong 已提交
998 999 1000
void DBImpl::SchedulePurge() {
  mutex_.AssertHeld();
  assert(opened_successfully_);
I
Islam AbdelRahman 已提交
1001

S
Siying Dong 已提交
1002 1003 1004
  // Purge operations are put into High priority queue
  bg_purge_scheduled_++;
  env_->Schedule(&DBImpl::BGWorkPurge, this, Env::Priority::HIGH, nullptr);
I
Islam AbdelRahman 已提交
1005 1006
}

1007 1008 1009
void DBImpl::BackgroundCallPurge() {
  mutex_.Lock();

1010 1011 1012 1013 1014 1015 1016
  // We use one single loop to clear both queues so that after existing the loop
  // both queues are empty. This is stricter than what is needed, but can make
  // it easier for us to reason the correctness.
  while (!purge_queue_.empty() || !logs_to_free_queue_.empty()) {
    if (!purge_queue_.empty()) {
      auto purge_file = purge_queue_.begin();
      auto fname = purge_file->fname;
1017
      auto dir_to_sync = purge_file->dir_to_sync;
1018 1019 1020 1021
      auto type = purge_file->type;
      auto number = purge_file->number;
      auto job_id = purge_file->job_id;
      purge_queue_.pop_front();
1022

1023
      mutex_.Unlock();
1024
      DeleteObsoleteFileImpl(job_id, fname, dir_to_sync, type, number);
1025 1026 1027 1028 1029 1030 1031 1032 1033
      mutex_.Lock();
    } else {
      assert(!logs_to_free_queue_.empty());
      log::Writer* log_writer = *(logs_to_free_queue_.begin());
      logs_to_free_queue_.pop_front();
      mutex_.Unlock();
      delete log_writer;
      mutex_.Lock();
    }
1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044
  }
  bg_purge_scheduled_--;

  bg_cv_.SignalAll();
  // IMPORTANT:there should be no code after calling SignalAll. This call may
  // signal the DB destructor that it's OK to proceed with destruction. In
  // that case, all DB variables will be dealloacated and referencing them
  // will cause trouble.
  mutex_.Unlock();
}

1045 1046
namespace {
struct IterState {
1047
  IterState(DBImpl* _db, InstrumentedMutex* _mu, SuperVersion* _super_version,
1048
            bool _background_purge)
1049 1050 1051
      : db(_db),
        mu(_mu),
        super_version(_super_version),
1052
        background_purge(_background_purge) {}
1053 1054

  DBImpl* db;
1055
  InstrumentedMutex* mu;
1056
  SuperVersion* super_version;
1057
  bool background_purge;
1058 1059
};

A
Andrew Kryczka 已提交
1060
static void CleanupIteratorState(void* arg1, void* /*arg2*/) {
1061
  IterState* state = reinterpret_cast<IterState*>(arg1);
1062

1063
  if (state->super_version->Unref()) {
1064 1065 1066
    // Job id == 0 means that this is not our background process, but rather
    // user thread
    JobContext job_context(0);
1067

1068 1069
    state->mu->Lock();
    state->super_version->Cleanup();
I
Igor Canadi 已提交
1070
    state->db->FindObsoleteFiles(&job_context, false, true);
1071 1072 1073
    if (state->background_purge) {
      state->db->ScheduleBgLogWriterClose(&job_context);
    }
1074 1075 1076
    state->mu->Unlock();

    delete state->super_version;
I
Igor Canadi 已提交
1077
    if (job_context.HaveSomethingToDelete()) {
1078
      if (state->background_purge) {
1079 1080 1081 1082 1083 1084 1085 1086 1087 1088
        // PurgeObsoleteFiles here does not delete files. Instead, it adds the
        // files to be deleted to a job queue, and deletes it in a separate
        // background thread.
        state->db->PurgeObsoleteFiles(job_context, true /* schedule only */);
        state->mu->Lock();
        state->db->SchedulePurge();
        state->mu->Unlock();
      } else {
        state->db->PurgeObsoleteFiles(job_context);
      }
1089
    }
I
Igor Canadi 已提交
1090
    job_context.Clean();
I
Igor Canadi 已提交
1091
  }
T
Tomislav Novak 已提交
1092

1093 1094
  delete state;
}
H
Hans Wennborg 已提交
1095
}  // namespace
1096

A
Andrew Kryczka 已提交
1097 1098 1099 1100
InternalIterator* DBImpl::NewInternalIterator(
    const ReadOptions& read_options, ColumnFamilyData* cfd,
    SuperVersion* super_version, Arena* arena,
    RangeDelAggregator* range_del_agg) {
S
sdong 已提交
1101
  InternalIterator* internal_iter;
1102
  assert(arena != nullptr);
A
Andrew Kryczka 已提交
1103
  assert(range_del_agg != nullptr);
1104
  // Need to create internal iterator from the arena.
1105 1106 1107
  MergeIteratorBuilder merge_iter_builder(
      &cfd->internal_comparator(), arena,
      !read_options.total_order_seek &&
1108
          super_version->mutable_cf_options.prefix_extractor != nullptr);
1109 1110
  // Collect iterator for mutable mem
  merge_iter_builder.AddIterator(
L
Lei Jin 已提交
1111
      super_version->mem->NewIterator(read_options, arena));
1112
  std::unique_ptr<InternalIterator> range_del_iter;
A
Andrew Kryczka 已提交
1113 1114
  Status s;
  if (!read_options.ignore_range_deletions) {
1115 1116
    range_del_iter.reset(
        super_version->mem->NewRangeTombstoneIterator(read_options));
A
Andrew Kryczka 已提交
1117 1118
    s = range_del_agg->AddTombstones(std::move(range_del_iter));
  }
1119
  // Collect all needed child iterators for immutable memtables
A
Andrew Kryczka 已提交
1120 1121 1122 1123 1124 1125 1126
  if (s.ok()) {
    super_version->imm->AddIterators(read_options, &merge_iter_builder);
    if (!read_options.ignore_range_deletions) {
      s = super_version->imm->AddRangeTombstoneIterators(read_options, arena,
                                                         range_del_agg);
    }
  }
1127
  TEST_SYNC_POINT_CALLBACK("DBImpl::NewInternalIterator:StatusCallback", &s);
A
Andrew Kryczka 已提交
1128 1129
  if (s.ok()) {
    // Collect iterators for files in L0 - Ln
S
Sagar Vemuri 已提交
1130 1131 1132 1133
    if (read_options.read_tier != kMemtableTier) {
      super_version->current->AddIterators(read_options, env_options_,
                                           &merge_iter_builder, range_del_agg);
    }
A
Andrew Kryczka 已提交
1134 1135 1136 1137 1138 1139 1140
    internal_iter = merge_iter_builder.Finish();
    IterState* cleanup =
        new IterState(this, &mutex_, super_version,
                      read_options.background_purge_on_iterator_cleanup);
    internal_iter->RegisterCleanup(CleanupIteratorState, cleanup, nullptr);

    return internal_iter;
1141 1142
  } else {
    CleanupSuperVersion(super_version);
A
Andrew Kryczka 已提交
1143
  }
1144
  return NewErrorInternalIterator<Slice>(s, arena);
J
jorlow@chromium.org 已提交
1145 1146
}

1147 1148 1149 1150
ColumnFamilyHandle* DBImpl::DefaultColumnFamily() const {
  return default_cf_handle_;
}

L
Lei Jin 已提交
1151
Status DBImpl::Get(const ReadOptions& read_options,
S
Siying Dong 已提交
1152 1153 1154
                   ColumnFamilyHandle* column_family, const Slice& key,
                   PinnableSlice* value) {
  return GetImpl(read_options, column_family, key, value);
L
Lei Jin 已提交
1155 1156 1157
}

Status DBImpl::GetImpl(const ReadOptions& read_options,
1158
                       ColumnFamilyHandle* column_family, const Slice& key,
1159
                       PinnableSlice* pinnable_val, bool* value_found,
Y
Yi Wu 已提交
1160
                       ReadCallback* callback, bool* is_blob_index) {
M
Maysam Yabandeh 已提交
1161
  assert(pinnable_val != nullptr);
L
Lei Jin 已提交
1162
  StopWatch sw(env_, stats_, DB_GET);
1163
  PERF_TIMER_GUARD(get_snapshot_time);
1164

1165 1166 1167
  auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
  auto cfd = cfh->cfd();

1168 1169 1170 1171 1172 1173 1174 1175 1176
  if (tracer_) {
    // TODO: This mutex should be removed later, to improve performance when
    // tracing is enabled.
    InstrumentedMutexLock lock(&trace_mutex_);
    if (tracer_) {
      tracer_->Get(column_family, key);
    }
  }

1177 1178 1179 1180 1181 1182
  // Acquire SuperVersion
  SuperVersion* sv = GetAndRefSuperVersion(cfd);

  TEST_SYNC_POINT("DBImpl::GetImpl:1");
  TEST_SYNC_POINT("DBImpl::GetImpl:2");

1183
  SequenceNumber snapshot;
L
Lei Jin 已提交
1184
  if (read_options.snapshot != nullptr) {
1185 1186 1187 1188 1189 1190 1191 1192
    // Note: In WritePrepared txns this is not necessary but not harmful
    // either.  Because prep_seq > snapshot => commit_seq > snapshot so if
    // a snapshot is specified we should be fine with skipping seq numbers
    // that are greater than that.
    //
    // In WriteUnprepared, we cannot set snapshot in the lookup key because we
    // may skip uncommitted data that should be visible to the transaction for
    // reading own writes.
1193 1194
    snapshot =
        reinterpret_cast<const SnapshotImpl*>(read_options.snapshot)->number_;
1195 1196 1197
    if (callback) {
      snapshot = std::max(snapshot, callback->MaxUnpreparedSequenceNumber());
    }
1198
  } else {
1199 1200 1201 1202 1203 1204
    // Since we get and reference the super version before getting
    // the snapshot number, without a mutex protection, it is possible
    // that a memtable switch happened in the middle and not all the
    // data for this snapshot is available. But it will contain all
    // the data available in the super version we have, which is also
    // a valid snapshot to read from.
1205 1206 1207 1208
    // We shouldn't get snapshot before finding and referencing the super
    // version because a flush happening in between may compact away data for
    // the snapshot, but the snapshot is earlier than the data overwriting it,
    // so users may see wrong results.
1209 1210 1211
    snapshot = last_seq_same_as_publish_seq_
                   ? versions_->LastSequence()
                   : versions_->LastPublishedSequence();
J
jorlow@chromium.org 已提交
1212
  }
1213 1214
  TEST_SYNC_POINT("DBImpl::GetImpl:3");
  TEST_SYNC_POINT("DBImpl::GetImpl:4");
1215

1216
  // Prepare to store a list of merge operations if merge occurs.
1217
  MergeContext merge_context;
1218
  RangeDelAggregator range_del_agg(cfd->internal_comparator(), snapshot);
1219

1220
  Status s;
1221
  // First look in the memtable, then in the immutable memtable (if any).
1222
  // s is both in/out. When in, s could either be OK or MergeInProgress.
1223
  // merge_operands will contain the sequence of merges in the latter case.
1224
  LookupKey lkey(key, snapshot);
L
Lei Jin 已提交
1225
  PERF_TIMER_STOP(get_snapshot_time);
1226

1227 1228
  bool skip_memtable = (read_options.read_tier == kPersistedTier &&
                        has_unpersisted_data_.load(std::memory_order_relaxed));
1229 1230
  bool done = false;
  if (!skip_memtable) {
M
Maysam Yabandeh 已提交
1231
    if (sv->mem->Get(lkey, pinnable_val->GetSelf(), &s, &merge_context,
Y
Yi Wu 已提交
1232
                     &range_del_agg, read_options, callback, is_blob_index)) {
1233
      done = true;
M
Maysam Yabandeh 已提交
1234
      pinnable_val->PinSelf();
1235
      RecordTick(stats_, MEMTABLE_HIT);
A
Andrew Kryczka 已提交
1236
    } else if ((s.ok() || s.IsMergeInProgress()) &&
M
Maysam Yabandeh 已提交
1237
               sv->imm->Get(lkey, pinnable_val->GetSelf(), &s, &merge_context,
Y
Yi Wu 已提交
1238 1239
                            &range_del_agg, read_options, callback,
                            is_blob_index)) {
1240
      done = true;
M
Maysam Yabandeh 已提交
1241
      pinnable_val->PinSelf();
1242 1243
      RecordTick(stats_, MEMTABLE_HIT);
    }
A
Andrew Kryczka 已提交
1244
    if (!done && !s.ok() && !s.IsMergeInProgress()) {
1245
      ReturnAndCleanupSuperVersion(cfd, sv);
A
Andrew Kryczka 已提交
1246 1247
      return s;
    }
1248 1249
  }
  if (!done) {
1250
    PERF_TIMER_GUARD(get_from_output_files_time);
M
Maysam Yabandeh 已提交
1251
    sv->current->Get(read_options, lkey, pinnable_val, &s, &merge_context,
Y
Yi Wu 已提交
1252 1253
                     &range_del_agg, value_found, nullptr, nullptr, callback,
                     is_blob_index);
L
Lei Jin 已提交
1254
    RecordTick(stats_, MEMTABLE_MISS);
1255
  }
1256

1257 1258
  {
    PERF_TIMER_GUARD(get_post_process_time);
1259

1260
    ReturnAndCleanupSuperVersion(cfd, sv);
1261

1262
    RecordTick(stats_, NUMBER_KEYS_READ);
M
Maysam Yabandeh 已提交
1263 1264 1265
    size_t size = pinnable_val->size();
    RecordTick(stats_, BYTES_READ, size);
    MeasureTime(stats_, BYTES_PER_READ, size);
1266
    PERF_COUNTER_ADD(get_read_bytes, size);
1267
  }
1268
  return s;
J
jorlow@chromium.org 已提交
1269 1270
}

1271
std::vector<Status> DBImpl::MultiGet(
L
Lei Jin 已提交
1272
    const ReadOptions& read_options,
1273
    const std::vector<ColumnFamilyHandle*>& column_family,
1274
    const std::vector<Slice>& keys, std::vector<std::string>* values) {
L
Lei Jin 已提交
1275
  StopWatch sw(env_, stats_, DB_MULTIGET);
1276
  PERF_TIMER_GUARD(get_snapshot_time);
K
Kai Liu 已提交
1277

1278
  SequenceNumber snapshot;
1279

1280
  struct MultiGetColumnFamilyData {
I
Igor Canadi 已提交
1281
    ColumnFamilyData* cfd;
1282 1283 1284 1285 1286
    SuperVersion* super_version;
  };
  std::unordered_map<uint32_t, MultiGetColumnFamilyData*> multiget_cf_data;
  // fill up and allocate outside of mutex
  for (auto cf : column_family) {
1287 1288 1289 1290 1291 1292
    auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(cf);
    auto cfd = cfh->cfd();
    if (multiget_cf_data.find(cfd->GetID()) == multiget_cf_data.end()) {
      auto mgcfd = new MultiGetColumnFamilyData();
      mgcfd->cfd = cfd;
      multiget_cf_data.insert({cfd->GetID(), mgcfd});
1293 1294 1295
    }
  }

1296
  mutex_.Lock();
L
Lei Jin 已提交
1297
  if (read_options.snapshot != nullptr) {
1298 1299
    snapshot =
        reinterpret_cast<const SnapshotImpl*>(read_options.snapshot)->number_;
1300
  } else {
1301 1302 1303
    snapshot = last_seq_same_as_publish_seq_
                   ? versions_->LastSequence()
                   : versions_->LastPublishedSequence();
1304
  }
1305
  for (auto mgd_iter : multiget_cf_data) {
1306 1307
    mgd_iter.second->super_version =
        mgd_iter.second->cfd->GetSuperVersion()->Ref();
1308
  }
1309
  mutex_.Unlock();
1310

1311 1312
  // Contain a list of merge operations if merge occurs.
  MergeContext merge_context;
1313

1314
  // Note: this always resizes the values array
1315 1316 1317
  size_t num_keys = keys.size();
  std::vector<Status> stat_list(num_keys);
  values->resize(num_keys);
1318 1319

  // Keep track of bytes that we read for statistics-recording later
1320
  uint64_t bytes_read = 0;
L
Lei Jin 已提交
1321
  PERF_TIMER_STOP(get_snapshot_time);
1322 1323 1324 1325

  // For each of the given keys, apply the entire "get" process as follows:
  // First look in the memtable, then in the immutable memtable (if any).
  // s is both in/out. When in, s could either be OK or MergeInProgress.
1326
  // merge_operands will contain the sequence of merges in the latter case.
1327
  size_t num_found = 0;
1328
  for (size_t i = 0; i < num_keys; ++i) {
1329
    merge_context.Clear();
1330
    Status& s = stat_list[i];
1331 1332 1333
    std::string* value = &(*values)[i];

    LookupKey lkey(keys[i], snapshot);
1334
    auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family[i]);
A
Andrew Kryczka 已提交
1335
    RangeDelAggregator range_del_agg(cfh->cfd()->internal_comparator(),
1336
                                     snapshot);
1337
    auto mgd_iter = multiget_cf_data.find(cfh->cfd()->GetID());
1338 1339 1340
    assert(mgd_iter != multiget_cf_data.end());
    auto mgd = mgd_iter->second;
    auto super_version = mgd->super_version;
1341
    bool skip_memtable =
1342 1343
        (read_options.read_tier == kPersistedTier &&
         has_unpersisted_data_.load(std::memory_order_relaxed));
1344 1345
    bool done = false;
    if (!skip_memtable) {
M
Maysam Yabandeh 已提交
1346
      if (super_version->mem->Get(lkey, value, &s, &merge_context,
A
Andrew Kryczka 已提交
1347
                                  &range_del_agg, read_options)) {
1348
        done = true;
1349
        RecordTick(stats_, MEMTABLE_HIT);
M
Maysam Yabandeh 已提交
1350
      } else if (super_version->imm->Get(lkey, value, &s, &merge_context,
A
Andrew Kryczka 已提交
1351
                                         &range_del_agg, read_options)) {
1352
        done = true;
1353
        RecordTick(stats_, MEMTABLE_HIT);
1354 1355 1356
      }
    }
    if (!done) {
M
Maysam Yabandeh 已提交
1357
      PinnableSlice pinnable_val;
1358
      PERF_TIMER_GUARD(get_from_output_files_time);
M
Maysam Yabandeh 已提交
1359 1360 1361
      super_version->current->Get(read_options, lkey, &pinnable_val, &s,
                                  &merge_context, &range_del_agg);
      value->assign(pinnable_val.data(), pinnable_val.size());
1362
      RecordTick(stats_, MEMTABLE_MISS);
1363 1364 1365
    }

    if (s.ok()) {
M
Maysam Yabandeh 已提交
1366
      bytes_read += value->size();
1367
      num_found++;
1368 1369 1370 1371
    }
  }

  // Post processing (decrement reference counts and record statistics)
1372
  PERF_TIMER_GUARD(get_post_process_time);
1373 1374
  autovector<SuperVersion*> superversions_to_delete;

I
Igor Canadi 已提交
1375
  // TODO(icanadi) do we need lock here or just around Cleanup()?
1376 1377 1378 1379 1380 1381
  mutex_.Lock();
  for (auto mgd_iter : multiget_cf_data) {
    auto mgd = mgd_iter.second;
    if (mgd->super_version->Unref()) {
      mgd->super_version->Cleanup();
      superversions_to_delete.push_back(mgd->super_version);
1382 1383
    }
  }
1384 1385 1386 1387 1388 1389 1390
  mutex_.Unlock();

  for (auto td : superversions_to_delete) {
    delete td;
  }
  for (auto mgd : multiget_cf_data) {
    delete mgd.second;
1391
  }
1392

L
Lei Jin 已提交
1393 1394
  RecordTick(stats_, NUMBER_MULTIGET_CALLS);
  RecordTick(stats_, NUMBER_MULTIGET_KEYS_READ, num_keys);
1395
  RecordTick(stats_, NUMBER_MULTIGET_KEYS_FOUND, num_found);
L
Lei Jin 已提交
1396
  RecordTick(stats_, NUMBER_MULTIGET_BYTES_READ, bytes_read);
1397
  MeasureTime(stats_, BYTES_PER_MULTIGET, bytes_read);
1398
  PERF_COUNTER_ADD(multiget_read_bytes, bytes_read);
L
Lei Jin 已提交
1399
  PERF_TIMER_STOP(get_post_process_time);
1400

1401
  return stat_list;
1402 1403
}

L
Lei Jin 已提交
1404
Status DBImpl::CreateColumnFamily(const ColumnFamilyOptions& cf_options,
Y
Yi Wu 已提交
1405
                                  const std::string& column_family,
1406
                                  ColumnFamilyHandle** handle) {
Y
Yi Wu 已提交
1407 1408 1409 1410 1411 1412 1413 1414 1415 1416 1417 1418 1419 1420 1421 1422 1423 1424 1425 1426 1427 1428 1429 1430 1431 1432 1433 1434 1435 1436 1437 1438 1439 1440 1441 1442 1443 1444 1445 1446 1447 1448 1449 1450 1451 1452 1453 1454 1455 1456 1457 1458 1459 1460 1461 1462 1463 1464 1465 1466 1467 1468 1469 1470 1471 1472 1473 1474
  assert(handle != nullptr);
  Status s = CreateColumnFamilyImpl(cf_options, column_family, handle);
  if (s.ok()) {
    s = WriteOptionsFile(true /*need_mutex_lock*/,
                         true /*need_enter_write_thread*/);
  }
  return s;
}

Status DBImpl::CreateColumnFamilies(
    const ColumnFamilyOptions& cf_options,
    const std::vector<std::string>& column_family_names,
    std::vector<ColumnFamilyHandle*>* handles) {
  assert(handles != nullptr);
  handles->clear();
  size_t num_cf = column_family_names.size();
  Status s;
  bool success_once = false;
  for (size_t i = 0; i < num_cf; i++) {
    ColumnFamilyHandle* handle;
    s = CreateColumnFamilyImpl(cf_options, column_family_names[i], &handle);
    if (!s.ok()) {
      break;
    }
    handles->push_back(handle);
    success_once = true;
  }
  if (success_once) {
    Status persist_options_status = WriteOptionsFile(
        true /*need_mutex_lock*/, true /*need_enter_write_thread*/);
    if (s.ok() && !persist_options_status.ok()) {
      s = persist_options_status;
    }
  }
  return s;
}

Status DBImpl::CreateColumnFamilies(
    const std::vector<ColumnFamilyDescriptor>& column_families,
    std::vector<ColumnFamilyHandle*>* handles) {
  assert(handles != nullptr);
  handles->clear();
  size_t num_cf = column_families.size();
  Status s;
  bool success_once = false;
  for (size_t i = 0; i < num_cf; i++) {
    ColumnFamilyHandle* handle;
    s = CreateColumnFamilyImpl(column_families[i].options,
                               column_families[i].name, &handle);
    if (!s.ok()) {
      break;
    }
    handles->push_back(handle);
    success_once = true;
  }
  if (success_once) {
    Status persist_options_status = WriteOptionsFile(
        true /*need_mutex_lock*/, true /*need_enter_write_thread*/);
    if (s.ok() && !persist_options_status.ok()) {
      s = persist_options_status;
    }
  }
  return s;
}

Status DBImpl::CreateColumnFamilyImpl(const ColumnFamilyOptions& cf_options,
                                      const std::string& column_family_name,
                                      ColumnFamilyHandle** handle) {
Y
Yueh-Hsuan Chiang 已提交
1475
  Status s;
1476
  Status persist_options_status;
I
Igor Canadi 已提交
1477
  *handle = nullptr;
1478 1479

  s = CheckCompressionSupported(cf_options);
1480
  if (s.ok() && immutable_db_options_.allow_concurrent_memtable_write) {
1481 1482
    s = CheckConcurrentWritesSupported(cf_options);
  }
1483 1484 1485 1486 1487 1488 1489 1490 1491 1492 1493
  if (s.ok()) {
    s = CheckCFPathsSupported(initial_db_options_, cf_options);
  }
  if (s.ok()) {
    for (auto& cf_path : cf_options.cf_paths) {
      s = env_->CreateDirIfMissing(cf_path.path);
      if (!s.ok()) {
        break;
      }
    }
  }
1494 1495 1496 1497
  if (!s.ok()) {
    return s;
  }

1498
  SuperVersionContext sv_context(/* create_superversion */ true);
Y
Yueh-Hsuan Chiang 已提交
1499
  {
1500
    InstrumentedMutexLock l(&mutex_);
I
Igor Canadi 已提交
1501

Y
Yueh-Hsuan Chiang 已提交
1502 1503 1504 1505 1506 1507 1508 1509 1510 1511 1512 1513 1514
    if (versions_->GetColumnFamilySet()->GetColumnFamily(column_family_name) !=
        nullptr) {
      return Status::InvalidArgument("Column family already exists");
    }
    VersionEdit edit;
    edit.AddColumnFamily(column_family_name);
    uint32_t new_id = versions_->GetColumnFamilySet()->GetNextColumnFamilyID();
    edit.SetColumnFamily(new_id);
    edit.SetLogNumber(logfile_number_);
    edit.SetComparatorName(cf_options.comparator->Name());

    // LogAndApply will both write the creation in MANIFEST and create
    // ColumnFamilyData object
I
Igor Canadi 已提交
1515
    {  // write thread
1516 1517
      WriteThread::Writer w;
      write_thread_.EnterUnbatched(&w, &mutex_);
I
Igor Canadi 已提交
1518 1519
      // LogAndApply will both write the creation in MANIFEST and create
      // ColumnFamilyData object
1520 1521 1522
      s = versions_->LogAndApply(nullptr, MutableCFOptions(cf_options), &edit,
                                 &mutex_, directories_.GetDbDir(), false,
                                 &cf_options);
1523
      write_thread_.ExitUnbatched(&w);
I
Igor Canadi 已提交
1524
    }
1525 1526 1527 1528 1529 1530
    if (s.ok()) {
      auto* cfd =
          versions_->GetColumnFamilySet()->GetColumnFamily(column_family_name);
      assert(cfd != nullptr);
      s = cfd->AddDirectories();
    }
Y
Yueh-Hsuan Chiang 已提交
1531 1532 1533 1534 1535
    if (s.ok()) {
      single_column_family_mode_ = false;
      auto* cfd =
          versions_->GetColumnFamilySet()->GetColumnFamily(column_family_name);
      assert(cfd != nullptr);
1536 1537
      InstallSuperVersionAndScheduleWork(cfd, &sv_context,
                                         *cfd->GetLatestMutableCFOptions());
S
sdong 已提交
1538 1539 1540 1541 1542

      if (!cfd->mem()->IsSnapshotSupported()) {
        is_snapshot_supported_ = false;
      }

1543 1544
      cfd->set_initialized();

Y
Yueh-Hsuan Chiang 已提交
1545
      *handle = new ColumnFamilyHandleImpl(cfd, this, &mutex_);
1546 1547 1548
      ROCKS_LOG_INFO(immutable_db_options_.info_log,
                     "Created column family [%s] (ID %u)",
                     column_family_name.c_str(), (unsigned)cfd->GetID());
Y
Yueh-Hsuan Chiang 已提交
1549
    } else {
1550 1551 1552
      ROCKS_LOG_ERROR(immutable_db_options_.info_log,
                      "Creating column family [%s] FAILED -- %s",
                      column_family_name.c_str(), s.ToString().c_str());
Y
Yueh-Hsuan Chiang 已提交
1553
    }
1554
  }  // InstrumentedMutexLock l(&mutex_)
Y
Yueh-Hsuan Chiang 已提交
1555

1556
  sv_context.Clean();
Y
Yueh-Hsuan Chiang 已提交
1557
  // this is outside the mutex
1558
  if (s.ok()) {
1559 1560
    NewThreadStatusCfInfo(
        reinterpret_cast<ColumnFamilyHandleImpl*>(*handle)->cfd());
1561
  }
1562
  return s;
1563 1564
}

1565
Status DBImpl::DropColumnFamily(ColumnFamilyHandle* column_family) {
Y
Yi Wu 已提交
1566 1567 1568 1569 1570 1571 1572 1573 1574 1575 1576 1577 1578 1579 1580 1581 1582 1583 1584 1585 1586 1587 1588 1589 1590 1591 1592 1593 1594 1595 1596
  assert(column_family != nullptr);
  Status s = DropColumnFamilyImpl(column_family);
  if (s.ok()) {
    s = WriteOptionsFile(true /*need_mutex_lock*/,
                         true /*need_enter_write_thread*/);
  }
  return s;
}

Status DBImpl::DropColumnFamilies(
    const std::vector<ColumnFamilyHandle*>& column_families) {
  Status s;
  bool success_once = false;
  for (auto* handle : column_families) {
    s = DropColumnFamilyImpl(handle);
    if (!s.ok()) {
      break;
    }
    success_once = true;
  }
  if (success_once) {
    Status persist_options_status = WriteOptionsFile(
        true /*need_mutex_lock*/, true /*need_enter_write_thread*/);
    if (s.ok() && !persist_options_status.ok()) {
      s = persist_options_status;
    }
  }
  return s;
}

Status DBImpl::DropColumnFamilyImpl(ColumnFamilyHandle* column_family) {
1597 1598 1599
  auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
  auto cfd = cfh->cfd();
  if (cfd->GetID() == 0) {
1600 1601
    return Status::InvalidArgument("Can't drop default column family");
  }
1602

S
sdong 已提交
1603 1604
  bool cf_support_snapshot = cfd->mem()->IsSnapshotSupported();

I
Igor Canadi 已提交
1605 1606
  VersionEdit edit;
  edit.DropColumnFamily();
1607 1608
  edit.SetColumnFamily(cfd->GetID());

1609
  Status s;
1610
  {
1611
    InstrumentedMutexLock l(&mutex_);
1612 1613 1614 1615
    if (cfd->IsDropped()) {
      s = Status::InvalidArgument("Column family already dropped!\n");
    }
    if (s.ok()) {
1616
      // we drop column family from a single write thread
1617 1618
      WriteThread::Writer w;
      write_thread_.EnterUnbatched(&w, &mutex_);
1619 1620
      s = versions_->LogAndApply(cfd, *cfd->GetLatestMutableCFOptions(), &edit,
                                 &mutex_);
1621
      write_thread_.ExitUnbatched(&w);
1622
    }
Y
Yi Wu 已提交
1623 1624 1625 1626 1627
    if (s.ok()) {
      auto* mutable_cf_options = cfd->GetLatestMutableCFOptions();
      max_total_in_memory_state_ -= mutable_cf_options->write_buffer_size *
                                    mutable_cf_options->max_write_buffer_number;
    }
S
sdong 已提交
1628 1629 1630 1631 1632 1633

    if (!cf_support_snapshot) {
      // Dropped Column Family doesn't support snapshot. Need to recalculate
      // is_snapshot_supported_.
      bool new_is_snapshot_supported = true;
      for (auto c : *versions_->GetColumnFamilySet()) {
1634
        if (!c->IsDropped() && !c->mem()->IsSnapshotSupported()) {
S
sdong 已提交
1635 1636 1637 1638 1639 1640
          new_is_snapshot_supported = false;
          break;
        }
      }
      is_snapshot_supported_ = new_is_snapshot_supported;
    }
1641
    bg_cv_.SignalAll();
1642
  }
1643

1644
  if (s.ok()) {
Y
Yueh-Hsuan Chiang 已提交
1645 1646 1647 1648
    // Note that here we erase the associated cf_info of the to-be-dropped
    // cfd before its ref-count goes to zero to avoid having to erase cf_info
    // later inside db_mutex.
    EraseThreadStatusCfInfo(cfd);
I
Igor Canadi 已提交
1649
    assert(cfd->IsDropped());
1650 1651
    ROCKS_LOG_INFO(immutable_db_options_.info_log,
                   "Dropped column family with id %u\n", cfd->GetID());
1652
  } else {
1653 1654 1655
    ROCKS_LOG_ERROR(immutable_db_options_.info_log,
                    "Dropping column family with id %u FAILED -- %s\n",
                    cfd->GetID(), s.ToString().c_str());
1656 1657
  }

1658
  return s;
1659 1660
}

L
Lei Jin 已提交
1661
bool DBImpl::KeyMayExist(const ReadOptions& read_options,
1662 1663
                         ColumnFamilyHandle* column_family, const Slice& key,
                         std::string* value, bool* value_found) {
M
Maysam Yabandeh 已提交
1664
  assert(value != nullptr);
1665
  if (value_found != nullptr) {
K
Kai Liu 已提交
1666 1667
    // falsify later if key-may-exist but can't fetch value
    *value_found = true;
1668
  }
L
Lei Jin 已提交
1669
  ReadOptions roptions = read_options;
1670
  roptions.read_tier = kBlockCacheTier;  // read from block cache only
M
Maysam Yabandeh 已提交
1671 1672 1673
  PinnableSlice pinnable_val;
  auto s = GetImpl(roptions, column_family, key, &pinnable_val, value_found);
  value->assign(pinnable_val.data(), pinnable_val.size());
K
Kai Liu 已提交
1674

1675
  // If block_cache is enabled and the index block of the table didn't
K
Kai Liu 已提交
1676 1677 1678
  // not present in block_cache, the return value will be Status::Incomplete.
  // In this case, key may still exist in the table.
  return s.ok() || s.IsIncomplete();
1679 1680
}

1681
Iterator* DBImpl::NewIterator(const ReadOptions& read_options,
1682
                              ColumnFamilyHandle* column_family) {
S
Siying Dong 已提交
1683 1684 1685 1686
  if (read_options.managed) {
    return NewErrorIterator(
        Status::NotSupported("Managed iterator is not supported anymore."));
  }
1687
  Iterator* result = nullptr;
1688 1689 1690 1691
  if (read_options.read_tier == kPersistedTier) {
    return NewErrorIterator(Status::NotSupported(
        "ReadTier::kPersistedData is not yet supported in iterators."));
  }
1692 1693 1694 1695 1696
  // if iterator wants internal keys, we can only proceed if
  // we can guarantee the deletes haven't been processed yet
  if (immutable_db_options_.preserve_deletes &&
      read_options.iter_start_seqnum > 0 &&
      read_options.iter_start_seqnum < preserve_deletes_seqnum_.load()) {
1697 1698 1699 1700
    return NewErrorIterator(Status::InvalidArgument(
        "Iterator requested internal keys which are too old and are not"
        " guaranteed to be preserved, try larger iter_start_seqnum opt."));
  }
1701 1702
  auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
  auto cfd = cfh->cfd();
Y
Yi Wu 已提交
1703
  ReadCallback* read_callback = nullptr;  // No read callback provided.
1704
  if (read_options.tailing) {
I
Igor Canadi 已提交
1705 1706
#ifdef ROCKSDB_LITE
    // not supported in lite version
1707 1708
    result = nullptr;

I
Igor Canadi 已提交
1709
#else
1710 1711
    SuperVersion* sv = cfd->GetReferencedSuperVersion(&mutex_);
    auto iter = new ForwardIterator(this, read_options, cfd, sv);
1712 1713 1714
    result = NewDBIterator(
        env_, read_options, *cfd->ioptions(), sv->mutable_cf_options,
        cfd->user_comparator(), iter, kMaxSequenceNumber,
1715 1716
        sv->mutable_cf_options.max_sequential_skip_in_iterations, read_callback,
        this, cfd);
I
Igor Canadi 已提交
1717
#endif
T
Tomislav Novak 已提交
1718
  } else {
1719
    // Note: no need to consider the special case of
1720
    // last_seq_same_as_publish_seq_==false since NewIterator is overridden in
1721
    // WritePreparedTxnDB
Y
Yi Wu 已提交
1722 1723 1724
    auto snapshot = read_options.snapshot != nullptr
                        ? read_options.snapshot->GetSequenceNumber()
                        : versions_->LastSequence();
1725
    result = NewIteratorImpl(read_options, cfd, snapshot, read_callback);
1726
  }
1727
  return result;
1728 1729
}

Y
Yi Wu 已提交
1730 1731 1732
ArenaWrappedDBIter* DBImpl::NewIteratorImpl(const ReadOptions& read_options,
                                            ColumnFamilyData* cfd,
                                            SequenceNumber snapshot,
Y
Yi Wu 已提交
1733
                                            ReadCallback* read_callback,
1734 1735
                                            bool allow_blob,
                                            bool allow_refresh) {
Y
Yi Wu 已提交
1736 1737 1738 1739 1740 1741 1742 1743 1744 1745 1746 1747 1748 1749 1750 1751 1752 1753 1754 1755 1756 1757 1758 1759 1760 1761 1762 1763 1764 1765 1766 1767 1768 1769 1770 1771 1772 1773 1774 1775 1776 1777 1778 1779 1780
  SuperVersion* sv = cfd->GetReferencedSuperVersion(&mutex_);

  // Try to generate a DB iterator tree in continuous memory area to be
  // cache friendly. Here is an example of result:
  // +-------------------------------+
  // |                               |
  // | ArenaWrappedDBIter            |
  // |  +                            |
  // |  +---> Inner Iterator   ------------+
  // |  |                            |     |
  // |  |    +-- -- -- -- -- -- -- --+     |
  // |  +--- | Arena                 |     |
  // |       |                       |     |
  // |          Allocated Memory:    |     |
  // |       |   +-------------------+     |
  // |       |   | DBIter            | <---+
  // |           |  +                |
  // |       |   |  +-> iter_  ------------+
  // |       |   |                   |     |
  // |       |   +-------------------+     |
  // |       |   | MergingIterator   | <---+
  // |           |  +                |
  // |       |   |  +->child iter1  ------------+
  // |       |   |  |                |          |
  // |           |  +->child iter2  ----------+ |
  // |       |   |  |                |        | |
  // |       |   |  +->child iter3  --------+ | |
  // |           |                   |      | | |
  // |       |   +-------------------+      | | |
  // |       |   | Iterator1         | <--------+
  // |       |   +-------------------+      | |
  // |       |   | Iterator2         | <------+
  // |       |   +-------------------+      |
  // |       |   | Iterator3         | <----+
  // |       |   +-------------------+
  // |       |                       |
  // +-------+-----------------------+
  //
  // ArenaWrappedDBIter inlines an arena area where all the iterators in
  // the iterator tree are allocated in the order of being accessed when
  // querying.
  // Laying out the iterators in the order of being accessed makes it more
  // likely that any iterator pointer is close to the iterator it points to so
  // that they are likely to be in the same cache line and/or page.
  ArenaWrappedDBIter* db_iter = NewArenaWrappedDbIterator(
1781
      env_, read_options, *cfd->ioptions(), sv->mutable_cf_options, snapshot,
Y
Yi Wu 已提交
1782
      sv->mutable_cf_options.max_sequential_skip_in_iterations,
1783 1784
      sv->version_number, read_callback, this, cfd, allow_blob,
      ((read_options.snapshot != nullptr) ? false : allow_refresh));
Y
Yi Wu 已提交
1785 1786 1787 1788 1789 1790 1791 1792 1793

  InternalIterator* internal_iter =
      NewInternalIterator(read_options, cfd, sv, db_iter->GetArena(),
                          db_iter->GetRangeDelAggregator());
  db_iter->SetIterUnderDBIter(internal_iter);

  return db_iter;
}

S
Siying Dong 已提交
1794 1795 1796 1797
Status DBImpl::NewIterators(
    const ReadOptions& read_options,
    const std::vector<ColumnFamilyHandle*>& column_families,
    std::vector<Iterator*>* iterators) {
S
Siying Dong 已提交
1798 1799 1800
  if (read_options.managed) {
    return Status::NotSupported("Managed iterator is not supported anymore.");
  }
S
Siying Dong 已提交
1801 1802 1803 1804
  if (read_options.read_tier == kPersistedTier) {
    return Status::NotSupported(
        "ReadTier::kPersistedData is not yet supported in iterators.");
  }
Y
Yi Wu 已提交
1805
  ReadCallback* read_callback = nullptr;  // No read callback provided.
S
Siying Dong 已提交
1806 1807
  iterators->clear();
  iterators->reserve(column_families.size());
S
Siying Dong 已提交
1808
  if (read_options.tailing) {
S
Siying Dong 已提交
1809 1810 1811 1812 1813 1814 1815 1816 1817
#ifdef ROCKSDB_LITE
    return Status::InvalidArgument(
        "Tailing interator not supported in RocksDB lite");
#else
    for (auto cfh : column_families) {
      auto cfd = reinterpret_cast<ColumnFamilyHandleImpl*>(cfh)->cfd();
      SuperVersion* sv = cfd->GetReferencedSuperVersion(&mutex_);
      auto iter = new ForwardIterator(this, read_options, cfd, sv);
      iterators->push_back(NewDBIterator(
1818 1819
          env_, read_options, *cfd->ioptions(), sv->mutable_cf_options,
          cfd->user_comparator(), iter, kMaxSequenceNumber,
Y
Yi Wu 已提交
1820
          sv->mutable_cf_options.max_sequential_skip_in_iterations,
1821
          read_callback, this, cfd));
1822
    }
S
Siying Dong 已提交
1823 1824
#endif
  } else {
1825
    // Note: no need to consider the special case of
1826
    // last_seq_same_as_publish_seq_==false since NewIterators is overridden in
1827
    // WritePreparedTxnDB
Y
Yi Wu 已提交
1828 1829 1830
    auto snapshot = read_options.snapshot != nullptr
                        ? read_options.snapshot->GetSequenceNumber()
                        : versions_->LastSequence();
S
Siying Dong 已提交
1831
    for (size_t i = 0; i < column_families.size(); ++i) {
1832 1833
      auto* cfd =
          reinterpret_cast<ColumnFamilyHandleImpl*>(column_families[i])->cfd();
Y
Yi Wu 已提交
1834 1835
      iterators->push_back(
          NewIteratorImpl(read_options, cfd, snapshot, read_callback));
1836 1837
    }
  }
1838

I
Igor Canadi 已提交
1839
  return Status::OK();
S
Stanislau Hlebik 已提交
1840 1841
}

S
Siying Dong 已提交
1842
const Snapshot* DBImpl::GetSnapshot() { return GetSnapshotImpl(false); }
1843

S
Siying Dong 已提交
1844 1845 1846
#ifndef ROCKSDB_LITE
const Snapshot* DBImpl::GetSnapshotForWriteConflictBoundary() {
  return GetSnapshotImpl(true);
1847 1848 1849
}
#endif  // ROCKSDB_LITE

1850
SnapshotImpl* DBImpl::GetSnapshotImpl(bool is_write_conflict_boundary) {
S
Siying Dong 已提交
1851 1852 1853 1854 1855 1856 1857 1858 1859
  int64_t unix_time = 0;
  env_->GetCurrentTime(&unix_time);  // Ignore error
  SnapshotImpl* s = new SnapshotImpl;

  InstrumentedMutexLock l(&mutex_);
  // returns null if the underlying memtable does not support snapshot.
  if (!is_snapshot_supported_) {
    delete s;
    return nullptr;
S
Sage Weil 已提交
1860
  }
1861
  auto snapshot_seq = last_seq_same_as_publish_seq_
1862
                          ? versions_->LastSequence()
1863
                          : versions_->LastPublishedSequence();
1864
  return snapshots_.New(s, snapshot_seq, unix_time, is_write_conflict_boundary);
S
Siying Dong 已提交
1865
}
1866

S
Siying Dong 已提交
1867 1868
void DBImpl::ReleaseSnapshot(const Snapshot* s) {
  const SnapshotImpl* casted_s = reinterpret_cast<const SnapshotImpl*>(s);
S
Stanislau Hlebik 已提交
1869
  {
S
Siying Dong 已提交
1870 1871
    InstrumentedMutexLock l(&mutex_);
    snapshots_.Delete(casted_s);
1872 1873
    uint64_t oldest_snapshot;
    if (snapshots_.empty()) {
1874
      oldest_snapshot = last_seq_same_as_publish_seq_
1875
                            ? versions_->LastSequence()
1876
                            : versions_->LastPublishedSequence();
1877 1878 1879 1880 1881 1882 1883 1884 1885 1886 1887 1888 1889
    } else {
      oldest_snapshot = snapshots_.oldest()->number_;
    }
    for (auto* cfd : *versions_->GetColumnFamilySet()) {
      cfd->current()->storage_info()->UpdateOldestSnapshot(oldest_snapshot);
      if (!cfd->current()
               ->storage_info()
               ->BottommostFilesMarkedForCompaction()
               .empty()) {
        SchedulePendingCompaction(cfd);
        MaybeScheduleFlushOrCompaction();
      }
    }
1890
  }
S
Siying Dong 已提交
1891
  delete casted_s;
1892 1893
}

I
Igor Canadi 已提交
1894
#ifndef ROCKSDB_LITE
I
Igor Canadi 已提交
1895 1896 1897 1898 1899
Status DBImpl::GetPropertiesOfAllTables(ColumnFamilyHandle* column_family,
                                        TablePropertiesCollection* props) {
  auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
  auto cfd = cfh->cfd();

1900 1901
  // Increment the ref count
  mutex_.Lock();
I
Igor Canadi 已提交
1902
  auto version = cfd->current();
1903 1904 1905 1906 1907 1908 1909 1910 1911 1912 1913 1914
  version->Ref();
  mutex_.Unlock();

  auto s = version->GetPropertiesOfAllTables(props);

  // Decrement the ref count
  mutex_.Lock();
  version->Unref();
  mutex_.Unlock();

  return s;
}
1915 1916

Status DBImpl::GetPropertiesOfTablesInRange(ColumnFamilyHandle* column_family,
1917
                                            const Range* range, std::size_t n,
1918 1919 1920 1921 1922 1923 1924 1925 1926 1927 1928 1929 1930 1931 1932 1933 1934 1935 1936 1937
                                            TablePropertiesCollection* props) {
  auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
  auto cfd = cfh->cfd();

  // Increment the ref count
  mutex_.Lock();
  auto version = cfd->current();
  version->Ref();
  mutex_.Unlock();

  auto s = version->GetPropertiesOfTablesInRange(range, n, props);

  // Decrement the ref count
  mutex_.Lock();
  version->Unref();
  mutex_.Unlock();

  return s;
}

I
Igor Canadi 已提交
1938
#endif  // ROCKSDB_LITE
1939

1940
const std::string& DBImpl::GetName() const { return dbname_; }
I
Igor Canadi 已提交
1941

1942
Env* DBImpl::GetEnv() const { return env_; }
1943

1944 1945
Options DBImpl::GetOptions(ColumnFamilyHandle* column_family) const {
  InstrumentedMutexLock l(&mutex_);
1946
  auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
1947 1948
  return Options(BuildDBOptions(immutable_db_options_, mutable_db_options_),
                 cfh->cfd()->GetLatestCFOptions());
I
Igor Canadi 已提交
1949 1950
}

1951 1952 1953 1954
DBOptions DBImpl::GetDBOptions() const {
  InstrumentedMutexLock l(&mutex_);
  return BuildDBOptions(immutable_db_options_, mutable_db_options_);
}
1955

1956
bool DBImpl::GetProperty(ColumnFamilyHandle* column_family,
1957
                         const Slice& property, std::string* value) {
1958
  const DBPropertyInfo* property_info = GetPropertyInfo(property);
1959
  value->clear();
1960
  auto cfd = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family)->cfd();
1961 1962 1963
  if (property_info == nullptr) {
    return false;
  } else if (property_info->handle_int) {
1964
    uint64_t int_value;
1965 1966
    bool ret_value =
        GetIntPropertyInternal(cfd, *property_info, false, &int_value);
1967
    if (ret_value) {
1968
      *value = ToString(int_value);
1969 1970
    }
    return ret_value;
1971
  } else if (property_info->handle_string) {
1972
    InstrumentedMutexLock l(&mutex_);
1973
    return cfd->internal_stats()->GetStringProperty(*property_info, property,
1974
                                                    value);
1975 1976 1977 1978 1979 1980 1981
  } else if (property_info->handle_string_dbimpl) {
    std::string tmp_value;
    bool ret_value = (this->*(property_info->handle_string_dbimpl))(&tmp_value);
    if (ret_value) {
      *value = tmp_value;
    }
    return ret_value;
1982
  }
1983 1984 1985 1986
  // Shouldn't reach here since exactly one of handle_string and handle_int
  // should be non-nullptr.
  assert(false);
  return false;
J
jorlow@chromium.org 已提交
1987 1988
}

1989 1990
bool DBImpl::GetMapProperty(ColumnFamilyHandle* column_family,
                            const Slice& property,
1991
                            std::map<std::string, std::string>* value) {
1992 1993 1994 1995 1996 1997 1998 1999 2000 2001 2002 2003 2004 2005 2006
  const DBPropertyInfo* property_info = GetPropertyInfo(property);
  value->clear();
  auto cfd = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family)->cfd();
  if (property_info == nullptr) {
    return false;
  } else if (property_info->handle_map) {
    InstrumentedMutexLock l(&mutex_);
    return cfd->internal_stats()->GetMapProperty(*property_info, property,
                                                 value);
  }
  // If we reach this point it means that handle_map is not provided for the
  // requested property
  return false;
}

2007 2008
bool DBImpl::GetIntProperty(ColumnFamilyHandle* column_family,
                            const Slice& property, uint64_t* value) {
2009 2010
  const DBPropertyInfo* property_info = GetPropertyInfo(property);
  if (property_info == nullptr || property_info->handle_int == nullptr) {
2011 2012
    return false;
  }
2013
  auto cfd = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family)->cfd();
2014
  return GetIntPropertyInternal(cfd, *property_info, false, value);
2015 2016
}

2017
bool DBImpl::GetIntPropertyInternal(ColumnFamilyData* cfd,
2018 2019 2020 2021
                                    const DBPropertyInfo& property_info,
                                    bool is_locked, uint64_t* value) {
  assert(property_info.handle_int != nullptr);
  if (!property_info.need_out_of_mutex) {
2022 2023
    if (is_locked) {
      mutex_.AssertHeld();
2024
      return cfd->internal_stats()->GetIntProperty(property_info, value, this);
2025 2026
    } else {
      InstrumentedMutexLock l(&mutex_);
2027
      return cfd->internal_stats()->GetIntProperty(property_info, value, this);
2028
    }
2029
  } else {
2030 2031 2032 2033 2034 2035
    SuperVersion* sv = nullptr;
    if (!is_locked) {
      sv = GetAndRefSuperVersion(cfd);
    } else {
      sv = cfd->GetSuperVersion();
    }
2036 2037

    bool ret = cfd->internal_stats()->GetIntPropertyOutOfMutex(
2038
        property_info, sv->current, value);
2039

2040 2041 2042
    if (!is_locked) {
      ReturnAndCleanupSuperVersion(cfd, sv);
    }
2043 2044 2045 2046 2047

    return ret;
  }
}

2048 2049 2050 2051 2052 2053 2054 2055 2056 2057
bool DBImpl::GetPropertyHandleOptionsStatistics(std::string* value) {
  assert(value != nullptr);
  Statistics* statistics = immutable_db_options_.statistics.get();
  if (!statistics) {
    return false;
  }
  *value = statistics->ToString();
  return true;
}

S
Siying Dong 已提交
2058 2059 2060 2061
#ifndef ROCKSDB_LITE
Status DBImpl::ResetStats() {
  InstrumentedMutexLock l(&mutex_);
  for (auto* cfd : *versions_->GetColumnFamilySet()) {
2062 2063 2064
    if (cfd->initialized()) {
      cfd->internal_stats()->Clear();
    }
S
Siying Dong 已提交
2065 2066 2067 2068 2069
  }
  return Status::OK();
}
#endif  // ROCKSDB_LITE

2070 2071
bool DBImpl::GetAggregatedIntProperty(const Slice& property,
                                      uint64_t* aggregated_value) {
2072 2073
  const DBPropertyInfo* property_info = GetPropertyInfo(property);
  if (property_info == nullptr || property_info->handle_int == nullptr) {
2074 2075 2076 2077 2078 2079 2080 2081 2082
    return false;
  }

  uint64_t sum = 0;
  {
    // Needs mutex to protect the list of column families.
    InstrumentedMutexLock l(&mutex_);
    uint64_t value;
    for (auto* cfd : *versions_->GetColumnFamilySet()) {
2083 2084 2085
      if (!cfd->initialized()) {
        continue;
      }
2086
      if (GetIntPropertyInternal(cfd, *property_info, true, &value)) {
2087 2088 2089 2090 2091 2092 2093 2094 2095 2096
        sum += value;
      } else {
        return false;
      }
    }
  }
  *aggregated_value = sum;
  return true;
}

2097 2098
SuperVersion* DBImpl::GetAndRefSuperVersion(ColumnFamilyData* cfd) {
  // TODO(ljin): consider using GetReferencedSuperVersion() directly
I
Igor Canadi 已提交
2099
  return cfd->GetThreadLocalSuperVersion(&mutex_);
2100 2101
}

A
agiardullo 已提交
2102 2103 2104 2105 2106 2107 2108 2109 2110 2111 2112 2113
// REQUIRED: this function should only be called on the write thread or if the
// mutex is held.
SuperVersion* DBImpl::GetAndRefSuperVersion(uint32_t column_family_id) {
  auto column_family_set = versions_->GetColumnFamilySet();
  auto cfd = column_family_set->GetColumnFamily(column_family_id);
  if (!cfd) {
    return nullptr;
  }

  return GetAndRefSuperVersion(cfd);
}

2114 2115 2116 2117 2118 2119 2120 2121 2122 2123 2124 2125 2126
void DBImpl::CleanupSuperVersion(SuperVersion* sv) {
  // Release SuperVersion
  if (sv->Unref()) {
    {
      InstrumentedMutexLock l(&mutex_);
      sv->Cleanup();
    }
    delete sv;
    RecordTick(stats_, NUMBER_SUPERVERSION_CLEANUPS);
  }
  RecordTick(stats_, NUMBER_SUPERVERSION_RELEASES);
}

2127 2128
void DBImpl::ReturnAndCleanupSuperVersion(ColumnFamilyData* cfd,
                                          SuperVersion* sv) {
2129 2130
  if (!cfd->ReturnThreadLocalSuperVersion(sv)) {
    CleanupSuperVersion(sv);
2131
  }
2132 2133
}

A
agiardullo 已提交
2134 2135 2136 2137 2138 2139 2140 2141 2142 2143 2144 2145 2146 2147 2148 2149 2150 2151 2152 2153 2154 2155 2156 2157
// REQUIRED: this function should only be called on the write thread.
void DBImpl::ReturnAndCleanupSuperVersion(uint32_t column_family_id,
                                          SuperVersion* sv) {
  auto column_family_set = versions_->GetColumnFamilySet();
  auto cfd = column_family_set->GetColumnFamily(column_family_id);

  // If SuperVersion is held, and we successfully fetched a cfd using
  // GetAndRefSuperVersion(), it must still exist.
  assert(cfd != nullptr);
  ReturnAndCleanupSuperVersion(cfd, sv);
}

// REQUIRED: this function should only be called on the write thread or if the
// mutex is held.
ColumnFamilyHandle* DBImpl::GetColumnFamilyHandle(uint32_t column_family_id) {
  ColumnFamilyMemTables* cf_memtables = column_family_memtables_.get();

  if (!cf_memtables->Seek(column_family_id)) {
    return nullptr;
  }

  return cf_memtables->GetColumnFamilyHandle();
}

A
Anirban Rahut 已提交
2158 2159 2160 2161 2162 2163 2164 2165 2166 2167 2168 2169 2170 2171
// REQUIRED: mutex is NOT held.
ColumnFamilyHandle* DBImpl::GetColumnFamilyHandleUnlocked(
    uint32_t column_family_id) {
  ColumnFamilyMemTables* cf_memtables = column_family_memtables_.get();

  InstrumentedMutexLock l(&mutex_);

  if (!cf_memtables->Seek(column_family_id)) {
    return nullptr;
  }

  return cf_memtables->GetColumnFamilyHandle();
}

2172 2173 2174 2175 2176 2177 2178 2179 2180 2181 2182 2183 2184 2185 2186 2187 2188 2189 2190 2191 2192 2193
void DBImpl::GetApproximateMemTableStats(ColumnFamilyHandle* column_family,
                                         const Range& range,
                                         uint64_t* const count,
                                         uint64_t* const size) {
  ColumnFamilyHandleImpl* cfh =
      reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
  ColumnFamilyData* cfd = cfh->cfd();
  SuperVersion* sv = GetAndRefSuperVersion(cfd);

  // Convert user_key into a corresponding internal key.
  InternalKey k1(range.start, kMaxSequenceNumber, kValueTypeForSeek);
  InternalKey k2(range.limit, kMaxSequenceNumber, kValueTypeForSeek);
  MemTable::MemTableStats memStats =
      sv->mem->ApproximateStats(k1.Encode(), k2.Encode());
  MemTable::MemTableStats immStats =
      sv->imm->ApproximateStats(k1.Encode(), k2.Encode());
  *count = memStats.count + immStats.count;
  *size = memStats.size + immStats.size;

  ReturnAndCleanupSuperVersion(cfd, sv);
}

2194
void DBImpl::GetApproximateSizes(ColumnFamilyHandle* column_family,
2195
                                 const Range* range, int n, uint64_t* sizes,
2196 2197 2198
                                 uint8_t include_flags) {
  assert(include_flags & DB::SizeApproximationFlags::INCLUDE_FILES ||
         include_flags & DB::SizeApproximationFlags::INCLUDE_MEMTABLES);
J
jorlow@chromium.org 已提交
2199
  Version* v;
2200 2201
  auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
  auto cfd = cfh->cfd();
2202 2203
  SuperVersion* sv = GetAndRefSuperVersion(cfd);
  v = sv->current;
J
jorlow@chromium.org 已提交
2204 2205 2206

  for (int i = 0; i < n; i++) {
    // Convert user_key into a corresponding internal key.
2207 2208
    InternalKey k1(range[i].start, kMaxSequenceNumber, kValueTypeForSeek);
    InternalKey k2(range[i].limit, kMaxSequenceNumber, kValueTypeForSeek);
V
Vitaliy Liptchinsky 已提交
2209
    sizes[i] = 0;
2210
    if (include_flags & DB::SizeApproximationFlags::INCLUDE_FILES) {
V
Vitaliy Liptchinsky 已提交
2211
      sizes[i] += versions_->ApproximateSize(v, k1.Encode(), k2.Encode());
2212 2213
    }
    if (include_flags & DB::SizeApproximationFlags::INCLUDE_MEMTABLES) {
2214 2215
      sizes[i] += sv->mem->ApproximateStats(k1.Encode(), k2.Encode()).size;
      sizes[i] += sv->imm->ApproximateStats(k1.Encode(), k2.Encode()).size;
2216
    }
J
jorlow@chromium.org 已提交
2217 2218
  }

2219
  ReturnAndCleanupSuperVersion(cfd, sv);
J
jorlow@chromium.org 已提交
2220 2221
}

I
Igor Canadi 已提交
2222 2223 2224 2225 2226 2227 2228 2229 2230 2231 2232 2233 2234 2235 2236 2237
std::list<uint64_t>::iterator
DBImpl::CaptureCurrentFileNumberInPendingOutputs() {
  // We need to remember the iterator of our insert, because after the
  // background job is done, we need to remove that element from
  // pending_outputs_.
  pending_outputs_.push_back(versions_->current_next_file_number());
  auto pending_outputs_inserted_elem = pending_outputs_.end();
  --pending_outputs_inserted_elem;
  return pending_outputs_inserted_elem;
}

void DBImpl::ReleaseFileNumberFromPendingOutputs(
    std::list<uint64_t>::iterator v) {
  pending_outputs_.erase(v);
}

I
Igor Canadi 已提交
2238 2239 2240 2241
#ifndef ROCKSDB_LITE
Status DBImpl::GetUpdatesSince(
    SequenceNumber seq, unique_ptr<TransactionLogIterator>* iter,
    const TransactionLogIterator::ReadOptions& read_options) {
D
Dmitri Smirnov 已提交
2242

L
Lei Jin 已提交
2243
  RecordTick(stats_, GET_UPDATES_SINCE_CALLS);
I
Igor Canadi 已提交
2244 2245 2246
  if (seq > versions_->LastSequence()) {
    return Status::NotFound("Requested sequence not yet written in the db");
  }
I
Igor Canadi 已提交
2247
  return wal_manager_.GetUpdatesSince(seq, iter, read_options, versions_.get());
I
Igor Canadi 已提交
2248 2249
}

2250 2251 2252
Status DBImpl::DeleteFile(std::string name) {
  uint64_t number;
  FileType type;
2253 2254 2255
  WalFileType log_type;
  if (!ParseFileName(name, &number, &type, &log_type) ||
      (type != kTableFile && type != kLogFile)) {
2256 2257
    ROCKS_LOG_ERROR(immutable_db_options_.info_log, "DeleteFile %s failed.\n",
                    name.c_str());
2258 2259 2260
    return Status::InvalidArgument("Invalid file name");
  }

2261 2262 2263 2264
  Status status;
  if (type == kLogFile) {
    // Only allow deleting archived log files
    if (log_type != kArchivedLogFile) {
2265 2266 2267
      ROCKS_LOG_ERROR(immutable_db_options_.info_log,
                      "DeleteFile %s failed - not archived log.\n",
                      name.c_str());
2268 2269
      return Status::NotSupported("Delete only supported for archived logs");
    }
C
Changli Gao 已提交
2270
    status = wal_manager_.DeleteFile(name, number);
2271
    if (!status.ok()) {
2272 2273 2274
      ROCKS_LOG_ERROR(immutable_db_options_.info_log,
                      "DeleteFile %s failed -- %s.\n", name.c_str(),
                      status.ToString().c_str());
2275 2276 2277 2278
    }
    return status;
  }

2279
  int level;
I
Igor Canadi 已提交
2280
  FileMetaData* metadata;
2281
  ColumnFamilyData* cfd;
2282
  VersionEdit edit;
2283
  JobContext job_context(next_job_id_.fetch_add(1), true);
D
Dhruba Borthakur 已提交
2284
  {
2285
    InstrumentedMutexLock l(&mutex_);
2286
    status = versions_->GetMetadataForFile(number, &level, &metadata, &cfd);
D
Dhruba Borthakur 已提交
2287
    if (!status.ok()) {
2288 2289
      ROCKS_LOG_WARN(immutable_db_options_.info_log,
                     "DeleteFile %s failed. File not found\n", name.c_str());
I
Igor Canadi 已提交
2290
      job_context.Clean();
D
Dhruba Borthakur 已提交
2291 2292
      return Status::InvalidArgument("File not found");
    }
2293
    assert(level < cfd->NumberLevels());
2294

D
Dhruba Borthakur 已提交
2295
    // If the file is being compacted no need to delete.
2296
    if (metadata->being_compacted) {
2297 2298 2299
      ROCKS_LOG_INFO(immutable_db_options_.info_log,
                     "DeleteFile %s Skipped. File about to be compacted\n",
                     name.c_str());
I
Igor Canadi 已提交
2300
      job_context.Clean();
D
Dhruba Borthakur 已提交
2301
      return Status::OK();
2302 2303
    }

D
Dhruba Borthakur 已提交
2304 2305 2306
    // Only the files in the last level can be deleted externally.
    // This is to make sure that any deletion tombstones are not
    // lost. Check that the level passed is the last level.
S
sdong 已提交
2307
    auto* vstoreage = cfd->current()->storage_info();
I
Igor Canadi 已提交
2308
    for (int i = level + 1; i < cfd->NumberLevels(); i++) {
S
sdong 已提交
2309
      if (vstoreage->NumLevelFiles(i) != 0) {
2310 2311 2312
        ROCKS_LOG_WARN(immutable_db_options_.info_log,
                       "DeleteFile %s FAILED. File not in last level\n",
                       name.c_str());
I
Igor Canadi 已提交
2313
        job_context.Clean();
D
Dhruba Borthakur 已提交
2314 2315 2316
        return Status::InvalidArgument("File not in last level");
      }
    }
2317
    // if level == 0, it has to be the oldest file
S
sdong 已提交
2318 2319
    if (level == 0 &&
        vstoreage->LevelFiles(0).back()->fd.GetNumber() != number) {
2320 2321 2322 2323
      ROCKS_LOG_WARN(immutable_db_options_.info_log,
                     "DeleteFile %s failed ---"
                     " target file in level 0 must be the oldest.",
                     name.c_str());
I
Igor Canadi 已提交
2324
      job_context.Clean();
2325 2326 2327
      return Status::InvalidArgument("File in level 0, but not oldest");
    }
    edit.SetColumnFamily(cfd->GetID());
D
Dhruba Borthakur 已提交
2328
    edit.DeleteFile(level, number);
2329
    status = versions_->LogAndApply(cfd, *cfd->GetLatestMutableCFOptions(),
2330
                                    &edit, &mutex_, directories_.GetDbDir());
I
Igor Canadi 已提交
2331
    if (status.ok()) {
Y
Yanqin Jin 已提交
2332 2333 2334
      InstallSuperVersionAndScheduleWork(
          cfd, &job_context.superversion_contexts[0],
          *cfd->GetLatestMutableCFOptions(), FlushReason::kDeleteFiles);
I
Igor Canadi 已提交
2335
    }
I
Igor Canadi 已提交
2336 2337
    FindObsoleteFiles(&job_context, false);
  }  // lock released here
2338

2339
  LogFlush(immutable_db_options_.info_log);
2340 2341 2342 2343 2344 2345 2346 2347 2348
  // remove files outside the db-lock
  if (job_context.HaveSomethingToDelete()) {
    // Call PurgeObsoleteFiles() without holding mutex.
    PurgeObsoleteFiles(job_context);
  }
  job_context.Clean();
  return status;
}

2349
Status DBImpl::DeleteFilesInRanges(ColumnFamilyHandle* column_family,
奏之章 已提交
2350
                                   const RangePtr* ranges, size_t n) {
2351 2352 2353 2354
  Status status;
  auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
  ColumnFamilyData* cfd = cfh->cfd();
  VersionEdit edit;
2355
  std::set<FileMetaData*> deleted_files;
2356 2357 2358
  JobContext job_context(next_job_id_.fetch_add(1), true);
  {
    InstrumentedMutexLock l(&mutex_);
2359
    Version* input_version = cfd->current();
2360
    edit.SetColumnFamily(cfd->GetID());
2361

2362
    auto* vstorage = input_version->storage_info();
2363 2364 2365 2366 2367 2368 2369 2370 2371 2372 2373 2374 2375 2376 2377 2378 2379 2380 2381 2382 2383 2384 2385 2386 2387 2388 2389 2390 2391 2392 2393
    if (cfd->GetCurrentMutableCFOptions()->enable_lazy_compaction) {
      const InternalKeyComparator& ic = cfd->ioptions()->internal_comparator;

      // deref nullptr of start/limit
      InternalKey* nullptr_start = nullptr;
      InternalKey* nullptr_limit = nullptr;
      for (int level = 0; level < cfd->NumberLevels(); level++) {
        auto& level_files = vstorage->LevelFiles(level);
        if (level == 0) {
          for (size_t i = 0; i < level_files.size(); ++i) {
            auto& f = level_files[i];
            if (nullptr_start == nullptr ||
                ic.Compare(f->smallest, *nullptr_start) < 0) {
              nullptr_start = &f->smallest;
            }
            if (nullptr_limit == nullptr ||
                ic.Compare(f->largest, *nullptr_limit) > 0) {
              nullptr_limit = &f->largest;
            }
          }
        } else if (!level_files.empty()) {
          auto& f0 = level_files.front();
          auto& fn = level_files.back();
          if (nullptr_start == nullptr ||
              ic.Compare(f0->smallest, *nullptr_start) < 0) {
            nullptr_start = &f0->smallest;
          }
          if (nullptr_limit == nullptr ||
              ic.Compare(fn->largest, *nullptr_limit) > 0) {
            nullptr_limit = &fn->largest;
          }
2394
        }
2395 2396 2397 2398 2399 2400
      }
      if (nullptr_start == nullptr || nullptr_limit == nullptr) {
        // empty vstorage ...
        job_context.Clean();
        return Status::OK();
      }
奏之章 已提交
2401
      // trans user_key to internal_key
奏之章 已提交
2402
      std::vector<std::pair<InternalKey, InternalKey>> deleted_range_storage;
2403
      std::vector<Range> deleted_range;
奏之章 已提交
2404
      deleted_range_storage.resize(n);
2405 2406
      deleted_range.resize(n);
      for (size_t i = 0; i < n; ++i) {
奏之章 已提交
2407 2408 2409 2410 2411 2412 2413 2414 2415 2416 2417 2418 2419 2420 2421 2422 2423 2424 2425 2426 2427 2428 2429 2430 2431
        deleted_range[i].include_start = ranges[i].include_start;
        deleted_range[i].include_limit = ranges[i].include_limit;
        auto& storage = deleted_range_storage[i];
        if (ranges[i].start == nullptr) {
          storage.first = *nullptr_start;
          deleted_range[i].include_start = true;
        } else {
          if (deleted_range[i].include_start) {
            storage.first.SetMinPossibleForUserKey(*ranges[i].start);
          } else {
            storage.first.SetMaxPossibleForUserKey(*ranges[i].start);
          }
        }
        if (ranges[i].limit == nullptr) {
          storage.second = *nullptr_limit;
          deleted_range[i].include_limit = true;
        } else {
          if (deleted_range[i].include_limit) {
            storage.second.SetMaxPossibleForUserKey(*ranges[i].limit);
          } else {
            storage.second.SetMinPossibleForUserKey(*ranges[i].limit);
          }
        }
        deleted_range[i].start = storage.first.Encode();
        deleted_range[i].limit = storage.second.Encode();
2432
      }
奏之章 已提交
2433
      // sort & merge ranges
2434
      std::sort(deleted_range.begin(), deleted_range.end(),
奏之章 已提交
2435 2436
                [&ic](const Range& rl, const Range& rr) {
                  return ic.Compare(rl.start, rr.start) < 0;
2437 2438 2439
                });
      size_t c = 0;
      for (size_t i = 1; i < n; ++i) {
奏之章 已提交
2440
        if (ic.Compare(deleted_range[c].limit, deleted_range[i].start) >= 0) {
2441
          deleted_range[c].include_start |= deleted_range[i].include_start;
奏之章 已提交
2442
          if (ic.Compare(deleted_range[c].limit, deleted_range[i].limit) <= 0) {
2443 2444 2445
            deleted_range[c].limit = deleted_range[i].limit;
            deleted_range[c].include_limit |= deleted_range[i].include_limit;
          }
2446
        } else {
奏之章 已提交
2447
          deleted_range[++c] = deleted_range[i];
2448 2449 2450 2451 2452 2453 2454 2455 2456 2457 2458 2459 2460 2461 2462 2463 2464
        }
      }
      deleted_range.resize(c + 1);
      MapBuilder map_builder(job_context.job_id, immutable_db_options_,
                             env_options_, versions_.get(), stats_,
                             table_cache_, dbname_);
      auto level_being_compacted = [vstorage](int level) {
        for (auto f : vstorage->LevelFiles(level)) {
          if (f->being_compacted) {
            return true;
          }
        }
        return false;
      };
      for (int i = 0; i < cfd->NumberLevels(); i++) {
        if (vstorage->LevelFiles(i).empty()) {
          continue;
2465
        }
2466 2467 2468 2469 2470
        if (cfd->ioptions()->compaction_style == kCompactionStyleUniversal &&
            !level_being_compacted(i)) {
          auto s = map_builder.Build(
              {CompactionInputFiles{i, vstorage->LevelFiles(i)}}, deleted_range,
              {}, kMapSst, i, vstorage->LevelFiles(i)[0]->fd.GetPathId(),
奏之章 已提交
2471
              vstorage, cfd, &edit);
2472 2473 2474
          if (!s.ok()) {
            return s;
          }
2475
        } else {
2476 2477 2478
          for (auto f : vstorage->LevelFiles(i)) {
            std::unique_ptr<TableProperties> prop;
            FileMetaData file_meta;
奏之章 已提交
2479 2480 2481
            auto s = map_builder.Build({CompactionInputFiles{i, {f}}},
                                       deleted_range, {}, kMapSst, i,
                                       f->fd.GetPathId(), vstorage, cfd, &edit);
2482 2483 2484 2485
            if (!s.ok()) {
              return s;
            }
          }
2486
        }
2487 2488 2489 2490 2491 2492 2493 2494 2495
      }
    } else {
      for (size_t r = 0; r < n; r++) {
        auto begin = ranges[r].start, end = ranges[r].limit;
        auto include_begin = ranges[r].include_start;
        auto include_end = ranges[r].include_limit;
        for (int i = 1; i < cfd->NumberLevels(); i++) {
          if (vstorage->LevelFiles(i).empty() ||
              !vstorage->OverlapInLevel(i, begin, end)) {
2496 2497
            continue;
          }
2498 2499 2500 2501 2502 2503 2504
          std::vector<FileMetaData*> level_files;
          InternalKey begin_storage, end_storage, *begin_key, *end_key;
          if (begin == nullptr) {
            begin_key = nullptr;
          } else {
            begin_storage.SetMinPossibleForUserKey(*begin);
            begin_key = &begin_storage;
2505
          }
2506 2507 2508 2509 2510
          if (end == nullptr) {
            end_key = nullptr;
          } else {
            end_storage.SetMaxPossibleForUserKey(*end);
            end_key = &end_storage;
奏之章 已提交
2511
          }
2512 2513 2514 2515 2516 2517 2518 2519 2520 2521 2522 2523 2524 2525 2526 2527 2528 2529 2530 2531 2532 2533 2534 2535 2536 2537

          vstorage->GetCleanInputsWithinInterval(
              i, begin_key, end_key, &level_files, -1 /* hint_index */,
              nullptr /* file_index */);
          FileMetaData* level_file;
          for (uint32_t j = 0; j < level_files.size(); j++) {
            level_file = level_files[j];
            if (level_file->being_compacted) {
              continue;
            }
            if (deleted_files.find(level_file) != deleted_files.end()) {
              continue;
            }
            if (!include_begin && begin != nullptr &&
                cfd->user_comparator()->Compare(level_file->smallest.user_key(),
                                                *begin) == 0) {
              continue;
            }
            if (!include_end && end != nullptr &&
                cfd->user_comparator()->Compare(level_file->largest.user_key(),
                                                *end) == 0) {
              continue;
            }
            edit.DeleteFile(i, level_file->fd.GetNumber());
            deleted_files.insert(level_file);
            level_file->being_compacted = true;
2538 2539
          }
        }
2540 2541
      }
    }
2542
    if (edit.GetNewFiles().empty() && edit.GetDeletedFiles().empty()) {
2543
      job_context.Clean();
2544 2545
      return Status::OK();
    }
2546
    input_version->Ref();
2547 2548 2549
    status = versions_->LogAndApply(cfd, *cfd->GetLatestMutableCFOptions(),
                                    &edit, &mutex_, directories_.GetDbDir());
    if (status.ok()) {
Y
Yanqin Jin 已提交
2550 2551 2552
      InstallSuperVersionAndScheduleWork(
          cfd, &job_context.superversion_contexts[0],
          *cfd->GetLatestMutableCFOptions(), FlushReason::kDeleteFiles);
2553
    }
2554 2555 2556 2557
    for (auto* deleted_file : deleted_files) {
      deleted_file->being_compacted = false;
    }
    input_version->Unref();
2558 2559 2560
    FindObsoleteFiles(&job_context, false);
  }  // lock released here

2561
  LogFlush(immutable_db_options_.info_log);
I
Igor Canadi 已提交
2562
  // remove files outside the db-lock
I
Igor Canadi 已提交
2563
  if (job_context.HaveSomethingToDelete()) {
2564
    // Call PurgeObsoleteFiles() without holding mutex.
I
Igor Canadi 已提交
2565
    PurgeObsoleteFiles(job_context);
2566
  }
I
Igor Canadi 已提交
2567
  job_context.Clean();
2568 2569 2570
  return status;
}

2571
void DBImpl::GetLiveFilesMetaData(std::vector<LiveFileMetaData>* metadata) {
2572
  InstrumentedMutexLock l(&mutex_);
2573
  versions_->GetLiveFilesMetaData(metadata);
2574
}
2575

2576 2577
void DBImpl::GetColumnFamilyMetaData(ColumnFamilyHandle* column_family,
                                     ColumnFamilyMetaData* cf_meta) {
2578 2579 2580 2581 2582 2583 2584
  assert(column_family);
  auto* cfd = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family)->cfd();
  auto* sv = GetAndRefSuperVersion(cfd);
  sv->current->GetColumnFamilyMetaData(cf_meta);
  ReturnAndCleanupSuperVersion(cfd, sv);
}

I
Igor Canadi 已提交
2585
#endif  // ROCKSDB_LITE
2586

I
Igor Canadi 已提交
2587 2588 2589 2590 2591 2592 2593
Status DBImpl::CheckConsistency() {
  mutex_.AssertHeld();
  std::vector<LiveFileMetaData> metadata;
  versions_->GetLiveFilesMetaData(&metadata);

  std::string corruption_messages;
  for (const auto& md : metadata) {
2594 2595
    // md.name has a leading "/".
    std::string file_path = md.db_path + md.name;
2596

I
Igor Canadi 已提交
2597 2598
    uint64_t fsize = 0;
    Status s = env_->GetFileSize(file_path, &fsize);
D
dyniusz 已提交
2599 2600 2601 2602
    if (!s.ok() &&
        env_->GetFileSize(Rocks2LevelTableFileName(file_path), &fsize).ok()) {
      s = Status::OK();
    }
I
Igor Canadi 已提交
2603 2604 2605 2606
    if (!s.ok()) {
      corruption_messages +=
          "Can't access " + md.name + ": " + s.ToString() + "\n";
    } else if (fsize != md.size) {
2607
      corruption_messages += "Sst file size mismatch: " + file_path +
I
Igor Canadi 已提交
2608
                             ". Size recorded in manifest " +
2609 2610
                             ToString(md.size) + ", actual size " +
                             ToString(fsize) + "\n";
I
Igor Canadi 已提交
2611 2612 2613 2614 2615 2616 2617 2618 2619
    }
  }
  if (corruption_messages.size() == 0) {
    return Status::OK();
  } else {
    return Status::Corruption(corruption_messages);
  }
}

2620
Status DBImpl::GetDbIdentity(std::string& identity) const {
2621 2622
  std::string idfilename = IdentityFileName(dbname_);
  const EnvOptions soptions;
2623 2624 2625 2626 2627 2628 2629 2630
  unique_ptr<SequentialFileReader> id_file_reader;
  Status s;
  {
    unique_ptr<SequentialFile> idfile;
    s = env_->NewSequentialFile(idfilename, &idfile, soptions);
    if (!s.ok()) {
      return s;
    }
2631 2632
    id_file_reader.reset(
        new SequentialFileReader(std::move(idfile), idfilename));
2633
  }
2634

2635 2636 2637 2638 2639
  uint64_t file_size;
  s = env_->GetFileSize(idfilename, &file_size);
  if (!s.ok()) {
    return s;
  }
奏之章 已提交
2640 2641
  char* buffer =
      reinterpret_cast<char*>(alloca(static_cast<size_t>(file_size)));
2642
  Slice id;
2643
  s = id_file_reader->Read(static_cast<size_t>(file_size), &id, buffer);
2644 2645 2646 2647 2648 2649 2650 2651 2652 2653 2654
  if (!s.ok()) {
    return s;
  }
  identity.assign(id.ToString());
  // If last character is '\n' remove it from identity
  if (identity.size() > 0 && identity.back() == '\n') {
    identity.pop_back();
  }
  return s;
}

2655
// Default implementation -- returns not supported status
A
Andrew Kryczka 已提交
2656 2657 2658
Status DB::CreateColumnFamily(const ColumnFamilyOptions& /*cf_options*/,
                              const std::string& /*column_family_name*/,
                              ColumnFamilyHandle** /*handle*/) {
2659
  return Status::NotSupported("");
2660
}
Y
Yi Wu 已提交
2661 2662

Status DB::CreateColumnFamilies(
A
Andrew Kryczka 已提交
2663 2664 2665
    const ColumnFamilyOptions& /*cf_options*/,
    const std::vector<std::string>& /*column_family_names*/,
    std::vector<ColumnFamilyHandle*>* /*handles*/) {
Y
Yi Wu 已提交
2666 2667 2668 2669
  return Status::NotSupported("");
}

Status DB::CreateColumnFamilies(
A
Andrew Kryczka 已提交
2670 2671
    const std::vector<ColumnFamilyDescriptor>& /*column_families*/,
    std::vector<ColumnFamilyHandle*>* /*handles*/) {
Y
Yi Wu 已提交
2672 2673 2674
  return Status::NotSupported("");
}

A
Andrew Kryczka 已提交
2675
Status DB::DropColumnFamily(ColumnFamilyHandle* /*column_family*/) {
2676
  return Status::NotSupported("");
2677
}
Y
Yi Wu 已提交
2678 2679

Status DB::DropColumnFamilies(
A
Andrew Kryczka 已提交
2680
    const std::vector<ColumnFamilyHandle*>& /*column_families*/) {
Y
Yi Wu 已提交
2681 2682 2683
  return Status::NotSupported("");
}

2684 2685 2686 2687
Status DB::DestroyColumnFamilyHandle(ColumnFamilyHandle* column_family) {
  delete column_family;
  return Status::OK();
}
2688

2689 2690 2691 2692 2693 2694 2695 2696 2697
DB::~DB() {}

Status DBImpl::Close() {
  if (!closed_) {
    closed_ = true;
    return CloseImpl();
  }
  return Status::OK();
}
J
jorlow@chromium.org 已提交
2698

2699 2700 2701
Status DB::ListColumnFamilies(const DBOptions& db_options,
                              const std::string& name,
                              std::vector<std::string>* column_families) {
I
Igor Canadi 已提交
2702
  return VersionSet::ListColumnFamilies(column_families, name, db_options.env);
2703 2704
}

2705
Snapshot::~Snapshot() {}
2706

2707 2708
Status DestroyDB(const std::string& dbname, const Options& options,
                 const std::vector<ColumnFamilyDescriptor>& column_families) {
D
Dmitri Smirnov 已提交
2709
  ImmutableDBOptions soptions(SanitizeOptions(dbname, options));
2710
  Env* env = soptions.env;
J
jorlow@chromium.org 已提交
2711
  std::vector<std::string> filenames;
2712

D
Dmitri Smirnov 已提交
2713 2714 2715
  // Reset the logger because it holds a handle to the
  // log file and prevents cleanup and directory removal
  soptions.info_log.reset();
J
jorlow@chromium.org 已提交
2716 2717
  // Ignore error in case directory does not exist
  env->GetChildren(dbname, &filenames);
2718

J
jorlow@chromium.org 已提交
2719
  FileLock* lock;
2720 2721
  const std::string lockname = LockFileName(dbname);
  Status result = env->LockFile(lockname, &lock);
J
jorlow@chromium.org 已提交
2722 2723 2724
  if (result.ok()) {
    uint64_t number;
    FileType type;
2725
    InfoLogPrefix info_log_prefix(!soptions.db_log_dir.empty(), dbname);
D
Dmitri Smirnov 已提交
2726 2727
    for (const auto& fname : filenames) {
      if (ParseFileName(fname, &number, info_log_prefix.prefix, &type) &&
2728
          type != kDBLockFile) {  // Lock file will be deleted at end
K
Kosie van der Merwe 已提交
2729
        Status del;
D
Dmitri Smirnov 已提交
2730
        std::string path_to_delete = dbname + "/" + fname;
K
Kosie van der Merwe 已提交
2731
        if (type == kMetaDatabase) {
2732 2733
          del = DestroyDB(path_to_delete, options);
        } else if (type == kTableFile) {
2734
          del = DeleteSSTFile(&soptions, path_to_delete, dbname);
K
Kosie van der Merwe 已提交
2735
        } else {
2736
          del = env->DeleteFile(path_to_delete);
K
Kosie van der Merwe 已提交
2737
        }
J
jorlow@chromium.org 已提交
2738 2739 2740 2741 2742
        if (result.ok() && !del.ok()) {
          result = del;
        }
      }
    }
2743

2744 2745
    std::vector<std::string> paths;

D
Dmitri Smirnov 已提交
2746 2747
    for (const auto& path : options.db_paths) {
      paths.emplace_back(path.path);
2748
    }
D
Dmitri Smirnov 已提交
2749 2750 2751
    for (const auto& cf : column_families) {
      for (const auto& path : cf.options.cf_paths) {
        paths.emplace_back(path.path);
2752 2753 2754 2755 2756 2757 2758 2759 2760 2761
      }
    }

    // Remove duplicate paths.
    // Note that we compare only the actual paths but not path ids.
    // This reason is that same path can appear at different path_ids
    // for different column families.
    std::sort(paths.begin(), paths.end());
    paths.erase(std::unique(paths.begin(), paths.end()), paths.end());

D
Dmitri Smirnov 已提交
2762 2763 2764 2765
    for (const auto& path : paths) {
      if (env->GetChildren(path, &filenames).ok()) {
        for (const auto& fname : filenames) {
          if (ParseFileName(fname, &number, &type) &&
2766
              type == kTableFile) {  // Lock file will be deleted at end
D
Dmitri Smirnov 已提交
2767 2768 2769 2770 2771
            std::string table_path = path + "/" + fname;
            Status del = DeleteSSTFile(&soptions, table_path, dbname);
            if (result.ok() && !del.ok()) {
              result = del;
            }
2772 2773
          }
        }
D
Dmitri Smirnov 已提交
2774
        env->DeleteDir(path);
2775 2776 2777
      }
    }

I
Igor Canadi 已提交
2778 2779
    std::vector<std::string> walDirFiles;
    std::string archivedir = ArchivalDirectory(dbname);
D
Dmitri Smirnov 已提交
2780
    bool wal_dir_exists = false;
I
Igor Canadi 已提交
2781
    if (dbname != soptions.wal_dir) {
D
Dmitri Smirnov 已提交
2782
      wal_dir_exists = env->GetChildren(soptions.wal_dir, &walDirFiles).ok();
I
Igor Canadi 已提交
2783 2784 2785
      archivedir = ArchivalDirectory(soptions.wal_dir);
    }

D
Dmitri Smirnov 已提交
2786 2787 2788 2789 2790 2791 2792
    // Archive dir may be inside wal dir or dbname and should be
    // processed and removed before those otherwise we have issues
    // removing them
    std::vector<std::string> archiveFiles;
    if (env->GetChildren(archivedir, &archiveFiles).ok()) {
      // Delete archival files.
      for (const auto& file : archiveFiles) {
2793
        if (ParseFileName(file, &number, &type) && type == kLogFile) {
D
Dmitri Smirnov 已提交
2794 2795 2796 2797
          Status del = env->DeleteFile(archivedir + "/" + file);
          if (result.ok() && !del.ok()) {
            result = del;
          }
I
Igor Canadi 已提交
2798 2799
        }
      }
D
Dmitri Smirnov 已提交
2800
      env->DeleteDir(archivedir);
I
Igor Canadi 已提交
2801 2802
    }

D
Dmitri Smirnov 已提交
2803 2804 2805 2806 2807 2808 2809 2810
    // Delete log files in the WAL dir
    if (wal_dir_exists) {
      for (const auto& file : walDirFiles) {
        if (ParseFileName(file, &number, &type) && type == kLogFile) {
          Status del = env->DeleteFile(LogFileName(soptions.wal_dir, number));
          if (result.ok() && !del.ok()) {
            result = del;
          }
2811 2812
        }
      }
D
Dmitri Smirnov 已提交
2813
      env->DeleteDir(soptions.wal_dir);
2814
    }
2815

J
jorlow@chromium.org 已提交
2816
    env->UnlockFile(lock);  // Ignore error since state is already gone
2817
    env->DeleteFile(lockname);
J
jorlow@chromium.org 已提交
2818 2819 2820 2821 2822
    env->DeleteDir(dbname);  // Ignore error in case dir contains other files
  }
  return result;
}

Y
Yi Wu 已提交
2823 2824
Status DBImpl::WriteOptionsFile(bool need_mutex_lock,
                                bool need_enter_write_thread) {
2825
#ifndef ROCKSDB_LITE
Y
Yi Wu 已提交
2826 2827 2828 2829 2830 2831 2832 2833 2834
  WriteThread::Writer w;
  if (need_mutex_lock) {
    mutex_.Lock();
  } else {
    mutex_.AssertHeld();
  }
  if (need_enter_write_thread) {
    write_thread_.EnterUnbatched(&w, &mutex_);
  }
2835 2836 2837

  std::vector<std::string> cf_names;
  std::vector<ColumnFamilyOptions> cf_opts;
2838 2839 2840 2841 2842

  // This part requires mutex to protect the column family options
  for (auto cfd : *versions_->GetColumnFamilySet()) {
    if (cfd->IsDropped()) {
      continue;
2843
    }
2844
    cf_names.push_back(cfd->GetName());
2845
    cf_opts.push_back(cfd->GetLatestCFOptions());
2846 2847
  }

2848 2849
  // Unlock during expensive operations.  New writes cannot get here
  // because the single write thread ensures all new writes get queued.
2850 2851
  DBOptions db_options =
      BuildDBOptions(immutable_db_options_, mutable_db_options_);
2852 2853
  mutex_.Unlock();

2854 2855 2856
  TEST_SYNC_POINT("DBImpl::WriteOptionsFile:1");
  TEST_SYNC_POINT("DBImpl::WriteOptionsFile:2");

2857 2858
  std::string file_name =
      TempOptionsFileName(GetName(), versions_->NewFileNumber());
2859 2860
  Status s =
      PersistRocksDBOptions(db_options, cf_names, cf_opts, file_name, GetEnv());
2861 2862 2863 2864

  if (s.ok()) {
    s = RenameTempFileToOptionsFile(file_name);
  }
Y
Yi Wu 已提交
2865 2866 2867 2868 2869 2870 2871 2872 2873 2874 2875 2876 2877 2878 2879
  // restore lock
  if (!need_mutex_lock) {
    mutex_.Lock();
  }
  if (need_enter_write_thread) {
    write_thread_.ExitUnbatched(&w);
  }
  if (!s.ok()) {
    ROCKS_LOG_WARN(immutable_db_options_.info_log,
                   "Unnable to persist options -- %s", s.ToString().c_str());
    if (immutable_db_options_.fail_if_options_file_error) {
      return Status::IOError("Unable to persist options.",
                             s.ToString().c_str());
    }
  }
2880 2881 2882
#else
  (void)need_mutex_lock;
  (void)need_enter_write_thread;
2883
#endif  // !ROCKSDB_LITE
Y
Yi Wu 已提交
2884
  return Status::OK();
2885 2886 2887 2888 2889 2890 2891 2892 2893 2894 2895 2896 2897 2898
}

#ifndef ROCKSDB_LITE
namespace {
void DeleteOptionsFilesHelper(const std::map<uint64_t, std::string>& filenames,
                              const size_t num_files_to_keep,
                              const std::shared_ptr<Logger>& info_log,
                              Env* env) {
  if (filenames.size() <= num_files_to_keep) {
    return;
  }
  for (auto iter = std::next(filenames.begin(), num_files_to_keep);
       iter != filenames.end(); ++iter) {
    if (!env->DeleteFile(iter->second).ok()) {
2899 2900
      ROCKS_LOG_WARN(info_log, "Unable to delete options file %s",
                     iter->second.c_str());
2901 2902 2903 2904 2905 2906 2907 2908 2909 2910 2911 2912 2913 2914 2915 2916 2917 2918 2919 2920 2921 2922 2923 2924 2925 2926 2927 2928 2929 2930
    }
  }
}
}  // namespace
#endif  // !ROCKSDB_LITE

Status DBImpl::DeleteObsoleteOptionsFiles() {
#ifndef ROCKSDB_LITE
  std::vector<std::string> filenames;
  // use ordered map to store keep the filenames sorted from the newest
  // to the oldest.
  std::map<uint64_t, std::string> options_filenames;
  Status s;
  s = GetEnv()->GetChildren(GetName(), &filenames);
  if (!s.ok()) {
    return s;
  }
  for (auto& filename : filenames) {
    uint64_t file_number;
    FileType type;
    if (ParseFileName(filename, &file_number, &type) && type == kOptionsFile) {
      options_filenames.insert(
          {std::numeric_limits<uint64_t>::max() - file_number,
           GetName() + "/" + filename});
    }
  }

  // Keeps the latest 2 Options file
  const size_t kNumOptionsFilesKept = 2;
  DeleteOptionsFilesHelper(options_filenames, kNumOptionsFilesKept,
2931
                           immutable_db_options_.info_log, GetEnv());
2932 2933 2934 2935 2936 2937 2938 2939 2940
  return Status::OK();
#else
  return Status::OK();
#endif  // !ROCKSDB_LITE
}

Status DBImpl::RenameTempFileToOptionsFile(const std::string& file_name) {
#ifndef ROCKSDB_LITE
  Status s;
W
Wanning Jiang 已提交
2941 2942

  versions_->options_file_number_ = versions_->NewFileNumber();
2943
  std::string options_file_name =
W
Wanning Jiang 已提交
2944
      OptionsFileName(GetName(), versions_->options_file_number_);
2945 2946 2947
  // Retry if the file name happen to conflict with an existing one.
  s = GetEnv()->RenameFile(file_name, options_file_name);

2948 2949 2950
  if (0 == disable_delete_obsolete_files_) {
    DeleteObsoleteOptionsFiles();
  }
2951 2952
  return s;
#else
2953
  (void)file_name;
2954 2955 2956 2957
  return Status::OK();
#endif  // !ROCKSDB_LITE
}

D
Daniel Black 已提交
2958
#ifdef ROCKSDB_USING_THREAD_STATUS
2959

2960
void DBImpl::NewThreadStatusCfInfo(ColumnFamilyData* cfd) const {
2961
  if (immutable_db_options_.enable_thread_tracking) {
2962 2963
    ThreadStatusUtil::NewColumnFamilyInfo(this, cfd, cfd->GetName(),
                                          cfd->ioptions()->env);
2964
  }
Y
Yueh-Hsuan Chiang 已提交
2965 2966
}

2967
void DBImpl::EraseThreadStatusCfInfo(ColumnFamilyData* cfd) const {
2968
  if (immutable_db_options_.enable_thread_tracking) {
2969
    ThreadStatusUtil::EraseColumnFamilyInfo(cfd);
2970
  }
Y
Yueh-Hsuan Chiang 已提交
2971 2972 2973
}

void DBImpl::EraseThreadStatusDbInfo() const {
2974
  if (immutable_db_options_.enable_thread_tracking) {
2975
    ThreadStatusUtil::EraseDatabaseInfo(this);
2976
  }
Y
Yueh-Hsuan Chiang 已提交
2977 2978 2979
}

#else
2980
void DBImpl::NewThreadStatusCfInfo(ColumnFamilyData* /*cfd*/) const {}
Y
Yueh-Hsuan Chiang 已提交
2981

2982
void DBImpl::EraseThreadStatusCfInfo(ColumnFamilyData* /*cfd*/) const {}
Y
Yueh-Hsuan Chiang 已提交
2983

2984
void DBImpl::EraseThreadStatusDbInfo() const {}
Y
Yueh-Hsuan Chiang 已提交
2985 2986
#endif  // ROCKSDB_USING_THREAD_STATUS

2987 2988
//
// A global method that can dump out the build version
2989
void DumpRocksDBBuildVersion(Logger* log) {
I
Igor Canadi 已提交
2990
#if !defined(IOS_CROSS_COMPILE)
H
hyunwoo 已提交
2991
  // if we compile with Xcode, we don't run build_detect_version, so we don't
I
Igor Canadi 已提交
2992
  // generate util/build_version.cc
2993 2994 2995 2996
  ROCKS_LOG_HEADER(log, "RocksDB version: %d.%d.%d\n", ROCKSDB_MAJOR,
                   ROCKSDB_MINOR, ROCKSDB_PATCH);
  ROCKS_LOG_HEADER(log, "Git sha %s", rocksdb_build_git_sha);
  ROCKS_LOG_HEADER(log, "Compile date %s", rocksdb_build_compile_date);
I
Igor Canadi 已提交
2997
#endif
2998 2999
}

A
agiardullo 已提交
3000 3001 3002 3003 3004 3005 3006 3007 3008 3009 3010 3011 3012 3013 3014 3015 3016
#ifndef ROCKSDB_LITE
SequenceNumber DBImpl::GetEarliestMemTableSequenceNumber(SuperVersion* sv,
                                                         bool include_history) {
  // Find the earliest sequence number that we know we can rely on reading
  // from the memtable without needing to check sst files.
  SequenceNumber earliest_seq =
      sv->imm->GetEarliestSequenceNumber(include_history);
  if (earliest_seq == kMaxSequenceNumber) {
    earliest_seq = sv->mem->GetEarliestSequenceNumber();
  }
  assert(sv->mem->GetEarliestSequenceNumber() >= earliest_seq);

  return earliest_seq;
}
#endif  // ROCKSDB_LITE

#ifndef ROCKSDB_LITE
3017 3018
Status DBImpl::GetLatestSequenceForKey(SuperVersion* sv, const Slice& key,
                                       bool cache_only, SequenceNumber* seq,
3019 3020
                                       bool* found_record_for_key,
                                       bool* is_blob_index) {
A
agiardullo 已提交
3021 3022
  Status s;
  MergeContext merge_context;
3023 3024
  RangeDelAggregator range_del_agg(sv->mem->GetInternalKeyComparator(),
                                   kMaxSequenceNumber);
A
agiardullo 已提交
3025

A
Andrew Kryczka 已提交
3026
  ReadOptions read_options;
A
agiardullo 已提交
3027 3028 3029 3030
  SequenceNumber current_seq = versions_->LastSequence();
  LookupKey lkey(key, current_seq);

  *seq = kMaxSequenceNumber;
3031 3032
  *found_record_for_key = false;

A
agiardullo 已提交
3033
  // Check if there is a record for this key in the latest memtable
M
Maysam Yabandeh 已提交
3034
  sv->mem->Get(lkey, nullptr, &s, &merge_context, &range_del_agg, seq,
3035
               read_options, nullptr /*read_callback*/, is_blob_index);
A
agiardullo 已提交
3036 3037 3038

  if (!(s.ok() || s.IsNotFound() || s.IsMergeInProgress())) {
    // unexpected error reading memtable.
3039 3040 3041
    ROCKS_LOG_ERROR(immutable_db_options_.info_log,
                    "Unexpected status returned from MemTable::Get: %s\n",
                    s.ToString().c_str());
A
agiardullo 已提交
3042 3043 3044 3045 3046 3047

    return s;
  }

  if (*seq != kMaxSequenceNumber) {
    // Found a sequence number, no need to check immutable memtables
3048
    *found_record_for_key = true;
A
agiardullo 已提交
3049 3050 3051 3052
    return Status::OK();
  }

  // Check if there is a record for this key in the immutable memtables
M
Maysam Yabandeh 已提交
3053
  sv->imm->Get(lkey, nullptr, &s, &merge_context, &range_del_agg, seq,
3054
               read_options, nullptr /*read_callback*/, is_blob_index);
A
agiardullo 已提交
3055 3056 3057

  if (!(s.ok() || s.IsNotFound() || s.IsMergeInProgress())) {
    // unexpected error reading memtable.
3058 3059 3060
    ROCKS_LOG_ERROR(immutable_db_options_.info_log,
                    "Unexpected status returned from MemTableList::Get: %s\n",
                    s.ToString().c_str());
A
agiardullo 已提交
3061 3062 3063 3064 3065 3066

    return s;
  }

  if (*seq != kMaxSequenceNumber) {
    // Found a sequence number, no need to check memtable history
3067
    *found_record_for_key = true;
A
agiardullo 已提交
3068 3069 3070 3071
    return Status::OK();
  }

  // Check if there is a record for this key in the immutable memtables
A
Andrew Kryczka 已提交
3072
  sv->imm->GetFromHistory(lkey, nullptr, &s, &merge_context, &range_del_agg,
3073
                          seq, read_options, is_blob_index);
A
agiardullo 已提交
3074 3075 3076

  if (!(s.ok() || s.IsNotFound() || s.IsMergeInProgress())) {
    // unexpected error reading memtable.
3077 3078
    ROCKS_LOG_ERROR(
        immutable_db_options_.info_log,
A
agiardullo 已提交
3079 3080 3081 3082 3083 3084
        "Unexpected status returned from MemTableList::GetFromHistory: %s\n",
        s.ToString().c_str());

    return s;
  }

3085 3086 3087 3088 3089 3090 3091 3092 3093 3094
  if (*seq != kMaxSequenceNumber) {
    // Found a sequence number, no need to check SST files
    *found_record_for_key = true;
    return Status::OK();
  }

  // TODO(agiardullo): possible optimization: consider checking cached
  // SST files if cache_only=true?
  if (!cache_only) {
    // Check tables
R
Reid Horuff 已提交
3095
    sv->current->Get(read_options, lkey, nullptr, &s, &merge_context,
A
Andrew Kryczka 已提交
3096
                     &range_del_agg, nullptr /* value_found */,
3097 3098
                     found_record_for_key, seq, nullptr /*read_callback*/,
                     is_blob_index);
3099 3100 3101

    if (!(s.ok() || s.IsNotFound() || s.IsMergeInProgress())) {
      // unexpected error reading SST files
3102 3103 3104
      ROCKS_LOG_ERROR(immutable_db_options_.info_log,
                      "Unexpected status returned from Version::Get: %s\n",
                      s.ToString().c_str());
3105 3106 3107 3108 3109

      return s;
    }
  }

A
agiardullo 已提交
3110 3111
  return Status::OK();
}
3112 3113 3114 3115 3116

Status DBImpl::IngestExternalFile(
    ColumnFamilyHandle* column_family,
    const std::vector<std::string>& external_files,
    const IngestExternalFileOptions& ingestion_options) {
3117 3118 3119 3120
  if (external_files.empty()) {
    return Status::InvalidArgument("external_files is empty");
  }

3121 3122 3123 3124
  Status status;
  auto cfh = reinterpret_cast<ColumnFamilyHandleImpl*>(column_family);
  auto cfd = cfh->cfd();

3125 3126 3127 3128 3129
  // Ingest should immediately fail if ingest_behind is requested,
  // but the DB doesn't support it.
  if (ingestion_options.ingest_behind) {
    if (!immutable_db_options_.allow_ingest_behind) {
      return Status::InvalidArgument(
3130
          "Can't ingest_behind file in DB with allow_ingest_behind=false");
3131 3132 3133
    }
  }

3134 3135 3136 3137
  ExternalSstFileIngestionJob ingestion_job(env_, versions_.get(), cfd,
                                            immutable_db_options_, env_options_,
                                            &snapshots_, ingestion_options);

3138 3139 3140
  SuperVersionContext dummy_sv_ctx(/* create_superversion */ true);
  VersionEdit dummy_edit;
  uint64_t next_file_number = 0;
3141 3142 3143
  std::list<uint64_t>::iterator pending_output_elem;
  {
    InstrumentedMutexLock l(&mutex_);
3144
    if (error_handler_.IsDBStopped()) {
3145
      // Don't ingest files when there is a bg_error
3146
      return error_handler_.GetBGError();
3147 3148 3149
    }

    // Make sure that bg cleanup wont delete the files that we are ingesting
3150
    pending_output_elem = CaptureCurrentFileNumberInPendingOutputs();
3151 3152 3153 3154 3155 3156 3157 3158 3159 3160 3161 3162 3163 3164 3165

    // If crash happen after a hard link established, Recover function may
    // reuse the file number that has already assigned to the internal file,
    // and this will overwrite the external file. To protect the external
    // file, we have to make sure the file number will never being reused.
    next_file_number = versions_->FetchAddFileNumber(external_files.size());
    auto cf_options = cfd->GetLatestMutableCFOptions();
    status = versions_->LogAndApply(cfd, *cf_options, &dummy_edit, &mutex_,
                                    directories_.GetDbDir());
    if (status.ok()) {
      InstallSuperVersionAndScheduleWork(cfd, &dummy_sv_ctx, *cf_options);
    }
  }
  dummy_sv_ctx.Clean();
  if (!status.ok()) {
3166 3167
    InstrumentedMutexLock l(&mutex_);
    ReleaseFileNumberFromPendingOutputs(pending_output_elem);
3168
    return status;
3169 3170
  }

3171
  SuperVersion* super_version = cfd->GetReferencedSuperVersion(&mutex_);
3172 3173
  status =
      ingestion_job.Prepare(external_files, next_file_number, super_version);
3174
  CleanupSuperVersion(super_version);
3175
  if (!status.ok()) {
3176
    InstrumentedMutexLock l(&mutex_);
3177
    ReleaseFileNumberFromPendingOutputs(pending_output_elem);
3178 3179 3180
    return status;
  }

3181
  SuperVersionContext sv_context(/* create_superversion */ true);
3182
  TEST_SYNC_POINT("DBImpl::AddFile:Start");
3183 3184 3185 3186 3187
  {
    // Lock db mutex
    InstrumentedMutexLock l(&mutex_);
    TEST_SYNC_POINT("DBImpl::AddFile:MutexLock");

3188
    // Stop writes to the DB by entering both write threads
3189 3190
    WriteThread::Writer w;
    write_thread_.EnterUnbatched(&w, &mutex_);
3191
    WriteThread::Writer nonmem_w;
3192
    if (two_write_queues_) {
3193 3194
      nonmem_write_thread_.EnterUnbatched(&nonmem_w, &mutex_);
    }
3195

3196 3197
    num_running_ingest_file_++;

3198 3199 3200 3201 3202 3203
    // We cannot ingest a file into a dropped CF
    if (cfd->IsDropped()) {
      status = Status::InvalidArgument(
          "Cannot ingest an external file into a dropped CF");
    }

3204
    // Figure out if we need to flush the memtable first
3205 3206
    if (status.ok()) {
      bool need_flush = false;
3207
      status = ingestion_job.NeedsFlush(&need_flush, cfd->GetSuperVersion());
3208 3209
      TEST_SYNC_POINT_CALLBACK("DBImpl::IngestExternalFile:NeedFlush",
                               &need_flush);
3210 3211
      if (status.ok() && need_flush) {
        mutex_.Unlock();
3212 3213 3214
        status = FlushMemTable(cfd, FlushOptions(),
                               FlushReason::kExternalFileIngestion,
                               true /* writes_stopped */);
3215 3216
        mutex_.Lock();
      }
3217 3218 3219 3220 3221 3222 3223 3224 3225 3226 3227 3228 3229 3230 3231
    }

    // Run the ingestion job
    if (status.ok()) {
      status = ingestion_job.Run();
    }

    // Install job edit [Mutex will be unlocked here]
    auto mutable_cf_options = cfd->GetLatestMutableCFOptions();
    if (status.ok()) {
      status =
          versions_->LogAndApply(cfd, *mutable_cf_options, ingestion_job.edit(),
                                 &mutex_, directories_.GetDbDir());
    }
    if (status.ok()) {
3232
      InstallSuperVersionAndScheduleWork(cfd, &sv_context, *mutable_cf_options,
Z
Zhongyi Xie 已提交
3233
                                         FlushReason::kExternalFileIngestion);
3234 3235 3236
    }

    // Resume writes to the DB
3237
    if (two_write_queues_) {
3238 3239
      nonmem_write_thread_.ExitUnbatched(&nonmem_w);
    }
3240 3241 3242 3243 3244 3245 3246 3247 3248 3249 3250 3251 3252 3253 3254 3255 3256 3257 3258
    write_thread_.ExitUnbatched(&w);

    // Update stats
    if (status.ok()) {
      ingestion_job.UpdateStats();
    }

    ReleaseFileNumberFromPendingOutputs(pending_output_elem);

    num_running_ingest_file_--;
    if (num_running_ingest_file_ == 0) {
      bg_cv_.SignalAll();
    }

    TEST_SYNC_POINT("DBImpl::AddFile:MutexUnlock");
  }
  // mutex_ is unlocked here

  // Cleanup
3259
  sv_context.Clean();
3260 3261
  ingestion_job.Cleanup(status);

3262 3263 3264 3265
  if (status.ok()) {
    NotifyOnExternalFileIngested(cfd, ingestion_job);
  }

3266 3267 3268
  return status;
}

A
Aaron G 已提交
3269 3270 3271 3272 3273 3274 3275 3276 3277 3278 3279 3280 3281 3282 3283 3284 3285 3286
Status DBImpl::VerifyChecksum() {
  Status s;
  std::vector<ColumnFamilyData*> cfd_list;
  {
    InstrumentedMutexLock l(&mutex_);
    for (auto cfd : *versions_->GetColumnFamilySet()) {
      if (!cfd->IsDropped() && cfd->initialized()) {
        cfd->Ref();
        cfd_list.push_back(cfd);
      }
    }
  }
  std::vector<SuperVersion*> sv_list;
  for (auto cfd : cfd_list) {
    sv_list.push_back(cfd->GetReferencedSuperVersion(&mutex_));
  }
  for (auto& sv : sv_list) {
    VersionStorageInfo* vstorage = sv->current->storage_info();
3287
    ColumnFamilyData* cfd = sv->current->cfd();
3288 3289 3290 3291
    Options opts;
    {
      InstrumentedMutexLock l(&mutex_);
      opts = Options(BuildDBOptions(immutable_db_options_,
3292
          mutable_db_options_), cfd->GetLatestCFOptions());
3293
    }
A
Aaron G 已提交
3294 3295 3296 3297
    for (int i = 0; i < vstorage->num_non_empty_levels() && s.ok(); i++) {
      for (size_t j = 0; j < vstorage->LevelFilesBrief(i).num_files && s.ok();
           j++) {
        const auto& fd = vstorage->LevelFilesBrief(i).files[j].fd;
3298
        std::string fname = TableFileName(cfd->ioptions()->cf_paths,
A
Aaron G 已提交
3299
                                          fd.GetNumber(), fd.GetPathId());
3300
        s = rocksdb::VerifySstFileChecksum(opts, env_options_, fname);
A
Aaron G 已提交
3301 3302 3303 3304 3305 3306 3307 3308 3309 3310 3311 3312 3313 3314 3315
      }
    }
    if (!s.ok()) {
      break;
    }
  }
  {
    InstrumentedMutexLock l(&mutex_);
    for (auto sv : sv_list) {
      if (sv && sv->Unref()) {
        sv->Cleanup();
        delete sv;
      }
    }
    for (auto cfd : cfd_list) {
3316
      cfd->Unref();
A
Aaron G 已提交
3317 3318 3319 3320 3321
    }
  }
  return s;
}

3322 3323 3324 3325 3326 3327 3328 3329 3330 3331 3332 3333 3334 3335 3336 3337 3338 3339 3340
void DBImpl::NotifyOnExternalFileIngested(
    ColumnFamilyData* cfd, const ExternalSstFileIngestionJob& ingestion_job) {
  if (immutable_db_options_.listeners.empty()) {
    return;
  }

  for (const IngestedFileInfo& f : ingestion_job.files_to_ingest()) {
    ExternalFileIngestionInfo info;
    info.cf_name = cfd->GetName();
    info.external_file_path = f.external_file_path;
    info.internal_file_path = f.internal_file_path;
    info.global_seqno = f.assigned_seqno;
    info.table_properties = f.table_properties;
    for (auto listener : immutable_db_options_.listeners) {
      listener->OnExternalFileIngested(this, info);
    }
  }
}

3341 3342 3343 3344 3345 3346 3347
void DBImpl::WaitForIngestFile() {
  mutex_.AssertHeld();
  while (num_running_ingest_file_ > 0) {
    bg_cv_.Wait();
  }
}

3348 3349 3350 3351 3352 3353 3354 3355 3356 3357 3358 3359 3360 3361
Status DBImpl::StartTrace(const TraceOptions& /* options */,
                          std::unique_ptr<TraceWriter>&& trace_writer) {
  InstrumentedMutexLock lock(&trace_mutex_);
  tracer_.reset(new Tracer(env_, std::move(trace_writer)));
  return Status::OK();
}

Status DBImpl::EndTrace() {
  InstrumentedMutexLock lock(&trace_mutex_);
  Status s = tracer_->Close();
  tracer_.reset();
  return s;
}

3362 3363 3364 3365 3366 3367 3368 3369 3370 3371 3372
Status DBImpl::TraceIteratorSeek(const uint32_t& cf_id, const Slice& key) {
  Status s;
  if (tracer_) {
    InstrumentedMutexLock lock(&trace_mutex_);
    if (tracer_) {
      s = tracer_->IteratorSeek(cf_id, key);
    }
  }
  return s;
}

3373 3374
Status DBImpl::TraceIteratorSeekForPrev(const uint32_t& cf_id,
                                        const Slice& key) {
3375 3376 3377 3378 3379 3380 3381 3382 3383 3384
  Status s;
  if (tracer_) {
    InstrumentedMutexLock lock(&trace_mutex_);
    if (tracer_) {
      s = tracer_->IteratorSeekForPrev(cf_id, key);
    }
  }
  return s;
}

A
agiardullo 已提交
3385
#endif  // ROCKSDB_LITE
3386

3387
}  // namespace rocksdb