compaction_job.cc 70.0 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).
I
Igor Canadi 已提交
5 6 7 8 9
//
// 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.

10 11
#include "db/compaction/compaction_job.h"

I
Igor Canadi 已提交
12
#include <algorithm>
13
#include <cinttypes>
14
#include <functional>
I
Igor Canadi 已提交
15
#include <list>
16 17
#include <memory>
#include <random>
18 19 20
#include <set>
#include <thread>
#include <utility>
21
#include <vector>
I
Igor Canadi 已提交
22 23

#include "db/builder.h"
24
#include "db/db_impl/db_impl.h"
I
Igor Canadi 已提交
25 26
#include "db/db_iter.h"
#include "db/dbformat.h"
27
#include "db/error_handler.h"
28
#include "db/event_helpers.h"
I
Igor Canadi 已提交
29 30 31 32 33
#include "db/log_reader.h"
#include "db/log_writer.h"
#include "db/memtable.h"
#include "db/memtable_list.h"
#include "db/merge_context.h"
34
#include "db/merge_helper.h"
35
#include "db/range_del_aggregator.h"
I
Igor Canadi 已提交
36
#include "db/version_set.h"
37
#include "file/filename.h"
38
#include "file/read_write_util.h"
39
#include "file/sst_file_manager_impl.h"
40
#include "file/writable_file_writer.h"
41 42
#include "logging/log_buffer.h"
#include "logging/logging.h"
43 44 45
#include "monitoring/iostats_context_imp.h"
#include "monitoring/perf_context_imp.h"
#include "monitoring/thread_status_util.h"
46
#include "port/port.h"
I
Igor Canadi 已提交
47 48
#include "rocksdb/db.h"
#include "rocksdb/env.h"
49
#include "rocksdb/sst_partitioner.h"
I
Igor Canadi 已提交
50 51 52
#include "rocksdb/statistics.h"
#include "rocksdb/status.h"
#include "rocksdb/table.h"
53 54
#include "table/block_based/block.h"
#include "table/block_based/block_based_table_factory.h"
55
#include "table/merging_iterator.h"
I
Igor Canadi 已提交
56
#include "table/table_builder.h"
57
#include "test_util/sync_point.h"
I
Igor Canadi 已提交
58
#include "util/coding.h"
59
#include "util/hash.h"
I
Igor Canadi 已提交
60
#include "util/mutexlock.h"
61
#include "util/random.h"
I
Igor Canadi 已提交
62
#include "util/stop_watch.h"
63
#include "util/string_util.h"
I
Igor Canadi 已提交
64

65
namespace ROCKSDB_NAMESPACE {
I
Igor Canadi 已提交
66

67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98
const char* GetCompactionReasonString(CompactionReason compaction_reason) {
  switch (compaction_reason) {
    case CompactionReason::kUnknown:
      return "Unknown";
    case CompactionReason::kLevelL0FilesNum:
      return "LevelL0FilesNum";
    case CompactionReason::kLevelMaxLevelSize:
      return "LevelMaxLevelSize";
    case CompactionReason::kUniversalSizeAmplification:
      return "UniversalSizeAmplification";
    case CompactionReason::kUniversalSizeRatio:
      return "UniversalSizeRatio";
    case CompactionReason::kUniversalSortedRunNum:
      return "UniversalSortedRunNum";
    case CompactionReason::kFIFOMaxSize:
      return "FIFOMaxSize";
    case CompactionReason::kFIFOReduceNumFiles:
      return "FIFOReduceNumFiles";
    case CompactionReason::kFIFOTtl:
      return "FIFOTtl";
    case CompactionReason::kManualCompaction:
      return "ManualCompaction";
    case CompactionReason::kFilesMarkedForCompaction:
      return "FilesMarkedForCompaction";
    case CompactionReason::kBottommostFiles:
      return "BottommostFiles";
    case CompactionReason::kTtl:
      return "Ttl";
    case CompactionReason::kFlush:
      return "Flush";
    case CompactionReason::kExternalSstIngestion:
      return "ExternalSstIngestion";
S
Sagar Vemuri 已提交
99 100
    case CompactionReason::kPeriodicCompaction:
      return "PeriodicCompaction";
101 102 103 104 105 106 107 108
    case CompactionReason::kNumOfReasons:
      // fall through
    default:
      assert(false);
      return "Invalid";
  }
}

109
// Maintains state for each sub-compaction
110
struct CompactionJob::SubcompactionState {
111
  const Compaction* compaction;
112
  std::unique_ptr<CompactionIterator> c_iter;
I
Igor Canadi 已提交
113

114 115 116 117 118
  // The boundaries of the key-range this compaction is interested in. No two
  // subcompactions may have overlapping key-ranges.
  // 'start' is inclusive, 'end' is exclusive, and nullptr means unbounded
  Slice *start, *end;

119
  // The return status of this subcompaction
120 121
  Status status;

122 123 124
  // The return IO Status of this subcompaction
  IOStatus io_status;

125
  // Files produced by this subcompaction
I
Igor Canadi 已提交
126
  struct Output {
127 128
    FileMetaData meta;
    bool finished;
129
    uint64_t paranoid_hash;
130
    std::shared_ptr<const TableProperties> table_properties;
I
Igor Canadi 已提交
131 132 133
  };

  // State kept for output being generated
134
  std::vector<Output> outputs;
135
  std::unique_ptr<WritableFileWriter> outfile;
I
Igor Canadi 已提交
136
  std::unique_ptr<TableBuilder> builder;
137

138
  Output* current_output() {
139
    if (outputs.empty()) {
140
      // This subcompaction's output could be empty if compaction was aborted
141 142
      // before this subcompaction had a chance to generate any output files.
      // When subcompactions are executed sequentially this is more likely and
143
      // will be particulalry likely for the later subcompactions to be empty.
144 145 146 147 148
      // Once they are run in parallel however it should be much rarer.
      return nullptr;
    } else {
      return &outputs.back();
    }
149
  }
I
Igor Canadi 已提交
150

151
  uint64_t current_output_file_size = 0;
152

153
  // State during the subcompaction
154 155
  uint64_t total_bytes = 0;
  uint64_t num_output_records = 0;
156
  CompactionJobStats compaction_job_stats;
157
  uint64_t approx_size = 0;
158
  // An index that used to speed up ShouldStopBefore().
159 160
  size_t grandparent_index = 0;
  // The number of bytes overlapping between the current output and
161
  // grandparent files used in ShouldStopBefore().
162
  uint64_t overlapped_bytes = 0;
163
  // A flag determine whether the key has been seen in ShouldStopBefore()
164
  bool seen_key = false;
165

166 167
  SubcompactionState(Compaction* c, Slice* _start, Slice* _end, uint64_t size)
      : compaction(c), start(_start), end(_end), approx_size(size) {
D
Dmitri Smirnov 已提交
168 169
    assert(compaction != nullptr);
  }
D
Dmitri Smirnov 已提交
170

171 172 173 174 175 176 177 178 179 180 181 182 183 184 185
  // Adds the key and value to the builder
  // If paranoid is true, adds the key-value to the paranoid hash
  void AddToBuilder(const Slice& key, const Slice& value, bool paranoid) {
    auto curr = current_output();
    assert(builder != nullptr);
    assert(curr != nullptr);
    if (paranoid) {
      // Generate a rolling 64-bit hash of the key and values
      curr->paranoid_hash = Hash64(key.data(), key.size(), curr->paranoid_hash);
      curr->paranoid_hash =
          Hash64(value.data(), value.size(), curr->paranoid_hash);
    }
    builder->Add(key, value);
  }

186 187
  // Returns true iff we should stop building the current output
  // before processing "internal_key".
188
  bool ShouldStopBefore(const Slice& internal_key, uint64_t curr_file_size) {
189 190 191 192 193 194 195 196 197 198 199 200 201 202 203
    const InternalKeyComparator* icmp =
        &compaction->column_family_data()->internal_comparator();
    const std::vector<FileMetaData*>& grandparents = compaction->grandparents();

    // Scan to find earliest grandparent file that contains key.
    while (grandparent_index < grandparents.size() &&
           icmp->Compare(internal_key,
                         grandparents[grandparent_index]->largest.Encode()) >
               0) {
      if (seen_key) {
        overlapped_bytes += grandparents[grandparent_index]->fd.GetFileSize();
      }
      assert(grandparent_index + 1 >= grandparents.size() ||
             icmp->Compare(
                 grandparents[grandparent_index]->largest.Encode(),
204
                 grandparents[grandparent_index + 1]->smallest.Encode()) <= 0);
205 206 207 208
      grandparent_index++;
    }
    seen_key = true;

209 210
    if (overlapped_bytes + curr_file_size >
        compaction->max_compaction_bytes()) {
211 212 213 214 215 216 217
      // Too much overlap for current output; start new output
      overlapped_bytes = 0;
      return true;
    }

    return false;
  }
218 219 220 221 222
};

// Maintains state for the entire compaction
struct CompactionJob::CompactionState {
  Compaction* const compaction;
I
Igor Canadi 已提交
223

224 225
  // REQUIRED: subcompaction states are stored in order of increasing
  // key-range
226
  std::vector<CompactionJob::SubcompactionState> sub_compact_states;
227 228 229 230
  Status status;

  uint64_t total_bytes;
  uint64_t num_output_records;
I
Igor Canadi 已提交
231

S
sdong 已提交
232 233 234 235
  explicit CompactionState(Compaction* c)
      : compaction(c),
        total_bytes(0),
        num_output_records(0) {}
I
Igor Canadi 已提交
236

237 238 239 240 241 242 243 244 245
  size_t NumOutputFiles() {
    size_t total = 0;
    for (auto& s : sub_compact_states) {
      total += s.outputs.size();
    }
    return total;
  }

  Slice SmallestUserKey() {
246 247 248 249
    for (const auto& sub_compact_state : sub_compact_states) {
      if (!sub_compact_state.outputs.empty() &&
          sub_compact_state.outputs[0].finished) {
        return sub_compact_state.outputs[0].meta.smallest.user_key();
250 251
      }
    }
252
    // If there is no finished output, return an empty slice.
253
    return Slice(nullptr, 0);
254 255 256
  }

  Slice LargestUserKey() {
257 258 259 260 261
    for (auto it = sub_compact_states.rbegin(); it < sub_compact_states.rend();
         ++it) {
      if (!it->outputs.empty() && it->current_output()->finished) {
        assert(it->current_output() != nullptr);
        return it->current_output()->meta.largest.user_key();
262 263
      }
    }
264
    // If there is no finished output, return an empty slice.
265
    return Slice(nullptr, 0);
266
  }
I
Igor Canadi 已提交
267 268
};

269
void CompactionJob::AggregateStatistics() {
270
  for (SubcompactionState& sc : compact_->sub_compact_states) {
271 272
    compact_->total_bytes += sc.total_bytes;
    compact_->num_output_records += sc.num_output_records;
273 274 275
  }
  if (compaction_job_stats_) {
    for (SubcompactionState& sc : compact_->sub_compact_states) {
276 277 278 279 280
      compaction_job_stats_->Add(sc.compaction_job_stats);
    }
  }
}

I
Igor Canadi 已提交
281
CompactionJob::CompactionJob(
282
    int job_id, Compaction* compaction, const ImmutableDBOptions& db_options,
283
    const FileOptions& file_options, VersionSet* versions,
284
    const std::atomic<bool>* shutting_down,
285
    const SequenceNumber preserve_deletes_seqnum, LogBuffer* log_buffer,
286
    FSDirectory* db_directory, FSDirectory* output_directory, Statistics* stats,
287
    InstrumentedMutex* db_mutex, ErrorHandler* db_error_handler,
I
Igor Canadi 已提交
288
    std::vector<SequenceNumber> existing_snapshots,
289
    SequenceNumber earliest_write_conflict_snapshot,
Y
Yi Wu 已提交
290 291
    const SnapshotChecker* snapshot_checker, std::shared_ptr<Cache> table_cache,
    EventLogger* event_logger, bool paranoid_file_checks, bool measure_io_stats,
292
    const std::string& dbname, CompactionJobStats* compaction_job_stats,
293
    Env::Priority thread_pri, const std::shared_ptr<IOTracer>& io_tracer,
294
    const std::atomic<int>* manual_compaction_paused, const std::string& db_id,
295
    const std::string& db_session_id)
296 297
    : job_id_(job_id),
      compact_(new CompactionState(compaction)),
298
      compaction_job_stats_(compaction_job_stats),
299
      compaction_stats_(compaction->compaction_reason(), 1),
300
      dbname_(dbname),
301 302
      db_id_(db_id),
      db_session_id_(db_session_id),
I
Igor Canadi 已提交
303
      db_options_(db_options),
304
      file_options_(file_options),
I
Igor Canadi 已提交
305
      env_(db_options.env),
306
      io_tracer_(io_tracer),
307
      fs_(db_options.fs, io_tracer),
308 309
      file_options_for_read_(
          fs_->OptimizeForCompactionTableRead(file_options, db_options_)),
I
Igor Canadi 已提交
310 311
      versions_(versions),
      shutting_down_(shutting_down),
312
      manual_compaction_paused_(manual_compaction_paused),
313
      preserve_deletes_seqnum_(preserve_deletes_seqnum),
I
Igor Canadi 已提交
314 315
      log_buffer_(log_buffer),
      db_directory_(db_directory),
316
      output_directory_(output_directory),
I
Igor Canadi 已提交
317
      stats_(stats),
318
      db_mutex_(db_mutex),
319
      db_error_handler_(db_error_handler),
I
Igor Canadi 已提交
320
      existing_snapshots_(std::move(existing_snapshots)),
321
      earliest_write_conflict_snapshot_(earliest_write_conflict_snapshot),
Y
Yi Wu 已提交
322
      snapshot_checker_(snapshot_checker),
I
Igor Canadi 已提交
323
      table_cache_(std::move(table_cache)),
I
Igor Canadi 已提交
324
      event_logger_(event_logger),
325
      bottommost_level_(false),
326
      paranoid_file_checks_(paranoid_file_checks),
S
Stream  
Shaohua Li 已提交
327
      measure_io_stats_(measure_io_stats),
328 329
      write_hint_(Env::WLTH_NOT_SET),
      thread_pri_(thread_pri) {
330
  assert(log_buffer_ != nullptr);
331 332
  const auto* cfd = compact_->compaction->column_family_data();
  ThreadStatusUtil::SetColumnFamily(cfd, cfd->ioptions()->env,
Y
Yi Wu 已提交
333
                                    db_options_.enable_thread_tracking);
334
  ThreadStatusUtil::SetThreadOperation(ThreadStatus::OP_COMPACTION);
335
  ReportStartedCompaction(compaction);
336 337 338 339 340 341
}

CompactionJob::~CompactionJob() {
  assert(compact_ == nullptr);
  ThreadStatusUtil::ResetThreadStatus();
}
I
Igor Canadi 已提交
342

343
void CompactionJob::ReportStartedCompaction(Compaction* compaction) {
344 345
  const auto* cfd = compact_->compaction->column_family_data();
  ThreadStatusUtil::SetColumnFamily(cfd, cfd->ioptions()->env,
Y
Yi Wu 已提交
346
                                    db_options_.enable_thread_tracking);
347

348 349
  ThreadStatusUtil::SetThreadOperationProperty(ThreadStatus::COMPACTION_JOB_ID,
                                               job_id_);
350 351 352 353 354 355

  ThreadStatusUtil::SetThreadOperationProperty(
      ThreadStatus::COMPACTION_INPUT_OUTPUT_LEVEL,
      (static_cast<uint64_t>(compact_->compaction->start_level()) << 32) +
          compact_->compaction->output_level());

356 357 358
  // In the current design, a CompactionJob is always created
  // for non-trivial compaction.
  assert(compaction->IsTrivialMove() == false ||
359
         compaction->is_manual_compaction() == true);
360

361 362
  ThreadStatusUtil::SetThreadOperationProperty(
      ThreadStatus::COMPACTION_PROP_FLAGS,
S
sdong 已提交
363
      compaction->is_manual_compaction() +
364
          (compaction->deletion_compaction() << 1));
365 366 367 368 369 370 371 372 373 374 375 376 377 378

  ThreadStatusUtil::SetThreadOperationProperty(
      ThreadStatus::COMPACTION_TOTAL_INPUT_BYTES,
      compaction->CalculateTotalInputSize());

  IOSTATS_RESET(bytes_written);
  IOSTATS_RESET(bytes_read);
  ThreadStatusUtil::SetThreadOperationProperty(
      ThreadStatus::COMPACTION_BYTES_WRITTEN, 0);
  ThreadStatusUtil::SetThreadOperationProperty(
      ThreadStatus::COMPACTION_BYTES_READ, 0);

  // Set the thread operation after operation properties
  // to ensure GetThreadList() can always show them all together.
379
  ThreadStatusUtil::SetThreadOperation(ThreadStatus::OP_COMPACTION);
380 381 382

  if (compaction_job_stats_) {
    compaction_job_stats_->is_manual_compaction =
S
sdong 已提交
383
        compaction->is_manual_compaction();
384
  }
385 386
}

I
Igor Canadi 已提交
387
void CompactionJob::Prepare() {
388 389
  AutoThreadOperationStageUpdater stage_updater(
      ThreadStatus::STAGE_COMPACTION_PREPARE);
I
Igor Canadi 已提交
390 391

  // Generate file_levels_ for compaction berfore making Iterator
392 393
  auto* c = compact_->compaction;
  assert(c->column_family_data() != nullptr);
394 395
  assert(c->column_family_data()->current()->storage_info()->NumLevelFiles(
             compact_->compaction->level()) > 0);
I
Igor Canadi 已提交
396

397 398
  write_hint_ =
      c->column_family_data()->CalculateSSTWriteHint(c->output_level());
399
  bottommost_level_ = c->bottommost_level();
400

401
  if (c->ShouldFormSubcompactions()) {
S
Siying Dong 已提交
402 403 404 405
    {
      StopWatch sw(env_, stats_, SUBCOMPACTION_SETUP_TIME);
      GenSubcompactionBoundaries();
    }
406 407 408 409 410 411 412
    assert(sizes_.size() == boundaries_.size() + 1);

    for (size_t i = 0; i <= boundaries_.size(); i++) {
      Slice* start = i == 0 ? nullptr : &boundaries_[i - 1];
      Slice* end = i == boundaries_.size() ? nullptr : &boundaries_[i];
      compact_->sub_compact_states.emplace_back(c, start, end, sizes_[i]);
    }
S
Siying Dong 已提交
413 414
    RecordInHistogram(stats_, NUM_SUBCOMPACTIONS_SCHEDULED,
                      compact_->sub_compact_states.size());
415
  } else {
416 417 418 419 420
    constexpr Slice* start = nullptr;
    constexpr Slice* end = nullptr;
    constexpr uint64_t size = 0;

    compact_->sub_compact_states.emplace_back(c, start, end, size);
421
  }
422 423
}

424 425 426
struct RangeWithSize {
  Range range;
  uint64_t size;
427

428 429 430 431 432 433 434
  RangeWithSize(const Slice& a, const Slice& b, uint64_t s = 0)
      : range(a, b), size(s) {}
};

void CompactionJob::GenSubcompactionBoundaries() {
  auto* c = compact_->compaction;
  auto* cfd = c->column_family_data();
435 436
  const Comparator* cfd_comparator = cfd->user_comparator();
  std::vector<Slice> bounds;
437 438 439 440
  int start_lvl = c->start_level();
  int out_lvl = c->output_level();

  // Add the starting and/or ending key of certain input files as a potential
441
  // boundary
442 443 444 445 446 447 448
  for (size_t lvl_idx = 0; lvl_idx < c->num_input_levels(); lvl_idx++) {
    int lvl = c->level(lvl_idx);
    if (lvl >= start_lvl && lvl <= out_lvl) {
      const LevelFilesBrief* flevel = c->input_levels(lvl_idx);
      size_t num_files = flevel->num_files;

      if (num_files == 0) {
449
        continue;
450
      }
A
Ari Ekmekji 已提交
451

452 453 454 455
      if (lvl == 0) {
        // For level 0 add the starting and ending key of each file since the
        // files may have greatly differing key ranges (not range-partitioned)
        for (size_t i = 0; i < num_files; i++) {
456 457
          bounds.emplace_back(flevel->files[i].smallest_key);
          bounds.emplace_back(flevel->files[i].largest_key);
458 459 460 461
        }
      } else {
        // For all other levels add the smallest/largest key in the level to
        // encompass the range covered by that level
462 463
        bounds.emplace_back(flevel->files[0].smallest_key);
        bounds.emplace_back(flevel->files[num_files - 1].largest_key);
464 465 466 467 468 469
        if (lvl == out_lvl) {
          // For the last level include the starting keys of all files since
          // the last level is the largest and probably has the widest key
          // range. Since it's range partitioned, the ending key of one file
          // and the starting key of the next are very close (or identical).
          for (size_t i = 1; i < num_files; i++) {
470
            bounds.emplace_back(flevel->files[i].smallest_key);
A
Ari Ekmekji 已提交
471
          }
472 473 474 475
        }
      }
    }
  }
476

477
  std::sort(bounds.begin(), bounds.end(),
478 479 480 481
            [cfd_comparator](const Slice& a, const Slice& b) -> bool {
              return cfd_comparator->Compare(ExtractUserKey(a),
                                             ExtractUserKey(b)) < 0;
            });
482
  // Remove duplicated entries from bounds
483 484 485 486 487 488 489
  bounds.erase(
      std::unique(bounds.begin(), bounds.end(),
                  [cfd_comparator](const Slice& a, const Slice& b) -> bool {
                    return cfd_comparator->Compare(ExtractUserKey(a),
                                                   ExtractUserKey(b)) == 0;
                  }),
      bounds.end());
490

491 492 493 494
  // Combine consecutive pairs of boundaries into ranges with an approximate
  // size of data covered by keys in that range
  uint64_t sum = 0;
  std::vector<RangeWithSize> ranges;
495 496 497 498
  // Get input version from CompactionState since it's already referenced
  // earlier in SetInputVersioCompaction::SetInputVersion and will not change
  // when db_mutex_ is released below
  auto* v = compact_->compaction->input_version();
499 500
  for (auto it = bounds.begin();;) {
    const Slice a = *it;
501
    ++it;
502 503 504 505 506 507

    if (it == bounds.end()) {
      break;
    }

    const Slice b = *it;
508 509 510 511 512

    // ApproximateSize could potentially create table reader iterator to seek
    // to the index block and may incur I/O cost in the process. Unlock db
    // mutex to reduce contention
    db_mutex_->Unlock();
513 514
    uint64_t size = versions_->ApproximateSize(SizeApproximationOptions(), v, a,
                                               b, start_lvl, out_lvl + 1,
515
                                               TableReaderCaller::kCompaction);
516
    db_mutex_->Lock();
517 518 519 520 521 522
    ranges.emplace_back(a, b, size);
    sum += size;
  }

  // Group the ranges into subcompactions
  const double min_file_fill_percent = 4.0 / 5;
523 524 525 526 527 528
  int base_level = v->storage_info()->base_level();
  uint64_t max_output_files = static_cast<uint64_t>(std::ceil(
      sum / min_file_fill_percent /
      MaxFileSizeForLevel(*(c->mutable_cf_options()), out_lvl,
          c->immutable_cf_options()->compaction_style, base_level,
          c->immutable_cf_options()->level_compaction_dynamic_level_bytes)));
529 530
  uint64_t subcompactions =
      std::min({static_cast<uint64_t>(ranges.size()),
531
                static_cast<uint64_t>(c->max_subcompactions()),
532 533 534
                max_output_files});

  if (subcompactions > 1) {
535
    double mean = sum * 1.0 / subcompactions;
536 537 538
    // Greedily add ranges to the subcompaction until the sum of the ranges'
    // sizes becomes >= the expected mean size of a subcompaction
    sum = 0;
539
    for (size_t i = 0; i + 1 < ranges.size(); i++) {
540
      sum += ranges[i].size;
541 542 543
      if (subcompactions == 1) {
        // If there's only one left to schedule then it goes to the end so no
        // need to put an end boundary
544
        continue;
545 546 547 548 549 550 551 552 553 554 555 556
      }
      if (sum >= mean) {
        boundaries_.emplace_back(ExtractUserKey(ranges[i].range.limit));
        sizes_.emplace_back(sum);
        subcompactions--;
        sum = 0;
      }
    }
    sizes_.emplace_back(sum + ranges.back().size);
  } else {
    // Only one range so its size is the total sum of sizes computed above
    sizes_.emplace_back(sum);
557
  }
I
Igor Canadi 已提交
558 559 560
}

Status CompactionJob::Run() {
561 562
  AutoThreadOperationStageUpdater stage_updater(
      ThreadStatus::STAGE_COMPACTION_RUN);
563
  TEST_SYNC_POINT("CompactionJob::Run():Start");
I
Igor Canadi 已提交
564
  log_buffer_->FlushBufferToLog();
565
  LogCompaction();
566

567 568
  const size_t num_threads = compact_->sub_compact_states.size();
  assert(num_threads > 0);
569
  const uint64_t start_micros = env_->NowMicros();
570 571

  // Launch a thread for each of subcompactions 1...num_threads-1
D
Dmitri Smirnov 已提交
572
  std::vector<port::Thread> thread_pool;
573 574 575 576
  thread_pool.reserve(num_threads - 1);
  for (size_t i = 1; i < compact_->sub_compact_states.size(); i++) {
    thread_pool.emplace_back(&CompactionJob::ProcessKeyValueCompaction, this,
                             &compact_->sub_compact_states[i]);
577
  }
578 579 580 581 582 583 584 585 586 587

  // Always schedule the first subcompaction (whether or not there are also
  // others) in the current thread to be efficient with resources
  ProcessKeyValueCompaction(&compact_->sub_compact_states[0]);

  // Wait for all other threads (if there are any) to finish execution
  for (auto& thread : thread_pool) {
    thread.join();
  }

588
  compaction_stats_.micros = env_->NowMicros() - start_micros;
589 590 591 592 593 594
  compaction_stats_.cpu_micros = 0;
  for (size_t i = 0; i < compact_->sub_compact_states.size(); i++) {
    compaction_stats_.cpu_micros +=
        compact_->sub_compact_states[i].compaction_job_stats.cpu_micros;
  }

S
Siying Dong 已提交
595 596 597
  RecordTimeToHistogram(stats_, COMPACTION_TIME, compaction_stats_.micros);
  RecordTimeToHistogram(stats_, COMPACTION_CPU_TIME,
                        compaction_stats_.cpu_micros);
598

599 600
  TEST_SYNC_POINT("CompactionJob::Run:BeforeVerify");

601
  // Check if any thread encountered an error during execution
602
  Status status;
603
  IOStatus io_s;
604 605 606
  for (const auto& state : compact_->sub_compact_states) {
    if (!state.status.ok()) {
      status = state.status;
607
      io_s = state.io_status;
608 609 610
      break;
    }
  }
611 612 613
  if (io_status_.ok()) {
    io_status_ = io_s;
  }
614
  if (status.ok() && output_directory_) {
615 616
    io_s = output_directory_->Fsync(IOOptions(), nullptr);
  }
617
  if (io_status_.ok()) {
618
    io_status_ = io_s;
619 620
  }
  if (status.ok()) {
621
    status = io_s;
622
  }
623 624
  if (status.ok()) {
    thread_pool.clear();
625
    std::vector<const CompactionJob::SubcompactionState::Output*> files_output;
626 627
    for (const auto& state : compact_->sub_compact_states) {
      for (const auto& output : state.outputs) {
628
        files_output.emplace_back(&output);
629 630 631 632 633
      }
    }
    ColumnFamilyData* cfd = compact_->compaction->column_family_data();
    auto prefix_extractor =
        compact_->compaction->mutable_cf_options()->prefix_extractor.get();
634
    std::atomic<size_t> next_file_idx(0);
635 636
    auto verify_table = [&](Status& output_status) {
      while (true) {
637 638
        size_t file_idx = next_file_idx.fetch_add(1);
        if (file_idx >= files_output.size()) {
639 640 641 642 643 644 645 646
          break;
        }
        // Verify that the table is usable
        // We set for_compaction to false and don't OptimizeForCompactionTableRead
        // here because this is a special case after we finish the table building
        // No matter whether use_direct_io_for_flush_and_compaction is true,
        // we will regard this verification as user reads since the goal is
        // to cache it here for further user reads
647
        ReadOptions read_options;
648
        InternalIterator* iter = cfd->table_cache()->NewIterator(
649
            read_options, file_options_, cfd->internal_comparator(),
650 651
            files_output[file_idx]->meta, /*range_del_agg=*/nullptr,
            prefix_extractor,
652
            /*table_reader_ptr=*/nullptr,
653 654
            cfd->internal_stats()->GetFileReadHist(
                compact_->compaction->output_level()),
655 656
            TableReaderCaller::kCompactionRefill, /*arena=*/nullptr,
            /*skip_filters=*/false, compact_->compaction->output_level(),
657 658
            MaxFileSizeForL0MetaPin(
                *compact_->compaction->mutable_cf_options()),
659
            /*smallest_compaction_key=*/nullptr,
660 661
            /*largest_compaction_key=*/nullptr,
            /*allow_unprepared_value=*/false);
662 663 664
        auto s = iter->status();

        if (s.ok() && paranoid_file_checks_) {
665 666 667 668 669 670
          uint64_t hash = 0;
          for (iter->SeekToFirst(); iter->Valid(); iter->Next()) {
            // Generate a rolling 64-bit hash of the key and values, using the
            hash = Hash64(iter->key().data(), iter->key().size(), hash);
            hash = Hash64(iter->value().data(), iter->value().size(), hash);
          }
671
          s = iter->status();
672
          if (s.ok() && hash != files_output[file_idx]->paranoid_hash) {
673
            s = Status::Corruption("Paranoid checksums do not match");
674
          }
675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700
        }

        delete iter;

        if (!s.ok()) {
          output_status = s;
          break;
        }
      }
    };
    for (size_t i = 1; i < compact_->sub_compact_states.size(); i++) {
      thread_pool.emplace_back(verify_table,
                               std::ref(compact_->sub_compact_states[i].status));
    }
    verify_table(compact_->sub_compact_states[0].status);
    for (auto& thread : thread_pool) {
      thread.join();
    }
    for (const auto& state : compact_->sub_compact_states) {
      if (!state.status.ok()) {
        status = state.status;
        break;
      }
    }
  }

701 702 703
  TablePropertiesCollection tp;
  for (const auto& state : compact_->sub_compact_states) {
    for (const auto& output : state.outputs) {
704 705 706
      auto fn =
          TableFileName(state.compaction->immutable_cf_options()->cf_paths,
                        output.meta.fd.GetNumber(), output.meta.fd.GetPathId());
707 708 709 710 711
      tp[fn] = output.table_properties;
    }
  }
  compact_->compaction->SetOutputTableProperties(std::move(tp));

712 713
  // Finish up all book-keeping to unify the subcompaction results
  AggregateStatistics();
714 715 716 717 718
  UpdateCompactionStats();
  RecordCompactionIOStats();
  LogFlush(db_options_.info_log);
  TEST_SYNC_POINT("CompactionJob::Run():End");

719
  compact_->status = status;
720 721 722
  return status;
}

723
Status CompactionJob::Install(const MutableCFOptions& mutable_cf_options) {
724 725
  AutoThreadOperationStageUpdater stage_updater(
      ThreadStatus::STAGE_COMPACTION_INSTALL);
726
  db_mutex_->AssertHeld();
727
  Status status = compact_->status;
I
Igor Canadi 已提交
728 729
  ColumnFamilyData* cfd = compact_->compaction->column_family_data();
  cfd->internal_stats()->AddCompactionStats(
730
      compact_->compaction->output_level(), thread_pri_, compaction_stats_);
I
Igor Canadi 已提交
731

732
  if (status.ok()) {
733
    status = InstallCompactionResults(mutable_cf_options);
I
Igor Canadi 已提交
734
  }
735 736 737
  if (!versions_->io_status().ok()) {
    io_status_ = versions_->io_status();
  }
I
Igor Canadi 已提交
738
  VersionStorageInfo::LevelSummaryStorage tmp;
739
  auto vstorage = cfd->current()->storage_info();
I
Igor Canadi 已提交
740
  const auto& stats = compaction_stats_;
741 742 743

  double read_write_amp = 0.0;
  double write_amp = 0.0;
744 745 746
  double bytes_read_per_sec = 0;
  double bytes_written_per_sec = 0;

747 748 749 750 751 752 753
  if (stats.bytes_read_non_output_levels > 0) {
    read_write_amp = (stats.bytes_written + stats.bytes_read_output_level +
                      stats.bytes_read_non_output_levels) /
                     static_cast<double>(stats.bytes_read_non_output_levels);
    write_amp = stats.bytes_written /
                static_cast<double>(stats.bytes_read_non_output_levels);
  }
754 755 756 757 758 759 760 761
  if (stats.micros > 0) {
    bytes_read_per_sec =
        (stats.bytes_read_non_output_levels + stats.bytes_read_output_level) /
        static_cast<double>(stats.micros);
    bytes_written_per_sec =
        stats.bytes_written / static_cast<double>(stats.micros);
  }

762
  ROCKS_LOG_BUFFER(
763 764 765 766
      log_buffer_,
      "[%s] compacted to: %s, MB/sec: %.1f rd, %.1f wr, level %d, "
      "files in(%d, %d) out(%d) "
      "MB in(%.1f, %.1f) out(%.1f), read-write-amplify(%.1f) "
S
Siying Dong 已提交
767 768
      "write-amplify(%.1f) %s, records in: %" PRIu64
      ", records dropped: %" PRIu64 " output_compression: %s\n",
769 770
      cfd->GetName().c_str(), vstorage->LevelSummary(&tmp), bytes_read_per_sec,
      bytes_written_per_sec, compact_->compaction->output_level(),
771
      stats.num_input_files_in_non_output_levels,
772
      stats.num_input_files_in_output_level, stats.num_output_files,
773 774
      stats.bytes_read_non_output_levels / 1048576.0,
      stats.bytes_read_output_level / 1048576.0,
775
      stats.bytes_written / 1048576.0, read_write_amp, write_amp,
776
      status.ToString().c_str(), stats.num_input_records,
777 778 779
      stats.num_dropped_records,
      CompressionTypeToString(compact_->compaction->output_compression())
          .c_str());
I
Igor Canadi 已提交
780

781 782
  UpdateCompactionJobStats(stats);

783
  auto stream = event_logger_->LogToBuffer(log_buffer_);
784 785
  stream << "job" << job_id_ << "event"
         << "compaction_finished"
786 787 788 789 790 791
         << "compaction_time_micros" << stats.micros
         << "compaction_time_cpu_micros" << stats.cpu_micros << "output_level"
         << compact_->compaction->output_level() << "num_output_files"
         << compact_->NumOutputFiles() << "total_output_size"
         << compact_->total_bytes << "num_input_records"
         << stats.num_input_records << "num_output_records"
792 793 794
         << compact_->num_output_records << "num_subcompactions"
         << compact_->sub_compact_states.size() << "output_compression"
         << CompressionTypeToString(compact_->compaction->output_compression());
795

796 797 798 799 800 801 802
  if (compaction_job_stats_ != nullptr) {
    stream << "num_single_delete_mismatches"
           << compaction_job_stats_->num_single_del_mismatch;
    stream << "num_single_delete_fallthrough"
           << compaction_job_stats_->num_single_del_fallthru;
  }

803 804 805 806 807 808 809 810 811
  if (measure_io_stats_ && compaction_job_stats_ != nullptr) {
    stream << "file_write_nanos" << compaction_job_stats_->file_write_nanos;
    stream << "file_range_sync_nanos"
           << compaction_job_stats_->file_range_sync_nanos;
    stream << "file_fsync_nanos" << compaction_job_stats_->file_fsync_nanos;
    stream << "file_prepare_write_nanos"
           << compaction_job_stats_->file_prepare_write_nanos;
  }

812 813 814 815 816 817 818
  stream << "lsm_state";
  stream.StartArray();
  for (int level = 0; level < vstorage->num_levels(); ++level) {
    stream << vstorage->NumLevelFiles(level);
  }
  stream.EndArray();

819 820
  CleanupCompaction();
  return status;
I
Igor Canadi 已提交
821 822
}

823
void CompactionJob::ProcessKeyValueCompaction(SubcompactionState* sub_compact) {
824
  assert(sub_compact != nullptr);
825 826 827

  uint64_t prev_cpu_micros = env_->NowCPUNanos() / 1000;

828
  ColumnFamilyData* cfd = sub_compact->compaction->column_family_data();
829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846

  // Create compaction filter and fail the compaction if
  // IgnoreSnapshots() = false because it is not supported anymore
  const CompactionFilter* compaction_filter =
      cfd->ioptions()->compaction_filter;
  std::unique_ptr<CompactionFilter> compaction_filter_from_factory = nullptr;
  if (compaction_filter == nullptr) {
    compaction_filter_from_factory =
        sub_compact->compaction->CreateCompactionFilter();
    compaction_filter = compaction_filter_from_factory.get();
  }
  if (compaction_filter != nullptr && !compaction_filter->IgnoreSnapshots()) {
    sub_compact->status = Status::NotSupported(
        "CompactionFilter::IgnoreSnapshots() = false is not supported "
        "anymore.");
    return;
  }

847 848
  CompactionRangeDelAggregator range_del_agg(&cfd->internal_comparator(),
                                             existing_snapshots_);
849 850 851 852 853 854 855 856
  ReadOptions read_options;
  read_options.verify_checksums = true;
  read_options.fill_cache = false;
  // Compaction iterators shouldn't be confined to a single prefix.
  // Compactions use Seek() for
  // (a) concurrent compactions,
  // (b) CompactionFilter::Decision::kRemoveAndSkipUntil.
  read_options.total_order_seek = true;
857 858 859

  // Although the v2 aggregator is what the level iterator(s) know about,
  // the AddTombstones calls will be propagated down to the v1 aggregator.
860 861 862
  std::unique_ptr<InternalIterator> input(
      versions_->MakeInputIterator(read_options, sub_compact->compaction,
                                   &range_del_agg, file_options_for_read_));
863

864 865
  AutoThreadOperationStageUpdater stage_updater(
      ThreadStatus::STAGE_COMPACTION_PROCESS_KV);
866 867 868

  // I/O measurement variables
  PerfLevel prev_perf_level = PerfLevel::kEnableTime;
869
  const uint64_t kRecordStatsEvery = 1000;
870 871 872 873
  uint64_t prev_write_nanos = 0;
  uint64_t prev_fsync_nanos = 0;
  uint64_t prev_range_sync_nanos = 0;
  uint64_t prev_prepare_write_nanos = 0;
874 875
  uint64_t prev_cpu_write_nanos = 0;
  uint64_t prev_cpu_read_nanos = 0;
876 877
  if (measure_io_stats_) {
    prev_perf_level = GetPerfLevel();
878
    SetPerfLevel(PerfLevel::kEnableTimeAndCPUTimeExceptForMutex);
I
Igor Canadi 已提交
879 880 881 882
    prev_write_nanos = IOSTATS(write_nanos);
    prev_fsync_nanos = IOSTATS(fsync_nanos);
    prev_range_sync_nanos = IOSTATS(range_sync_nanos);
    prev_prepare_write_nanos = IOSTATS(prepare_write_nanos);
883 884
    prev_cpu_write_nanos = IOSTATS(cpu_write_nanos);
    prev_cpu_read_nanos = IOSTATS(cpu_read_nanos);
885 886
  }

I
Igor Canadi 已提交
887 888 889 890 891
  MergeHelper merge(
      env_, cfd->user_comparator(), cfd->ioptions()->merge_operator,
      compaction_filter, db_options_.info_log.get(),
      false /* internal key corruption is expected */,
      existing_snapshots_.empty() ? 0 : existing_snapshots_.back(),
892
      snapshot_checker_, compact_->compaction->level(),
893
      db_options_.statistics.get());
I
Igor Canadi 已提交
894

895
  TEST_SYNC_POINT("CompactionJob::Run():Inprogress");
896 897 898
  TEST_SYNC_POINT_CALLBACK(
      "CompactionJob::Run():PausingManualCompaction:1",
      reinterpret_cast<void*>(
899
          const_cast<std::atomic<int>*>(manual_compaction_paused_)));
900

901 902
  Slice* start = sub_compact->start;
  Slice* end = sub_compact->end;
903 904 905
  if (start != nullptr) {
    IterKey start_iter;
    start_iter.SetInternalKey(*start, kMaxSequenceNumber, kValueTypeForSeek);
906
    input->Seek(start_iter.GetInternalKey());
907 908 909 910
  } else {
    input->SeekToFirst();
  }

911 912 913
  Status status;
  sub_compact->c_iter.reset(new CompactionIterator(
      input.get(), cfd->user_comparator(), &merge, versions_->LastSequence(),
Y
Yi Wu 已提交
914
      &existing_snapshots_, earliest_write_conflict_snapshot_,
915 916
      snapshot_checker_, env_, ShouldReportDetailedTime(env_, stats_),
      /*expect_valid_internal_key=*/true, &range_del_agg,
917 918 919 920
      /* blob_file_builder */ nullptr, db_options_.allow_data_in_errors,
      sub_compact->compaction, compaction_filter, shutting_down_,
      preserve_deletes_seqnum_, manual_compaction_paused_,
      db_options_.info_log));
921 922
  auto c_iter = sub_compact->c_iter.get();
  c_iter->SeekToFirst();
923
  if (c_iter->Valid() && sub_compact->compaction->output_level() != 0) {
924 925 926
    // ShouldStopBefore() maintains state based on keys processed so far. The
    // compaction loop always calls it on the "next" key, thus won't tell it the
    // first key. So we do that here.
927 928
    sub_compact->ShouldStopBefore(c_iter->key(),
                                  sub_compact->current_output_file_size);
929
  }
930
  const auto& c_iter_stats = c_iter->iter_stats();
931

932 933 934 935 936 937
  std::unique_ptr<SstPartitioner> partitioner =
      sub_compact->compaction->output_level() == 0
          ? nullptr
          : sub_compact->compaction->CreateSstPartitioner();
  std::string last_key_for_partitioner;

938
  while (status.ok() && !cfd->IsDropped() && c_iter->Valid()) {
939 940 941 942
    // Invariant: c_iter.status() is guaranteed to be OK if c_iter->Valid()
    // returns true.
    const Slice& key = c_iter->key();
    const Slice& value = c_iter->value();
943

944 945
    // If an end key (exclusive) is specified, check if the current key is
    // >= than it and exit if it is because the iterator is out of its range
946
    if (end != nullptr &&
947
        cfd->user_comparator()->Compare(c_iter->user_key(), *end) >= 0) {
948
      break;
I
Igor Canadi 已提交
949
    }
950 951 952 953 954
    if (c_iter_stats.num_input_records % kRecordStatsEvery ==
        kRecordStatsEvery - 1) {
      RecordDroppedKeys(c_iter_stats, &sub_compact->compaction_job_stats);
      c_iter->ResetRecordCounts();
      RecordCompactionIOStats();
I
Igor Canadi 已提交
955 956
    }

957 958 959 960
    // Open output file if necessary
    if (sub_compact->builder == nullptr) {
      status = OpenCompactionOutputFile(sub_compact);
      if (!status.ok()) {
961 962
        break;
      }
963
    }
964 965
    sub_compact->AddToBuilder(key, value, paranoid_file_checks_);

966 967
    sub_compact->current_output_file_size =
        sub_compact->builder->EstimatedFileSize();
968
    const ParsedInternalKey& ikey = c_iter->ikey();
969
    sub_compact->current_output()->meta.UpdateBoundaries(
970
        key, value, ikey.sequence, ikey.type);
971 972
    sub_compact->num_output_records++;

973 974 975 976
    // Close output file if it is big enough. Two possibilities determine it's
    // time to close it: (1) the current key should be this file's last key, (2)
    // the next key should not be in this file.
    //
977 978 979 980
    // TODO(aekmekji): determine if file should be closed earlier than this
    // during subcompactions (i.e. if output size, estimated by input size, is
    // going to be 1.2MB and max_output_file_size = 1MB, prefer to have 0.6MB
    // and 0.6MB instead of 1MB and 0.2MB)
981
    bool output_file_ended = false;
982 983 984
    if (sub_compact->compaction->output_level() != 0 &&
        sub_compact->current_output_file_size >=
            sub_compact->compaction->max_output_file_size()) {
985 986 987 988
      // (1) this key terminates the file. For historical reasons, the iterator
      // status before advancing will be given to FinishCompactionOutputFile().
      output_file_ended = true;
    }
989 990 991
    TEST_SYNC_POINT_CALLBACK(
        "CompactionJob::Run():PausingManualCompaction:2",
        reinterpret_cast<void*>(
992
            const_cast<std::atomic<int>*>(manual_compaction_paused_)));
993 994 995 996
    if (partitioner.get()) {
      last_key_for_partitioner.assign(c_iter->user_key().data_,
                                      c_iter->user_key().size_);
    }
997
    c_iter->Next();
998 999 1000
    if (c_iter->status().IsManualCompactionPaused()) {
      break;
    }
1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014
    if (!output_file_ended && c_iter->Valid()) {
      if (((partitioner.get() &&
            partitioner->ShouldPartition(PartitionerRequest(
                last_key_for_partitioner, c_iter->user_key(),
                sub_compact->current_output_file_size)) == kRequired) ||
           (sub_compact->compaction->output_level() != 0 &&
            sub_compact->ShouldStopBefore(
                c_iter->key(), sub_compact->current_output_file_size))) &&
          sub_compact->builder != nullptr) {
        // (2) this key belongs to the next file. For historical reasons, the
        // iterator status after advancing will be given to
        // FinishCompactionOutputFile().
        output_file_ended = true;
      }
1015 1016
    }
    if (output_file_ended) {
1017 1018 1019 1020
      const Slice* next_key = nullptr;
      if (c_iter->Valid()) {
        next_key = &c_iter->key();
      }
A
Andrew Kryczka 已提交
1021
      CompactionIterationStats range_del_out_stats;
1022 1023 1024
      status = FinishCompactionOutputFile(input->status(), sub_compact,
                                          &range_del_agg, &range_del_out_stats,
                                          next_key);
A
Andrew Kryczka 已提交
1025 1026
      RecordDroppedKeys(range_del_out_stats,
                        &sub_compact->compaction_job_stats);
I
Igor Canadi 已提交
1027 1028
    }
  }
1029

1030 1031 1032 1033
  sub_compact->compaction_job_stats.num_input_deletion_records =
      c_iter_stats.num_input_deletion_records;
  sub_compact->compaction_job_stats.num_corrupt_keys =
      c_iter_stats.num_input_corrupt_records;
1034 1035 1036 1037
  sub_compact->compaction_job_stats.num_single_del_fallthru =
      c_iter_stats.num_single_del_fallthru;
  sub_compact->compaction_job_stats.num_single_del_mismatch =
      c_iter_stats.num_single_del_mismatch;
1038 1039 1040 1041 1042 1043 1044 1045
  sub_compact->compaction_job_stats.total_input_raw_key_bytes +=
      c_iter_stats.total_input_raw_key_bytes;
  sub_compact->compaction_job_stats.total_input_raw_value_bytes +=
      c_iter_stats.total_input_raw_value_bytes;

  RecordTick(stats_, FILTER_OPERATION_TOTAL_TIME,
             c_iter_stats.total_filter_time);
  RecordDroppedKeys(c_iter_stats, &sub_compact->compaction_job_stats);
I
Igor Canadi 已提交
1046 1047
  RecordCompactionIOStats();

1048 1049 1050 1051 1052 1053 1054
  if (status.ok() && cfd->IsDropped()) {
    status =
        Status::ColumnFamilyDropped("Column family dropped during compaction");
  }
  if ((status.ok() || status.IsColumnFamilyDropped()) &&
      shutting_down_->load(std::memory_order_relaxed)) {
    status = Status::ShutdownInProgress("Database shutdown");
1055
  }
1056 1057
  if ((status.ok() || status.IsColumnFamilyDropped()) &&
      (manual_compaction_paused_ &&
1058
       manual_compaction_paused_->load(std::memory_order_relaxed) > 0)) {
1059 1060
    status = Status::Incomplete(Status::SubCode::kManualCompactionPaused);
  }
1061 1062 1063 1064 1065 1066 1067
  if (status.ok()) {
    status = input->status();
  }
  if (status.ok()) {
    status = c_iter->status();
  }

1068
  if (status.ok() && sub_compact->builder == nullptr &&
1069
      sub_compact->outputs.size() == 0 && !range_del_agg.IsEmpty()) {
1070
    // handle subcompaction containing only range deletions
1071 1072
    status = OpenCompactionOutputFile(sub_compact);
  }
1073 1074 1075 1076

  // Call FinishCompactionOutputFile() even if status is not ok: it needs to
  // close the output file.
  if (sub_compact->builder != nullptr) {
A
Andrew Kryczka 已提交
1077
    CompactionIterationStats range_del_out_stats;
1078
    Status s = FinishCompactionOutputFile(status, sub_compact, &range_del_agg,
1079
                                          &range_del_out_stats);
1080
    if (!s.ok() && status.ok()) {
1081 1082
      status = s;
    }
A
Andrew Kryczka 已提交
1083
    RecordDroppedKeys(range_del_out_stats, &sub_compact->compaction_job_stats);
1084 1085
  }

1086 1087 1088
  sub_compact->compaction_job_stats.cpu_micros =
      env_->NowCPUNanos() / 1000 - prev_cpu_micros;

1089 1090
  if (measure_io_stats_) {
    sub_compact->compaction_job_stats.file_write_nanos +=
I
Igor Canadi 已提交
1091
        IOSTATS(write_nanos) - prev_write_nanos;
1092
    sub_compact->compaction_job_stats.file_fsync_nanos +=
I
Igor Canadi 已提交
1093
        IOSTATS(fsync_nanos) - prev_fsync_nanos;
1094
    sub_compact->compaction_job_stats.file_range_sync_nanos +=
I
Igor Canadi 已提交
1095
        IOSTATS(range_sync_nanos) - prev_range_sync_nanos;
1096
    sub_compact->compaction_job_stats.file_prepare_write_nanos +=
I
Igor Canadi 已提交
1097
        IOSTATS(prepare_write_nanos) - prev_prepare_write_nanos;
1098
    sub_compact->compaction_job_stats.cpu_micros -=
1099 1100 1101
        (IOSTATS(cpu_write_nanos) - prev_cpu_write_nanos +
         IOSTATS(cpu_read_nanos) - prev_cpu_read_nanos) /
        1000;
1102
    if (prev_perf_level != PerfLevel::kEnableTimeAndCPUTimeExceptForMutex) {
1103 1104 1105
      SetPerfLevel(prev_perf_level);
    }
  }
1106 1107 1108 1109 1110 1111 1112 1113 1114 1115
#ifdef ROCKSDB_ASSERT_STATUS_CHECKED
  if (!status.ok()) {
    if (sub_compact->c_iter) {
      sub_compact->c_iter->status().PermitUncheckedError();
    }
    if (input) {
      input->status().PermitUncheckedError();
    }
  }
#endif  // ROCKSDB_ASSERT_STATUS_CHECKED
1116

1117 1118
  sub_compact->c_iter.reset();
  input.reset();
1119
  sub_compact->status = status;
I
Igor Canadi 已提交
1120 1121
}

1122
void CompactionJob::RecordDroppedKeys(
A
Andrew Kryczka 已提交
1123
    const CompactionIterationStats& c_iter_stats,
1124
    CompactionJobStats* compaction_job_stats) {
1125 1126 1127
  if (c_iter_stats.num_record_drop_user > 0) {
    RecordTick(stats_, COMPACTION_KEY_DROP_USER,
               c_iter_stats.num_record_drop_user);
1128
  }
1129 1130 1131
  if (c_iter_stats.num_record_drop_hidden > 0) {
    RecordTick(stats_, COMPACTION_KEY_DROP_NEWER_ENTRY,
               c_iter_stats.num_record_drop_hidden);
1132
    if (compaction_job_stats) {
1133 1134
      compaction_job_stats->num_records_replaced +=
          c_iter_stats.num_record_drop_hidden;
1135 1136
    }
  }
1137 1138 1139
  if (c_iter_stats.num_record_drop_obsolete > 0) {
    RecordTick(stats_, COMPACTION_KEY_DROP_OBSOLETE,
               c_iter_stats.num_record_drop_obsolete);
1140
    if (compaction_job_stats) {
1141 1142
      compaction_job_stats->num_expired_deletion_records +=
          c_iter_stats.num_record_drop_obsolete;
1143
    }
1144
  }
A
Andrew Kryczka 已提交
1145 1146 1147 1148 1149 1150 1151 1152
  if (c_iter_stats.num_record_drop_range_del > 0) {
    RecordTick(stats_, COMPACTION_KEY_DROP_RANGE_DEL,
               c_iter_stats.num_record_drop_range_del);
  }
  if (c_iter_stats.num_range_del_drop_obsolete > 0) {
    RecordTick(stats_, COMPACTION_RANGE_DEL_DROP_OBSOLETE,
               c_iter_stats.num_range_del_drop_obsolete);
  }
1153 1154 1155 1156
  if (c_iter_stats.num_optimized_del_drop_obsolete > 0) {
    RecordTick(stats_, COMPACTION_OPTIMIZED_DEL_DROP_OBSOLETE,
               c_iter_stats.num_optimized_del_drop_obsolete);
  }
1157 1158
}

1159
Status CompactionJob::FinishCompactionOutputFile(
1160
    const Status& input_status, SubcompactionState* sub_compact,
1161
    CompactionRangeDelAggregator* range_del_agg,
A
Andrew Kryczka 已提交
1162
    CompactionIterationStats* range_del_out_stats,
1163
    const Slice* next_table_min_key /* = nullptr */) {
1164 1165
  AutoThreadOperationStageUpdater stage_updater(
      ThreadStatus::STAGE_COMPACTION_SYNC_FILE);
1166 1167 1168
  assert(sub_compact != nullptr);
  assert(sub_compact->outfile);
  assert(sub_compact->builder != nullptr);
1169
  assert(sub_compact->current_output() != nullptr);
I
Igor Canadi 已提交
1170

1171
  uint64_t output_number = sub_compact->current_output()->meta.fd.GetNumber();
I
Igor Canadi 已提交
1172 1173
  assert(output_number != 0);

1174 1175
  ColumnFamilyData* cfd = sub_compact->compaction->column_family_data();
  const Comparator* ucmp = cfd->user_comparator();
1176 1177
  std::string file_checksum = kUnknownFileChecksum;
  std::string file_checksum_func_name = kUnknownFileChecksumFuncName;
1178

I
Igor Canadi 已提交
1179
  // Check for iterator errors
1180
  Status s = input_status;
1181
  auto meta = &sub_compact->current_output()->meta;
1182
  assert(meta != nullptr);
1183
  if (s.ok()) {
1184
    Slice lower_bound_guard, upper_bound_guard;
Z
zhangjinpeng1987 已提交
1185
    std::string smallest_user_key;
1186
    const Slice *lower_bound, *upper_bound;
1187
    bool lower_bound_from_sub_compact = false;
1188 1189 1190 1191
    if (sub_compact->outputs.size() == 1) {
      // For the first output table, include range tombstones before the min key
      // but after the subcompaction boundary.
      lower_bound = sub_compact->start;
1192
      lower_bound_from_sub_compact = true;
1193 1194 1195 1196
    } else if (meta->smallest.size() > 0) {
      // For subsequent output tables, only include range tombstones from min
      // key onwards since the previous file was extended to contain range
      // tombstones falling before min key.
Z
zhangjinpeng1987 已提交
1197 1198
      smallest_user_key = meta->smallest.user_key().ToString(false /*hex*/);
      lower_bound_guard = Slice(smallest_user_key);
1199 1200 1201 1202 1203
      lower_bound = &lower_bound_guard;
    } else {
      lower_bound = nullptr;
    }
    if (next_table_min_key != nullptr) {
1204 1205 1206 1207 1208 1209
      // This may be the last file in the subcompaction in some cases, so we
      // need to compare the end key of subcompaction with the next file start
      // key. When the end key is chosen by the subcompaction, we know that
      // it must be the biggest key in output file. Therefore, it is safe to
      // use the smaller key as the upper bound of the output file, to ensure
      // that there is no overlapping between different output files.
1210
      upper_bound_guard = ExtractUserKey(*next_table_min_key);
1211 1212 1213 1214 1215 1216
      if (sub_compact->end != nullptr &&
          ucmp->Compare(upper_bound_guard, *sub_compact->end) >= 0) {
        upper_bound = sub_compact->end;
      } else {
        upper_bound = &upper_bound_guard;
      }
1217 1218 1219 1220 1221
    } else {
      // This is the last file in the subcompaction, so extend until the
      // subcompaction ends.
      upper_bound = sub_compact->end;
    }
1222 1223 1224 1225
    auto earliest_snapshot = kMaxSequenceNumber;
    if (existing_snapshots_.size() > 0) {
      earliest_snapshot = existing_snapshots_[0];
    }
1226 1227 1228 1229 1230 1231 1232
    bool has_overlapping_endpoints;
    if (upper_bound != nullptr && meta->largest.size() > 0) {
      has_overlapping_endpoints =
          ucmp->Compare(meta->largest.user_key(), *upper_bound) == 0;
    } else {
      has_overlapping_endpoints = false;
    }
1233

1234 1235 1236 1237 1238 1239 1240
    // The end key of the subcompaction must be bigger or equal to the upper
    // bound. If the end of subcompaction is null or the upper bound is null,
    // it means that this file is the last file in the compaction. So there
    // will be no overlapping between this file and others.
    assert(sub_compact->end == nullptr ||
           upper_bound == nullptr ||
           ucmp->Compare(*upper_bound , *sub_compact->end) <= 0);
1241 1242 1243 1244 1245 1246 1247 1248 1249 1250
    auto it = range_del_agg->NewIterator(lower_bound, upper_bound,
                                         has_overlapping_endpoints);
    // Position the range tombstone output iterator. There may be tombstone
    // fragments that are entirely out of range, so make sure that we do not
    // include those.
    if (lower_bound != nullptr) {
      it->Seek(*lower_bound);
    } else {
      it->SeekToFirst();
    }
1251
    TEST_SYNC_POINT("CompactionJob::FinishCompactionOutputFile1");
1252 1253
    for (; it->Valid(); it->Next()) {
      auto tombstone = it->Tombstone();
1254 1255 1256 1257 1258 1259 1260 1261 1262 1263 1264 1265
      if (upper_bound != nullptr) {
        int cmp = ucmp->Compare(*upper_bound, tombstone.start_key_);
        if ((has_overlapping_endpoints && cmp < 0) ||
            (!has_overlapping_endpoints && cmp <= 0)) {
          // Tombstones starting after upper_bound only need to be included in
          // the next table. If the current SST ends before upper_bound, i.e.,
          // `has_overlapping_endpoints == false`, we can also skip over range
          // tombstones that start exactly at upper_bound. Such range tombstones
          // will be included in the next file and are not relevant to the point
          // keys or endpoints of the current file.
          break;
        }
1266 1267 1268 1269 1270 1271 1272 1273 1274 1275 1276
      }

      if (bottommost_level_ && tombstone.seq_ <= earliest_snapshot) {
        // TODO(andrewkr): tombstones that span multiple output files are
        // counted for each compaction output file, so lots of double counting.
        range_del_out_stats->num_range_del_drop_obsolete++;
        range_del_out_stats->num_record_drop_obsolete++;
        continue;
      }

      auto kv = tombstone.Serialize();
1277 1278
      assert(lower_bound == nullptr ||
             ucmp->Compare(*lower_bound, kv.second) < 0);
1279 1280
      sub_compact->AddToBuilder(kv.first.Encode(), kv.second,
                                paranoid_file_checks_);
1281 1282 1283 1284 1285 1286 1287
      InternalKey smallest_candidate = std::move(kv.first);
      if (lower_bound != nullptr &&
          ucmp->Compare(smallest_candidate.user_key(), *lower_bound) <= 0) {
        // Pretend the smallest key has the same user key as lower_bound
        // (the max key in the previous table or subcompaction) in order for
        // files to appear key-space partitioned.
        //
1288 1289 1290 1291 1292 1293 1294 1295 1296 1297 1298 1299 1300 1301 1302 1303 1304 1305
        // When lower_bound is chosen by a subcompaction, we know that
        // subcompactions over smaller keys cannot contain any keys at
        // lower_bound. We also know that smaller subcompactions exist, because
        // otherwise the subcompaction woud be unbounded on the left. As a
        // result, we know that no other files on the output level will contain
        // actual keys at lower_bound (an output file may have a largest key of
        // lower_bound@kMaxSequenceNumber, but this only indicates a large range
        // tombstone was truncated). Therefore, it is safe to use the
        // tombstone's sequence number, to ensure that keys at lower_bound at
        // lower levels are covered by truncated tombstones.
        //
        // If lower_bound was chosen by the smallest data key in the file,
        // choose lowest seqnum so this file's smallest internal key comes after
        // the previous file's largest. The fake seqnum is OK because the read
        // path's file-picking code only considers user key.
        smallest_candidate = InternalKey(
            *lower_bound, lower_bound_from_sub_compact ? tombstone.seq_ : 0,
            kTypeRangeDeletion);
1306 1307 1308 1309 1310 1311 1312 1313 1314 1315 1316 1317 1318 1319 1320 1321 1322 1323 1324 1325 1326
      }
      InternalKey largest_candidate = tombstone.SerializeEndKey();
      if (upper_bound != nullptr &&
          ucmp->Compare(*upper_bound, largest_candidate.user_key()) <= 0) {
        // Pretend the largest key has the same user key as upper_bound (the
        // min key in the following table or subcompaction) in order for files
        // to appear key-space partitioned.
        //
        // Choose highest seqnum so this file's largest internal key comes
        // before the next file's/subcompaction's smallest. The fake seqnum is
        // OK because the read path's file-picking code only considers the user
        // key portion.
        //
        // Note Seek() also creates InternalKey with (user_key,
        // kMaxSequenceNumber), but with kTypeDeletion (0x7) instead of
        // kTypeRangeDeletion (0xF), so the range tombstone comes before the
        // Seek() key in InternalKey's ordering. So Seek() will look in the
        // next file for the user key.
        largest_candidate =
            InternalKey(*upper_bound, kMaxSequenceNumber, kTypeRangeDeletion);
      }
1327 1328 1329 1330 1331 1332
#ifndef NDEBUG
      SequenceNumber smallest_ikey_seqnum = kMaxSequenceNumber;
      if (meta->smallest.size() > 0) {
        smallest_ikey_seqnum = GetInternalKeySeqno(meta->smallest.Encode());
      }
#endif
1333 1334 1335
      meta->UpdateBoundariesForRange(smallest_candidate, largest_candidate,
                                     tombstone.seq_,
                                     cfd->internal_comparator());
1336 1337 1338 1339 1340 1341 1342 1343

      // The smallest key in a file is used for range tombstone truncation, so
      // it cannot have a seqnum of 0 (unless the smallest data key in a file
      // has a seqnum of 0). Otherwise, the truncated tombstone may expose
      // deleted keys at lower levels.
      assert(smallest_ikey_seqnum == 0 ||
             ExtractInternalKeyFooter(meta->smallest.Encode()) !=
                 PackSequenceAndType(0, kTypeRangeDeletion));
1344
    }
S
Siying Dong 已提交
1345
    meta->marked_for_compaction = sub_compact->builder->NeedCompact();
1346
  }
1347
  const uint64_t current_entries = sub_compact->builder->NumEntries();
I
Igor Canadi 已提交
1348
  if (s.ok()) {
1349
    s = sub_compact->builder->Finish();
I
Igor Canadi 已提交
1350
  } else {
1351
    sub_compact->builder->Abandon();
I
Igor Canadi 已提交
1352
  }
1353 1354 1355
  IOStatus io_s = sub_compact->builder->io_status();
  if (s.ok()) {
    s = io_s;
1356
  }
1357
  const uint64_t current_bytes = sub_compact->builder->FileSize();
S
Siying Dong 已提交
1358 1359 1360
  if (s.ok()) {
    meta->fd.file_size = current_bytes;
  }
1361
  sub_compact->current_output()->finished = true;
1362
  sub_compact->total_bytes += current_bytes;
I
Igor Canadi 已提交
1363 1364

  // Finish and check for file errors
S
Sagar Vemuri 已提交
1365
  if (s.ok()) {
1366
    StopWatch sw(env_, stats_, COMPACTION_OUTFILE_SYNC_MICROS);
1367
    io_s = sub_compact->outfile->Sync(db_options_.use_fsync);
I
Igor Canadi 已提交
1368
  }
1369
  if (s.ok() && io_s.ok()) {
1370 1371
    io_s = sub_compact->outfile->Close();
  }
1372
  if (s.ok() && io_s.ok()) {
1373 1374 1375 1376
    // Add the checksum information to file metadata.
    meta->file_checksum = sub_compact->outfile->GetFileChecksum();
    meta->file_checksum_func_name =
        sub_compact->outfile->GetFileChecksumFuncName();
1377 1378
    file_checksum = meta->file_checksum;
    file_checksum_func_name = meta->file_checksum_func_name;
1379
  }
1380
  if (s.ok()) {
1381
    s = io_s;
I
Igor Canadi 已提交
1382
  }
1383 1384
  if (sub_compact->io_status.ok()) {
    sub_compact->io_status = io_s;
1385 1386 1387
    // Since this error is really a copy of the
    // "normal" status, it does not also need to be checked
    sub_compact->io_status.PermitUncheckedError();
1388
  }
1389
  sub_compact->outfile.reset();
I
Igor Canadi 已提交
1390

1391 1392 1393 1394 1395 1396
  TableProperties tp;
  if (s.ok()) {
    tp = sub_compact->builder->GetTableProperties();
  }

  if (s.ok() && current_entries == 0 && tp.num_range_deletions == 0) {
Z
zhangjinpeng1987 已提交
1397 1398 1399
    // If there is nothing to output, no necessary to generate a sst file.
    // This happens when the output level is bottom level, at the same time
    // the sub_compact output nothing.
1400 1401 1402
    std::string fname =
        TableFileName(sub_compact->compaction->immutable_cf_options()->cf_paths,
                      meta->fd.GetNumber(), meta->fd.GetPathId());
Z
zhangjinpeng1987 已提交
1403 1404 1405 1406 1407 1408
    env_->DeleteFile(fname);

    // Also need to remove the file from outputs, or it will be added to the
    // VersionEdit.
    assert(!sub_compact->outputs.empty());
    sub_compact->outputs.pop_back();
1409
    meta = nullptr;
Z
zhangjinpeng1987 已提交
1410 1411
  }

1412
  if (s.ok() && (current_entries > 0 || tp.num_range_deletions > 0)) {
1413
    // Output to event logger and fire events.
1414 1415 1416 1417 1418 1419 1420 1421
    sub_compact->current_output()->table_properties =
        std::make_shared<TableProperties>(tp);
    ROCKS_LOG_INFO(db_options_.info_log,
                   "[%s] [JOB %d] Generated table #%" PRIu64 ": %" PRIu64
                   " keys, %" PRIu64 " bytes%s",
                   cfd->GetName().c_str(), job_id_, output_number,
                   current_entries, current_bytes,
                   meta->marked_for_compaction ? " (need compaction)" : "");
I
Igor Canadi 已提交
1422
  }
S
Siying Dong 已提交
1423 1424
  std::string fname;
  FileDescriptor output_fd;
1425
  uint64_t oldest_blob_file_number = kInvalidBlobFileNumber;
S
Siying Dong 已提交
1426
  if (meta != nullptr) {
1427 1428 1429
    fname =
        TableFileName(sub_compact->compaction->immutable_cf_options()->cf_paths,
                      meta->fd.GetNumber(), meta->fd.GetPathId());
S
Siying Dong 已提交
1430
    output_fd = meta->fd;
1431
    oldest_blob_file_number = meta->oldest_blob_file_number;
S
Siying Dong 已提交
1432 1433 1434
  } else {
    fname = "(nil)";
  }
1435 1436
  EventHelpers::LogAndNotifyTableFileCreationFinished(
      event_logger_, cfd->ioptions()->listeners, dbname_, cfd->GetName(), fname,
1437
      job_id_, output_fd, oldest_blob_file_number, tp,
1438 1439
      TableFileCreationReason::kCompaction, s, file_checksum,
      file_checksum_func_name);
1440

1441
#ifndef ROCKSDB_LITE
1442 1443 1444
  // Report new file to SstFileManagerImpl
  auto sfm =
      static_cast<SstFileManagerImpl*>(db_options_.sst_file_manager.get());
S
Siying Dong 已提交
1445
  if (sfm && meta != nullptr && meta->fd.GetPathId() == 0) {
1446 1447 1448 1449
    Status add_s = sfm->OnAddFile(fname);
    if (!add_s.ok() && s.ok()) {
      s = add_s;
    }
1450
    if (sfm->IsMaxAllowedSpaceReached()) {
1451 1452
      // TODO(ajkr): should we return OK() if max space was reached by the final
      // compaction output file (similarly to how flush works when full)?
1453
      s = Status::SpaceLimit("Max allowed space was reached");
1454 1455 1456
      TEST_SYNC_POINT(
          "CompactionJob::FinishCompactionOutputFile:"
          "MaxAllowedSpaceReached");
1457
      InstrumentedMutexLock l(db_mutex_);
1458
      db_error_handler_->SetBGError(s, BackgroundErrorReason::kCompaction);
1459 1460
    }
  }
1461
#endif
1462

1463
  sub_compact->builder.reset();
1464
  sub_compact->current_output_file_size = 0;
I
Igor Canadi 已提交
1465 1466 1467
  return s;
}

I
Igor Canadi 已提交
1468
Status CompactionJob::InstallCompactionResults(
1469 1470
    const MutableCFOptions& mutable_cf_options) {
  db_mutex_->AssertHeld();
I
Igor Canadi 已提交
1471

1472
  auto* compaction = compact_->compaction;
I
Igor Canadi 已提交
1473 1474 1475 1476
  // paranoia: verify that the files that we started with
  // still exist in the current version and in the same original level.
  // This ensures that a concurrent compaction did not erroneously
  // pick the same files to compact_.
1477 1478 1479
  if (!versions_->VerifyCompactionFileConsistency(compaction)) {
    Compaction::InputLevelSummaryBuffer inputs_summary;

1480 1481 1482
    ROCKS_LOG_ERROR(db_options_.info_log, "[%s] [JOB %d] Compaction %s aborted",
                    compaction->column_family_data()->GetName().c_str(),
                    job_id_, compaction->InputLevelSummary(&inputs_summary));
I
Igor Canadi 已提交
1483 1484 1485
    return Status::Corruption("Compaction input files inconsistent");
  }

1486 1487
  {
    Compaction::InputLevelSummaryBuffer inputs_summary;
1488 1489
    ROCKS_LOG_INFO(
        db_options_.info_log, "[%s] [JOB %d] Compacted %s => %" PRIu64 " bytes",
1490 1491 1492
        compaction->column_family_data()->GetName().c_str(), job_id_,
        compaction->InputLevelSummary(&inputs_summary), compact_->total_bytes);
  }
I
Igor Canadi 已提交
1493

G
Gihwan Oh 已提交
1494
  // Add compaction inputs
1495
  compaction->AddInputDeletions(compact_->compaction->edit());
1496

1497 1498 1499
  for (const auto& sub_compact : compact_->sub_compact_states) {
    for (const auto& out : sub_compact.outputs) {
      compaction->edit()->AddFile(compaction->output_level(), out.meta);
1500
    }
1501 1502
  }
  return versions_->LogAndApply(compaction->column_family_data(),
I
Igor Canadi 已提交
1503
                                mutable_cf_options, compaction->edit(),
1504
                                db_mutex_, db_directory_);
I
Igor Canadi 已提交
1505 1506 1507 1508
}

void CompactionJob::RecordCompactionIOStats() {
  RecordTick(stats_, COMPACT_READ_BYTES, IOSTATS(bytes_read));
1509 1510 1511 1512 1513 1514 1515 1516 1517 1518 1519 1520 1521
  RecordTick(stats_, COMPACT_WRITE_BYTES, IOSTATS(bytes_written));
  CompactionReason compaction_reason =
      compact_->compaction->compaction_reason();
  if (compaction_reason == CompactionReason::kFilesMarkedForCompaction) {
    RecordTick(stats_, COMPACT_READ_BYTES_MARKED, IOSTATS(bytes_read));
    RecordTick(stats_, COMPACT_WRITE_BYTES_MARKED, IOSTATS(bytes_written));
  } else if (compaction_reason == CompactionReason::kPeriodicCompaction) {
    RecordTick(stats_, COMPACT_READ_BYTES_PERIODIC, IOSTATS(bytes_read));
    RecordTick(stats_, COMPACT_WRITE_BYTES_PERIODIC, IOSTATS(bytes_written));
  } else if (compaction_reason == CompactionReason::kTtl) {
    RecordTick(stats_, COMPACT_READ_BYTES_TTL, IOSTATS(bytes_read));
    RecordTick(stats_, COMPACT_WRITE_BYTES_TTL, IOSTATS(bytes_written));
  }
1522 1523
  ThreadStatusUtil::IncreaseThreadOperationProperty(
      ThreadStatus::COMPACTION_BYTES_READ, IOSTATS(bytes_read));
I
Igor Canadi 已提交
1524
  IOSTATS_RESET(bytes_read);
1525 1526
  ThreadStatusUtil::IncreaseThreadOperationProperty(
      ThreadStatus::COMPACTION_BYTES_WRITTEN, IOSTATS(bytes_written));
I
Igor Canadi 已提交
1527 1528 1529
  IOSTATS_RESET(bytes_written);
}

1530 1531
Status CompactionJob::OpenCompactionOutputFile(
    SubcompactionState* sub_compact) {
1532 1533
  assert(sub_compact != nullptr);
  assert(sub_compact->builder == nullptr);
1534 1535
  // no need to lock because VersionSet::next_file_number_ is atomic
  uint64_t file_number = versions_->NewFileNumber();
1536 1537 1538
  std::string fname =
      TableFileName(sub_compact->compaction->immutable_cf_options()->cf_paths,
                    file_number, sub_compact->compaction->output_path_id());
1539 1540 1541 1542 1543 1544 1545 1546
  // Fire events.
  ColumnFamilyData* cfd = sub_compact->compaction->column_family_data();
#ifndef ROCKSDB_LITE
  EventHelpers::NotifyTableFileCreationStarted(
      cfd->ioptions()->listeners, dbname_, cfd->GetName(), fname, job_id_,
      TableFileCreationReason::kCompaction);
#endif  // !ROCKSDB_LITE
  // Make the output file
1547
  std::unique_ptr<FSWritableFile> writable_file;
A
Aliaksei Sandryhaila 已提交
1548
#ifndef NDEBUG
1549
  bool syncpoint_arg = file_options_.use_direct_writes;
1550
  TEST_SYNC_POINT_CALLBACK("CompactionJob::OpenCompactionOutputFile",
1551
                           &syncpoint_arg);
A
Aliaksei Sandryhaila 已提交
1552
#endif
1553
  Status s;
1554 1555
  IOStatus io_s =
      NewWritableFile(fs_.get(), fname, &writable_file, file_options_);
1556 1557 1558
  s = io_s;
  if (sub_compact->io_status.ok()) {
    sub_compact->io_status = io_s;
1559 1560 1561
    // Since this error is really a copy of the io_s that is checked below as s,
    // it does not also need to be checked.
    sub_compact->io_status.PermitUncheckedError();
1562
  }
I
Igor Canadi 已提交
1563
  if (!s.ok()) {
1564 1565
    ROCKS_LOG_ERROR(
        db_options_.info_log,
1566
        "[%s] [JOB %d] OpenCompactionOutputFiles for table #%" PRIu64
1567
        " fails at NewWritableFile with status %s",
1568 1569
        sub_compact->compaction->column_family_data()->GetName().c_str(),
        job_id_, file_number, s.ToString().c_str());
I
Igor Canadi 已提交
1570
    LogFlush(db_options_.info_log);
1571 1572
    EventHelpers::LogAndNotifyTableFileCreationFinished(
        event_logger_, cfd->ioptions()->listeners, dbname_, cfd->GetName(),
1573
        fname, job_id_, FileDescriptor(), kInvalidBlobFileNumber,
1574 1575
        TableProperties(), TableFileCreationReason::kCompaction, s,
        kUnknownFileChecksum, kUnknownFileChecksumFuncName);
I
Igor Canadi 已提交
1576 1577
    return s;
  }
1578

1579 1580 1581 1582 1583 1584 1585 1586 1587 1588 1589 1590 1591 1592 1593 1594 1595 1596 1597 1598 1599 1600
  // Try to figure out the output file's oldest ancester time.
  int64_t temp_current_time = 0;
  auto get_time_status = env_->GetCurrentTime(&temp_current_time);
  // Safe to proceed even if GetCurrentTime fails. So, log and proceed.
  if (!get_time_status.ok()) {
    ROCKS_LOG_WARN(db_options_.info_log,
                   "Failed to get current time. Status: %s",
                   get_time_status.ToString().c_str());
  }
  uint64_t current_time = static_cast<uint64_t>(temp_current_time);
  uint64_t oldest_ancester_time =
      sub_compact->compaction->MinInputFileOldestAncesterTime();
  if (oldest_ancester_time == port::kMaxUint64) {
    oldest_ancester_time = current_time;
  }

  // Initialize a SubcompactionState::Output and add it to sub_compact->outputs
  {
    SubcompactionState::Output out;
    out.meta.fd = FileDescriptor(file_number,
                                 sub_compact->compaction->output_path_id(), 0);
    out.meta.oldest_ancester_time = oldest_ancester_time;
1601
    out.meta.file_creation_time = current_time;
1602
    out.finished = false;
1603
    out.paranoid_hash = 0;
1604 1605
    sub_compact->outputs.push_back(out);
  }
I
Igor Canadi 已提交
1606

1607
  writable_file->SetIOPriority(Env::IOPriority::IO_LOW);
S
Stream  
Shaohua Li 已提交
1608
  writable_file->SetWriteLifeTimeHint(write_hint_);
1609 1610
  writable_file->SetPreallocationBlockSize(static_cast<size_t>(
      sub_compact->compaction->OutputFilePreallocationSize()));
1611 1612
  const auto& listeners =
      sub_compact->compaction->immutable_cf_options()->listeners;
1613 1614 1615 1616
  sub_compact->outfile.reset(new WritableFileWriter(
      std::move(writable_file), fname, file_options_, env_, io_tracer_,
      db_options_.statistics.get(), listeners,
      db_options_.file_checksum_gen_factory.get()));
I
Igor Canadi 已提交
1617

1618 1619 1620
  // If the Column family flag is to only optimize filters for hits,
  // we can skip creating filters if this is the bottommost_level where
  // data is going to be found
1621 1622
  bool skip_filters =
      cfd->ioptions()->optimize_filters_for_hits && bottommost_level_;
S
Sagar Vemuri 已提交
1623

1624
  sub_compact->builder.reset(NewTableBuilder(
1625 1626 1627 1628
      *cfd->ioptions(), *(sub_compact->compaction->mutable_cf_options()),
      cfd->internal_comparator(), cfd->int_tbl_prop_collector_factories(),
      cfd->GetID(), cfd->GetName(), sub_compact->outfile.get(),
      sub_compact->compaction->output_compression(),
1629
      0 /*sample_for_compression */,
1630
      sub_compact->compaction->output_compression_opts(),
1631 1632
      sub_compact->compaction->output_level(), skip_filters,
      oldest_ancester_time, 0 /* oldest_key_time */,
1633 1634
      sub_compact->compaction->max_output_file_size(), current_time, db_id_,
      db_session_id_));
I
Igor Canadi 已提交
1635 1636 1637 1638
  LogFlush(db_options_.info_log);
  return s;
}

1639
void CompactionJob::CleanupCompaction() {
1640
  for (SubcompactionState& sub_compact : compact_->sub_compact_states) {
1641
    const auto& sub_status = sub_compact.status;
I
Igor Canadi 已提交
1642

1643 1644 1645 1646 1647 1648 1649
    if (sub_compact.builder != nullptr) {
      // May happen if we get a shutdown call in the middle of compaction
      sub_compact.builder->Abandon();
      sub_compact.builder.reset();
    } else {
      assert(!sub_status.ok() || sub_compact.outfile == nullptr);
    }
1650
    for (const auto& out : sub_compact.outputs) {
1651 1652 1653
      // If this file was inserted into the table cache then remove
      // them here because this compaction was not committed.
      if (!sub_status.ok()) {
1654
        TableCache::Evict(table_cache_.get(), out.meta.fd.GetNumber());
1655
      }
I
Igor Canadi 已提交
1656
    }
1657 1658 1659
    // TODO: sub_compact.io_status is not checked like status. Not sure if thats
    // intentional. So ignoring the io_status as of now.
    sub_compact.io_status.PermitUncheckedError();
I
Igor Canadi 已提交
1660 1661 1662 1663 1664
  }
  delete compact_;
  compact_ = nullptr;
}

1665 1666
#ifndef ROCKSDB_LITE
namespace {
1667
void CopyPrefix(const Slice& src, size_t prefix_length, std::string* dst) {
1668 1669 1670
  assert(prefix_length > 0);
  size_t length = src.size() > prefix_length ? prefix_length : src.size();
  dst->assign(src.data(), length);
1671 1672 1673 1674 1675
}
}  // namespace

#endif  // !ROCKSDB_LITE

A
Andres Notzli 已提交
1676
void CompactionJob::UpdateCompactionStats() {
1677 1678 1679 1680 1681 1682
  Compaction* compaction = compact_->compaction;
  compaction_stats_.num_input_files_in_non_output_levels = 0;
  compaction_stats_.num_input_files_in_output_level = 0;
  for (int input_level = 0;
       input_level < static_cast<int>(compaction->num_input_levels());
       ++input_level) {
1683
    if (compaction->level(input_level) != compaction->output_level()) {
1684 1685
      UpdateCompactionInputStatsHelper(
          &compaction_stats_.num_input_files_in_non_output_levels,
1686
          &compaction_stats_.bytes_read_non_output_levels, input_level);
1687 1688 1689
    } else {
      UpdateCompactionInputStatsHelper(
          &compaction_stats_.num_input_files_in_output_level,
1690
          &compaction_stats_.bytes_read_output_level, input_level);
1691 1692
    }
  }
A
Andres Notzli 已提交
1693

1694 1695
  uint64_t num_output_records = 0;

1696 1697 1698 1699 1700 1701 1702 1703 1704
  for (const auto& sub_compact : compact_->sub_compact_states) {
    size_t num_output_files = sub_compact.outputs.size();
    if (sub_compact.builder != nullptr) {
      // An error occurred so ignore the last output.
      assert(num_output_files > 0);
      --num_output_files;
    }
    compaction_stats_.num_output_files += static_cast<int>(num_output_files);

1705 1706
    num_output_records += sub_compact.num_output_records;

1707 1708
    for (const auto& out : sub_compact.outputs) {
      compaction_stats_.bytes_written += out.meta.fd.file_size;
1709
    }
1710 1711 1712 1713 1714
  }

  if (compaction_stats_.num_input_records > num_output_records) {
    compaction_stats_.num_dropped_records =
        compaction_stats_.num_input_records - num_output_records;
A
Andres Notzli 已提交
1715
  }
1716 1717
}

1718 1719 1720
void CompactionJob::UpdateCompactionInputStatsHelper(int* num_files,
                                                     uint64_t* bytes_read,
                                                     int input_level) {
1721 1722 1723 1724 1725 1726 1727 1728 1729 1730 1731 1732
  const Compaction* compaction = compact_->compaction;
  auto num_input_files = compaction->num_input_files(input_level);
  *num_files += static_cast<int>(num_input_files);

  for (size_t i = 0; i < num_input_files; ++i) {
    const auto* file_meta = compaction->input(input_level, i);
    *bytes_read += file_meta->fd.GetFileSize();
    compaction_stats_.num_input_records +=
        static_cast<uint64_t>(file_meta->num_entries);
  }
}

1733 1734 1735 1736 1737 1738 1739 1740
void CompactionJob::UpdateCompactionJobStats(
    const InternalStats::CompactionStats& stats) const {
#ifndef ROCKSDB_LITE
  if (compaction_job_stats_) {
    compaction_job_stats_->elapsed_micros = stats.micros;

    // input information
    compaction_job_stats_->total_input_bytes =
1741
        stats.bytes_read_non_output_levels + stats.bytes_read_output_level;
1742
    compaction_job_stats_->num_input_records = stats.num_input_records;
1743
    compaction_job_stats_->num_input_files =
1744 1745
        stats.num_input_files_in_non_output_levels +
        stats.num_input_files_in_output_level;
1746
    compaction_job_stats_->num_input_files_at_output_level =
1747
        stats.num_input_files_in_output_level;
1748 1749 1750

    // output information
    compaction_job_stats_->total_output_bytes = stats.bytes_written;
1751
    compaction_job_stats_->num_output_records = compact_->num_output_records;
1752
    compaction_job_stats_->num_output_files = stats.num_output_files;
1753

1754
    if (compact_->NumOutputFiles() > 0U) {
1755 1756 1757 1758 1759 1760
      CopyPrefix(compact_->SmallestUserKey(),
                 CompactionJobStats::kMaxPrefixLength,
                 &compaction_job_stats_->smallest_output_key_prefix);
      CopyPrefix(compact_->LargestUserKey(),
                 CompactionJobStats::kMaxPrefixLength,
                 &compaction_job_stats_->largest_output_key_prefix);
1761 1762
    }
  }
1763 1764
#else
  (void)stats;
1765 1766 1767
#endif  // !ROCKSDB_LITE
}

1768 1769 1770 1771
void CompactionJob::LogCompaction() {
  Compaction* compaction = compact_->compaction;
  ColumnFamilyData* cfd = compaction->column_family_data();

A
Andres Notzli 已提交
1772 1773 1774 1775
  // Let's check if anything will get logged. Don't prepare all the info if
  // we're not logging
  if (db_options_.info_log_level <= InfoLogLevel::INFO_LEVEL) {
    Compaction::InputLevelSummaryBuffer inputs_summary;
1776 1777 1778 1779
    ROCKS_LOG_INFO(
        db_options_.info_log, "[%s] [JOB %d] Compacting %s, score %.2f",
        cfd->GetName().c_str(), job_id_,
        compaction->InputLevelSummary(&inputs_summary), compaction->score());
A
Andres Notzli 已提交
1780 1781
    char scratch[2345];
    compaction->Summary(scratch, sizeof(scratch));
1782 1783
    ROCKS_LOG_INFO(db_options_.info_log, "[%s] Compaction start summary: %s\n",
                   cfd->GetName().c_str(), scratch);
A
Andres Notzli 已提交
1784 1785
    // build event logger report
    auto stream = event_logger_->Log();
1786 1787 1788 1789
    stream << "job" << job_id_ << "event"
           << "compaction_started"
           << "compaction_reason"
           << GetCompactionReasonString(compaction->compaction_reason());
A
Andres Notzli 已提交
1790 1791 1792 1793 1794 1795 1796 1797 1798 1799 1800 1801 1802
    for (size_t i = 0; i < compaction->num_input_levels(); ++i) {
      stream << ("files_L" + ToString(compaction->level(i)));
      stream.StartArray();
      for (auto f : *compaction->inputs(i)) {
        stream << f->fd.GetNumber();
      }
      stream.EndArray();
    }
    stream << "score" << compaction->score() << "input_data_size"
           << compaction->CalculateTotalInputSize();
  }
}

1803
}  // namespace ROCKSDB_NAMESPACE