Coverage Report

Created: 2026-09-28 07:52

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/rocksdb/db/flush_job.cc
Line
Count
Source
1
//  Copyright (c) 2011-present, Facebook, Inc.  All rights reserved.
2
//  This source code is licensed under both the GPLv2 (found in the
3
//  COPYING file in the root directory) and Apache 2.0 License
4
//  (found in the LICENSE.Apache file in the root directory).
5
//
6
// Copyright (c) 2011 The LevelDB Authors. All rights reserved.
7
// Use of this source code is governed by a BSD-style license that can be
8
// found in the LICENSE file. See the AUTHORS file for names of contributors.
9
10
#include "db/flush_job.h"
11
12
#include <algorithm>
13
#include <cinttypes>
14
#include <vector>
15
16
#include "db/builder.h"
17
#include "db/db_iter.h"
18
#include "db/dbformat.h"
19
#include "db/event_helpers.h"
20
#include "db/log_reader.h"
21
#include "db/log_writer.h"
22
#include "db/memtable.h"
23
#include "db/memtable_list.h"
24
#include "db/merge_context.h"
25
#include "db/range_tombstone_fragmenter.h"
26
#include "db/table_cache.h"
27
#include "db/version_edit.h"
28
#include "db/version_set.h"
29
#include "file/file_util.h"
30
#include "file/filename.h"
31
#include "logging/event_logger.h"
32
#include "logging/log_buffer.h"
33
#include "logging/logging.h"
34
#include "monitoring/iostats_context_imp.h"
35
#include "monitoring/perf_context_imp.h"
36
#include "monitoring/statistics_impl.h"
37
#include "monitoring/thread_status_util.h"
38
#include "port/port.h"
39
#include "rocksdb/db.h"
40
#include "rocksdb/env.h"
41
#include "rocksdb/statistics.h"
42
#include "rocksdb/status.h"
43
#include "rocksdb/table.h"
44
#include "table/merging_iterator.h"
45
#include "table/table_builder.h"
46
#include "table/two_level_iterator.h"
47
#include "test_util/sync_point.h"
48
#include "util/coding.h"
49
#include "util/mutexlock.h"
50
#include "util/stop_watch.h"
51
52
namespace ROCKSDB_NAMESPACE {
53
54
2.33k
const char* GetFlushReasonString(FlushReason flush_reason) {
55
2.33k
  switch (flush_reason) {
56
0
    case FlushReason::kOthers:
57
0
      return "Other Reasons";
58
0
    case FlushReason::kGetLiveFiles:
59
0
      return "Get Live Files";
60
0
    case FlushReason::kShutDown:
61
0
      return "Shut down";
62
0
    case FlushReason::kExternalFileIngestion:
63
0
      return "External File Ingestion";
64
2.33k
    case FlushReason::kManualCompaction:
65
2.33k
      return "Manual Compaction";
66
0
    case FlushReason::kWriteBufferManager:
67
0
      return "Write Buffer Manager";
68
0
    case FlushReason::kWriteBufferFull:
69
0
      return "Write Buffer Full";
70
0
    case FlushReason::kTest:
71
0
      return "Test";
72
0
    case FlushReason::kDeleteFiles:
73
0
      return "Delete Files";
74
0
    case FlushReason::kAutoCompaction:
75
0
      return "Auto Compaction";
76
0
    case FlushReason::kManualFlush:
77
0
      return "Manual Flush";
78
0
    case FlushReason::kErrorRecovery:
79
0
      return "Error Recovery";
80
0
    case FlushReason::kErrorRecoveryRetryFlush:
81
0
      return "Error Recovery Retry Flush";
82
0
    case FlushReason::kWalFull:
83
0
      return "WAL Full";
84
0
    case FlushReason::kCatchUpAfterErrorRecovery:
85
0
      return "Catch Up After Error Recovery";
86
0
    case FlushReason::kMemtableMaxRangeDeletions:
87
0
      return "Memtable Max Range Deletions";
88
0
    default:
89
0
      return "Invalid";
90
2.33k
  }
91
2.33k
}
92
93
FlushJob::FlushJob(
94
    const std::string& dbname, ColumnFamilyData* cfd,
95
    const ImmutableDBOptions& db_options,
96
    const MutableCFOptions& mutable_cf_options, uint64_t max_memtable_id,
97
    const FileOptions& file_options, VersionSet* versions,
98
    InstrumentedMutex* db_mutex, std::atomic<bool>* shutting_down,
99
    JobContext* job_context, FlushReason flush_reason, LogBuffer* log_buffer,
100
    FSDirectory* db_directory, FSDirectory* output_file_directory,
101
    CompressionType output_compression, Statistics* stats,
102
    EventLogger* event_logger, bool measure_io_stats,
103
    const bool sync_output_directory, const bool write_manifest,
104
    Env::Priority thread_pri, const std::shared_ptr<IOTracer>& io_tracer,
105
    std::shared_ptr<const SeqnoToTimeMapping> seqno_to_time_mapping,
106
    const std::string& db_id, const std::string& db_session_id,
107
    std::string full_history_ts_low, BlobFileCompletionCallback* blob_callback,
108
    bool fast_sst_open)
109
2.33k
    : dbname_(dbname),
110
2.33k
      db_id_(db_id),
111
2.33k
      db_session_id_(db_session_id),
112
2.33k
      cfd_(cfd),
113
2.33k
      db_options_(db_options),
114
2.33k
      mutable_cf_options_(mutable_cf_options),
115
2.33k
      max_memtable_id_(max_memtable_id),
116
2.33k
      file_options_(file_options),
117
2.33k
      versions_(versions),
118
2.33k
      db_mutex_(db_mutex),
119
2.33k
      shutting_down_(shutting_down),
120
2.33k
      earliest_snapshot_(job_context->GetEarliestSnapshotSequence()),
121
2.33k
      job_context_(job_context),
122
2.33k
      flush_reason_(flush_reason),
123
2.33k
      log_buffer_(log_buffer),
124
2.33k
      db_directory_(db_directory),
125
2.33k
      output_file_directory_(output_file_directory),
126
2.33k
      output_compression_(output_compression),
127
2.33k
      stats_(stats),
128
2.33k
      event_logger_(event_logger),
129
2.33k
      measure_io_stats_(measure_io_stats),
130
2.33k
      sync_output_directory_(sync_output_directory),
131
2.33k
      write_manifest_(write_manifest),
132
2.33k
      edit_(nullptr),
133
2.33k
      base_(nullptr),
134
2.33k
      pick_memtable_called(false),
135
2.33k
      thread_pri_(thread_pri),
136
2.33k
      io_tracer_(io_tracer),
137
2.33k
      clock_(db_options_.clock),
138
2.33k
      full_history_ts_low_(std::move(full_history_ts_low)),
139
2.33k
      blob_callback_(blob_callback),
140
2.33k
      fast_sst_open_(fast_sst_open),
141
2.33k
      seqno_to_time_mapping_(std::move(seqno_to_time_mapping)) {
142
2.33k
  assert(job_context->snapshot_context_initialized);
143
  // Update the thread status to indicate flush.
144
2.33k
  ReportStartedFlush();
145
2.33k
  TEST_SYNC_POINT("FlushJob::FlushJob()");
146
2.33k
}
147
148
2.33k
FlushJob::~FlushJob() { ThreadStatusUtil::ResetThreadStatus(); }
149
150
2.33k
void FlushJob::ReportStartedFlush() {
151
2.33k
  ThreadStatusUtil::SetEnableTracking(db_options_.enable_thread_tracking);
152
2.33k
  ThreadStatusUtil::SetColumnFamily(cfd_);
153
2.33k
  ThreadStatusUtil::SetThreadOperation(ThreadStatus::OP_FLUSH);
154
2.33k
  ThreadStatusUtil::SetThreadOperationProperty(ThreadStatus::COMPACTION_JOB_ID,
155
2.33k
                                               job_context_->job_id);
156
157
2.33k
  IOSTATS_RESET(bytes_written);
158
2.33k
}
159
160
2.33k
void FlushJob::ReportFlushInputSize(const autovector<ReadOnlyMemTable*>& mems) {
161
2.33k
  uint64_t input_size = 0;
162
2.33k
  for (auto* mem : mems) {
163
2.33k
    input_size += mem->ApproximateMemoryUsage();
164
2.33k
  }
165
2.33k
  ThreadStatusUtil::IncreaseThreadOperationProperty(
166
2.33k
      ThreadStatus::FLUSH_BYTES_MEMTABLES, input_size);
167
2.33k
}
168
169
4.67k
void FlushJob::RecordFlushIOStats() {
170
4.67k
  RecordTick(stats_, FLUSH_WRITE_BYTES, IOSTATS(bytes_written));
171
4.67k
  ThreadStatusUtil::IncreaseThreadOperationProperty(
172
4.67k
      ThreadStatus::FLUSH_BYTES_WRITTEN, IOSTATS(bytes_written));
173
4.67k
  IOSTATS_RESET(bytes_written);
174
4.67k
}
175
2.33k
void FlushJob::PickMemTable() {
176
2.33k
  db_mutex_->AssertHeld();
177
2.33k
  assert(!pick_memtable_called);
178
2.33k
  pick_memtable_called = true;
179
180
  // Maximum "NextLogNumber" of the memtables to flush.
181
  // When mempurge feature is turned off, this variable is useless
182
  // because the memtables are implicitly sorted by increasing order of creation
183
  // time. Therefore mems_->back()->GetNextLogNumber() is already equal to
184
  // max_next_log_number. However when Mempurge is on, the memtables are no
185
  // longer sorted by increasing order of creation time. Therefore this variable
186
  // becomes necessary because mems_->back()->GetNextLogNumber() is no longer
187
  // necessarily equal to max_next_log_number.
188
2.33k
  uint64_t max_next_log_number = 0;
189
190
  // Save the contents of the earliest memtable as a new Table
191
2.33k
  cfd_->imm()->PickMemtablesToFlush(max_memtable_id_, &mems_,
192
2.33k
                                    &max_next_log_number);
193
2.33k
  if (mems_.empty()) {
194
0
    return;
195
0
  }
196
197
  // Track effective cutoff user-defined timestamp during flush if
198
  // user-defined timestamps can be stripped.
199
2.33k
  GetEffectiveCutoffUDTForPickedMemTables();
200
2.33k
  GetPrecludeLastLevelMinSeqno();
201
202
2.33k
  ReportFlushInputSize(mems_);
203
204
  // entries mems are (implicitly) sorted in ascending order by their created
205
  // time. We will use the first memtable's `edit` to keep the meta info for
206
  // this flush.
207
2.33k
  ReadOnlyMemTable* m = mems_[0];
208
2.33k
  edit_ = m->GetEdits();
209
2.33k
  edit_->SetPrevLogNumber(0);
210
  // SetLogNumber(log_num) indicates logs with number smaller than log_num
211
  // will no longer be picked up for recovery.
212
2.33k
  edit_->SetLogNumber(max_next_log_number);
213
2.33k
  edit_->SetColumnFamily(cfd_->GetID());
214
215
  // path 0 for level 0 file.
216
2.33k
  meta_.fd = FileDescriptor(versions_->NewFileNumber(), 0, 0);
217
2.33k
  meta_.epoch_number = cfd_->NewEpochNumber();
218
219
2.33k
  base_ = cfd_->current();
220
2.33k
  base_->Ref();  // it is likely that we do not need this reference
221
2.33k
}
222
223
Status FlushJob::Run(LogsWithPrepTracker* prep_tracker, FileMetaData* file_meta,
224
                     bool* switched_to_mempurge, bool* skipped_since_bg_error,
225
2.33k
                     ErrorHandler* error_handler) {
226
2.33k
  TEST_SYNC_POINT("FlushJob::Start");
227
2.33k
  db_mutex_->AssertHeld();
228
2.33k
  assert(pick_memtable_called);
229
  // Mempurge threshold can be dynamically changed.
230
  // For sake of consistency, mempurge_threshold is
231
  // saved locally to maintain consistency in each
232
  // FlushJob::Run call.
233
2.33k
  double mempurge_threshold =
234
2.33k
      mutable_cf_options_.experimental_mempurge_threshold;
235
236
2.33k
  AutoThreadOperationStageUpdater stage_run(ThreadStatus::STAGE_FLUSH_RUN);
237
2.33k
  if (mems_.empty()) {
238
0
    ROCKS_LOG_BUFFER(log_buffer_, "[%s] No memtable to flush",
239
0
                     cfd_->GetName().c_str());
240
0
    return Status::OK();
241
0
  }
242
243
  // I/O measurement variables
244
2.33k
  PerfLevel prev_perf_level = PerfLevel::kEnableTime;
245
2.33k
  uint64_t prev_write_nanos = 0;
246
2.33k
  uint64_t prev_fsync_nanos = 0;
247
2.33k
  uint64_t prev_range_sync_nanos = 0;
248
2.33k
  uint64_t prev_prepare_write_nanos = 0;
249
2.33k
  uint64_t prev_cpu_write_nanos = 0;
250
2.33k
  uint64_t prev_cpu_read_nanos = 0;
251
2.33k
  if (measure_io_stats_) {
252
0
    prev_perf_level = GetPerfLevel();
253
0
    SetPerfLevel(PerfLevel::kEnableTime);
254
0
    prev_write_nanos = IOSTATS(write_nanos);
255
0
    prev_fsync_nanos = IOSTATS(fsync_nanos);
256
0
    prev_range_sync_nanos = IOSTATS(range_sync_nanos);
257
0
    prev_prepare_write_nanos = IOSTATS(prepare_write_nanos);
258
0
    prev_cpu_write_nanos = IOSTATS(cpu_write_nanos);
259
0
    prev_cpu_read_nanos = IOSTATS(cpu_read_nanos);
260
0
  }
261
2.33k
  Status mempurge_s = Status::NotFound("No MemPurge.");
262
2.33k
  if ((mempurge_threshold > 0.0) &&
263
0
      (flush_reason_ == FlushReason::kWriteBufferFull) && (!mems_.empty()) &&
264
0
      MemPurgeDecider(mempurge_threshold) && !(db_options_.atomic_flush)) {
265
0
    cfd_->SetMempurgeUsed();
266
0
    mempurge_s = MemPurge();
267
0
    if (!mempurge_s.ok()) {
268
      // Mempurge is typically aborted when the output
269
      // bytes cannot be contained onto a single output memtable.
270
0
      if (mempurge_s.IsAborted()) {
271
0
        ROCKS_LOG_INFO(db_options_.info_log, "Mempurge process aborted: %s\n",
272
0
                       mempurge_s.ToString().c_str());
273
0
      } else {
274
        // However the mempurge process can also fail for
275
        // other reasons (eg: new_mem->Add() fails).
276
0
        ROCKS_LOG_WARN(db_options_.info_log, "Mempurge process failed: %s\n",
277
0
                       mempurge_s.ToString().c_str());
278
0
      }
279
0
    } else {
280
0
      if (switched_to_mempurge) {
281
0
        *switched_to_mempurge = true;
282
0
      } else {
283
        // The mempurge process was successful, but no switch_to_mempurge
284
        // pointer provided so no way to propagate the state of flush job.
285
0
        ROCKS_LOG_WARN(db_options_.info_log,
286
0
                       "Mempurge process succeeded"
287
0
                       "but no 'switched_to_mempurge' ptr provided.\n");
288
0
      }
289
0
    }
290
0
  }
291
2.33k
  Status s;
292
2.33k
  if (mempurge_s.ok()) {
293
0
    base_->Unref();
294
0
    s = Status::OK();
295
2.33k
  } else {
296
    // This will release and re-acquire the mutex.
297
2.33k
    s = WriteLevel0Table();
298
2.33k
  }
299
300
2.33k
  if (s.ok() && cfd_->IsDropped()) {
301
0
    s = Status::ColumnFamilyDropped("Column family dropped during compaction");
302
0
  }
303
2.33k
  if ((s.ok() || s.IsColumnFamilyDropped()) &&
304
2.33k
      shutting_down_->load(std::memory_order_acquire)) {
305
0
    s = Status::ShutdownInProgress("Database shutdown");
306
0
  }
307
308
2.33k
  if (s.ok()) {
309
2.33k
    s = MaybeIncreaseFullHistoryTsLowToAboveCutoffUDT();
310
2.33k
  }
311
312
2.33k
  TEST_SYNC_POINT_CALLBACK("FlushJob::Run:PostBuildTable", &s);
313
314
2.33k
  if (!s.ok()) {
315
0
    cfd_->imm()->RollbackMemtableFlush(
316
0
        mems_, /*rollback_succeeding_memtables=*/!db_options_.atomic_flush);
317
2.33k
  } else if (write_manifest_) {
318
2.33k
    assert(!db_options_.atomic_flush);
319
2.33k
    if (!db_options_.atomic_flush &&
320
2.33k
        flush_reason_ != FlushReason::kErrorRecovery &&
321
2.33k
        flush_reason_ != FlushReason::kErrorRecoveryRetryFlush &&
322
2.33k
        error_handler && !error_handler->GetBGError().ok() &&
323
0
        error_handler->IsBGWorkStopped()) {
324
0
      cfd_->imm()->RollbackMemtableFlush(
325
0
          mems_, /*rollback_succeeding_memtables=*/!db_options_.atomic_flush);
326
0
      s = error_handler->GetBGError();
327
0
      if (skipped_since_bg_error) {
328
0
        *skipped_since_bg_error = true;
329
0
      }
330
2.33k
    } else {
331
2.33k
      TEST_SYNC_POINT("FlushJob::InstallResults");
332
      // Replace immutable memtable with the generated Table
333
2.33k
      s = cfd_->imm()->TryInstallMemtableFlushResults(
334
2.33k
              cfd_, mems_, prep_tracker, versions_, db_mutex_,
335
2.33k
              meta_.fd.GetNumber(), &job_context_->memtables_to_free, db_directory_,
336
2.33k
              log_buffer_, &committed_flush_jobs_info_,
337
2.33k
              !(mempurge_s.ok()) /* write_edit : true if no mempurge happened (or if aborted),
338
                              but 'false' if mempurge successful: no new min log number
339
2.33k
                              or new level 0 file path to write to manifest. */);
340
2.33k
    }
341
2.33k
  }
342
343
2.33k
  if (!s.ok() && meta_.fd.GetFileSize() > 0) {
344
    // If BuildTable succeeded (file was cached in table cache for user reads)
345
    // but the flush failed to install (e.g., LogAndApply MANIFEST I/O error,
346
    // CF dropped, shutdown), evict the cached entry. Without this, the file
347
    // is cached but not in any Version, and the FindObsoleteFiles backstop
348
    // can fail under metadata read fault injection, causing
349
    // TEST_VerifyNoObsoleteFilesCached to fire during Close().
350
0
    TableCache::ReleaseObsolete(cfd_->table_cache()->get_cache().get(),
351
0
                                meta_.fd.GetNumber(), nullptr /*handle*/,
352
0
                                mutable_cf_options_.uncache_aggressiveness);
353
0
  }
354
355
2.33k
  if (s.ok() && file_meta != nullptr) {
356
2.33k
    *file_meta = meta_;
357
2.33k
  }
358
2.33k
  RecordFlushIOStats();
359
360
  // When measure_io_stats_ is true, the default 512 bytes is not enough.
361
2.33k
  auto stream = event_logger_->LogToBuffer(log_buffer_, 1024);
362
2.33k
  stream << "job" << job_context_->job_id << "event" << "flush_finished";
363
2.33k
  stream << "output_compression"
364
2.33k
         << CompressionTypeToString(output_compression_);
365
2.33k
  stream << "lsm_state";
366
2.33k
  stream.StartArray();
367
2.33k
  auto vstorage = cfd_->current()->storage_info();
368
18.7k
  for (int level = 0; level < vstorage->num_levels(); ++level) {
369
16.3k
    stream << vstorage->NumLevelFiles(level);
370
16.3k
  }
371
2.33k
  stream.EndArray();
372
373
2.33k
  const auto& blob_files = vstorage->GetBlobFiles();
374
2.33k
  if (!blob_files.empty()) {
375
0
    assert(blob_files.front());
376
0
    stream << "blob_file_head" << blob_files.front()->GetBlobFileNumber();
377
378
0
    assert(blob_files.back());
379
0
    stream << "blob_file_tail" << blob_files.back()->GetBlobFileNumber();
380
0
  }
381
382
2.33k
  stream << "immutable_memtables" << cfd_->imm()->NumNotFlushed();
383
384
2.33k
  if (measure_io_stats_) {
385
0
    if (prev_perf_level != PerfLevel::kEnableTime) {
386
0
      SetPerfLevel(prev_perf_level);
387
0
    }
388
0
    stream << "file_write_nanos" << (IOSTATS(write_nanos) - prev_write_nanos);
389
0
    stream << "file_range_sync_nanos"
390
0
           << (IOSTATS(range_sync_nanos) - prev_range_sync_nanos);
391
0
    stream << "file_fsync_nanos" << (IOSTATS(fsync_nanos) - prev_fsync_nanos);
392
0
    stream << "file_prepare_write_nanos"
393
0
           << (IOSTATS(prepare_write_nanos) - prev_prepare_write_nanos);
394
0
    stream << "file_cpu_write_nanos"
395
0
           << (IOSTATS(cpu_write_nanos) - prev_cpu_write_nanos);
396
0
    stream << "file_cpu_read_nanos"
397
0
           << (IOSTATS(cpu_read_nanos) - prev_cpu_read_nanos);
398
0
  }
399
400
2.33k
  TEST_SYNC_POINT("FlushJob::End");
401
2.33k
  return s;
402
2.33k
}
403
404
0
void FlushJob::Cancel() {
405
0
  db_mutex_->AssertHeld();
406
0
  assert(base_ != nullptr);
407
0
  base_->Unref();
408
0
}
409
410
0
Status FlushJob::MemPurge() {
411
0
  Status s;
412
0
  db_mutex_->AssertHeld();
413
0
  db_mutex_->Unlock();
414
0
  assert(!mems_.empty());
415
416
  // Measure purging time.
417
0
  const uint64_t start_micros = clock_->NowMicros();
418
0
  const uint64_t start_cpu_micros = clock_->CPUMicros();
419
420
0
  MemTable* new_mem = nullptr;
421
  // For performance/log investigation purposes:
422
  // look at how much useful payload we harvest in the new_mem.
423
  // This value is then printed to the DB log.
424
0
  double new_mem_capacity = 0.0;
425
426
  // Create two iterators, one for the memtable data (contains
427
  // info from puts + deletes), and one for the memtable
428
  // Range Tombstones (from DeleteRanges).
429
  // TODO: plumb Env::IOActivity, Env::IOPriority
430
0
  ReadOptions ro;
431
0
  ro.total_order_seek = true;
432
0
  Arena arena;
433
0
  std::vector<InternalIterator*> memtables;
434
0
  std::vector<std::unique_ptr<FragmentedRangeTombstoneIterator>>
435
0
      range_del_iters;
436
0
  for (ReadOnlyMemTable* m : mems_) {
437
0
    memtables.push_back(m->NewIterator(ro, /*seqno_to_time_mapping=*/nullptr,
438
0
                                       &arena, /*prefix_extractor=*/nullptr,
439
0
                                       /*for_flush=*/true));
440
0
    auto* range_del_iter = m->NewRangeTombstoneIterator(
441
0
        ro, kMaxSequenceNumber, true /* immutable_memtable */);
442
0
    if (range_del_iter != nullptr) {
443
0
      range_del_iters.emplace_back(range_del_iter);
444
0
    }
445
0
  }
446
447
0
  assert(!memtables.empty());
448
0
  SequenceNumber first_seqno = kMaxSequenceNumber;
449
0
  SequenceNumber earliest_seqno = kMaxSequenceNumber;
450
  // Pick first and earliest seqno as min of all first_seqno
451
  // and earliest_seqno of the mempurged memtables.
452
0
  for (const auto& mem : mems_) {
453
0
    first_seqno = mem->GetFirstSequenceNumber() < first_seqno
454
0
                      ? mem->GetFirstSequenceNumber()
455
0
                      : first_seqno;
456
0
    earliest_seqno = mem->GetEarliestSequenceNumber() < earliest_seqno
457
0
                         ? mem->GetEarliestSequenceNumber()
458
0
                         : earliest_seqno;
459
0
  }
460
461
0
  ScopedArenaPtr<InternalIterator> iter(
462
0
      NewMergingIterator(&(cfd_->internal_comparator()), memtables.data(),
463
0
                         static_cast<int>(memtables.size()), &arena));
464
465
0
  const auto& ioptions = cfd_->ioptions();
466
467
  // Place iterator at the First (meaning most recent) key node.
468
0
  iter->SeekToFirst();
469
470
0
  const std::string* const full_history_ts_low = &(cfd_->GetFullHistoryTsLow());
471
0
  std::unique_ptr<CompactionRangeDelAggregator> range_del_agg(
472
0
      new CompactionRangeDelAggregator(&(cfd_->internal_comparator()),
473
0
                                       job_context_->snapshot_seqs,
474
0
                                       full_history_ts_low));
475
0
  for (auto& rd_iter : range_del_iters) {
476
0
    range_del_agg->AddTombstones(std::move(rd_iter));
477
0
  }
478
479
  // If there is valid data in the memtable,
480
  // or at least range tombstones, copy over the info
481
  // to the new memtable.
482
0
  if (iter->Valid() || !range_del_agg->IsEmpty()) {
483
    // MaxSize is the size of a memtable.
484
0
    size_t maxSize = mutable_cf_options_.write_buffer_size;
485
0
    std::unique_ptr<CompactionFilter> compaction_filter;
486
0
    if (ioptions.compaction_filter_factory != nullptr &&
487
0
        ioptions.compaction_filter_factory->ShouldFilterTableFileCreation(
488
0
            TableFileCreationReason::kFlush)) {
489
0
      CompactionFilter::Context ctx;
490
0
      ctx.is_full_compaction = false;
491
0
      ctx.is_manual_compaction = false;
492
0
      ctx.column_family_id = cfd_->GetID();
493
0
      ctx.reason = TableFileCreationReason::kFlush;
494
0
      compaction_filter =
495
0
          ioptions.compaction_filter_factory->CreateCompactionFilter(ctx);
496
0
      if (compaction_filter != nullptr &&
497
0
          !compaction_filter->IgnoreSnapshots()) {
498
0
        s = Status::NotSupported(
499
0
            "CompactionFilter::IgnoreSnapshots() = false is not supported "
500
0
            "anymore.");
501
0
        return s;
502
0
      }
503
0
    }
504
505
0
    new_mem = new MemTable(cfd_->internal_comparator(), cfd_->ioptions(),
506
0
                           mutable_cf_options_, cfd_->write_buffer_mgr(),
507
0
                           earliest_seqno, cfd_->GetID());
508
0
    assert(new_mem != nullptr);
509
510
0
    Env* env = db_options_.env;
511
0
    assert(env);
512
0
    MergeHelper merge(env, (cfd_->internal_comparator()).user_comparator(),
513
0
                      (ioptions.merge_operator).get(), compaction_filter.get(),
514
0
                      ioptions.logger,
515
0
                      true /* internal key corruption is not ok */,
516
0
                      job_context_->GetLatestSnapshotSequence(),
517
0
                      job_context_->snapshot_checker);
518
0
    assert(job_context_);
519
0
    const std::atomic<bool> kManualCompactionCanceledFalse{false};
520
0
    CompactionIterator c_iter(
521
0
        iter.get(), (cfd_->internal_comparator()).user_comparator(), &merge,
522
0
        kMaxSequenceNumber, &job_context_->snapshot_seqs, earliest_snapshot_,
523
0
        job_context_->earliest_write_conflict_snapshot,
524
0
        job_context_->GetJobSnapshotSequence(), job_context_->snapshot_checker,
525
0
        env, ShouldReportDetailedTime(env, ioptions.stats), range_del_agg.get(),
526
0
        nullptr, ioptions.allow_data_in_errors,
527
0
        ioptions.enforce_single_del_contracts,
528
0
        /*manual_compaction_canceled=*/kManualCompactionCanceledFalse,
529
0
        false /* must_count_input_entries */,
530
0
        /*compaction=*/nullptr, compaction_filter.get(),
531
0
        /*shutting_down=*/nullptr, ioptions.info_log, full_history_ts_low);
532
533
    // Set earliest sequence number in the new memtable
534
    // to be equal to the earliest sequence number of the
535
    // memtable being flushed (See later if there is a need
536
    // to update this number!).
537
0
    new_mem->SetEarliestSequenceNumber(earliest_seqno);
538
    // Likewise for first seq number.
539
0
    new_mem->SetFirstSequenceNumber(first_seqno);
540
0
    SequenceNumber new_first_seqno = kMaxSequenceNumber;
541
542
0
    c_iter.SeekToFirst();
543
544
    // Key transfer
545
0
    for (; c_iter.Valid(); c_iter.Next()) {
546
0
      const ParsedInternalKey ikey = c_iter.ikey();
547
0
      const Slice value = c_iter.value();
548
0
      new_first_seqno =
549
0
          ikey.sequence < new_first_seqno ? ikey.sequence : new_first_seqno;
550
551
      // Should we update "OldestKeyTime" ???? -> timestamp appear
552
      // to still be an "experimental" feature.
553
0
      s = new_mem->Add(
554
0
          ikey.sequence, ikey.type, ikey.user_key, value,
555
0
          nullptr,   // KV protection info set as nullptr since it
556
                     // should only be useful for the first add to
557
                     // the original memtable.
558
0
          false,     // : allow concurrent_memtable_writes_
559
                     // Not seen as necessary for now.
560
0
          nullptr,   // get_post_process_info(m) must be nullptr
561
                     // when concurrent_memtable_writes is switched off.
562
0
          nullptr);  // hint, only used when concurrent_memtable_writes_
563
                     // is switched on.
564
0
      if (!s.ok()) {
565
0
        break;
566
0
      }
567
568
      // If new_mem has size greater than maxSize,
569
      // then rollback to regular flush operation,
570
      // and destroy new_mem.
571
0
      if (new_mem->ApproximateMemoryUsage() > maxSize) {
572
0
        s = Status::Aborted("Mempurge filled more than one memtable.");
573
0
        new_mem_capacity = 1.0;
574
0
        break;
575
0
      }
576
0
    }
577
578
    // Check status and propagate
579
    // potential error status from c_iter
580
0
    if (!s.ok()) {
581
0
      c_iter.status().PermitUncheckedError();
582
0
    } else if (!c_iter.status().ok()) {
583
0
      s = c_iter.status();
584
0
    }
585
586
    // Range tombstone transfer.
587
0
    if (s.ok()) {
588
0
      auto range_del_it = range_del_agg->NewIterator();
589
0
      for (range_del_it->SeekToFirst(); range_del_it->Valid();
590
0
           range_del_it->Next()) {
591
0
        auto tombstone = range_del_it->Tombstone();
592
0
        new_first_seqno =
593
0
            tombstone.seq_ < new_first_seqno ? tombstone.seq_ : new_first_seqno;
594
0
        MemTablePostProcessInfo post_process_info;
595
0
        s = new_mem->Add(
596
0
            tombstone.seq_,        // Sequence number
597
0
            kTypeRangeDeletion,    // KV type
598
0
            tombstone.start_key_,  // Key is start key.
599
0
            tombstone.end_key_,    // Value is end key.
600
0
            nullptr,               // KV protection info set as nullptr since it
601
                                   // should only be useful for the first add to
602
                                   // the original memtable.
603
0
            true,                  // allow_concurrent: required for range
604
                                   // deletions since range_del_table_ uses
605
                                   // concurrent-safe SkipList.
606
0
            &post_process_info,
607
0
            nullptr);  // hint, only used when concurrent_memtable_writes_
608
                       // is switched on.
609
0
        if (s.ok()) {
610
0
          new_mem->BatchPostProcess(post_process_info);
611
0
        }
612
613
0
        if (!s.ok()) {
614
0
          break;
615
0
        }
616
617
        // If new_mem has size greater than maxSize,
618
        // then rollback to regular flush operation,
619
        // and destroy new_mem.
620
0
        if (new_mem->ApproximateMemoryUsage() > maxSize) {
621
0
          s = Status::Aborted(Slice("Mempurge filled more than one memtable."));
622
0
          new_mem_capacity = 1.0;
623
0
          break;
624
0
        }
625
0
      }
626
0
    }
627
628
    // If everything happened smoothly and new_mem contains valid data,
629
    // decide if it is flushed to storage or kept in the imm()
630
    // memtable list (memory).
631
0
    if (s.ok() && (new_first_seqno != kMaxSequenceNumber)) {
632
      // Rectify the first sequence number, which (unlike the earliest seq
633
      // number) needs to be present in the new memtable.
634
0
      new_mem->SetFirstSequenceNumber(new_first_seqno);
635
636
      // The new_mem is added to the list of immutable memtables
637
      // only if it filled at less than 100% capacity and isn't flagged
638
      // as in need of being flushed.
639
0
      if (new_mem->ApproximateMemoryUsage() < maxSize &&
640
0
          !(new_mem->ShouldFlushNow())) {
641
        // Construct fragmented memtable range tombstones without mutex
642
0
        new_mem->ConstructFragmentedRangeTombstones();
643
0
        TEST_SYNC_POINT("FlushJob::MemPurge:BeforeReacquireMutex");
644
0
        TEST_SYNC_POINT("FlushJob::MemPurge:AfterWaitForTest");
645
0
        db_mutex_->Lock();
646
        // Take the newest id, so that memtables in MemtableList don't have
647
        // out-of-order memtable ids. While the db mutex was released during
648
        // MemPurge, new memtables may have been switched to the immutable
649
        // list with higher IDs, so we must use the maximum of the original
650
        // flush batch ID and the current latest immutable memtable ID.
651
0
        uint64_t new_mem_id = std::max(
652
0
            mems_.back()->GetID(),
653
0
            cfd_->imm()->GetLatestMemTableID(false /*for_atomic_flush*/));
654
655
0
        new_mem->SetID(new_mem_id);
656
        // Take the latest memtable's next log number.
657
0
        new_mem->SetNextLogNumber(mems_.back()->GetNextLogNumber());
658
659
        // This addition will not trigger another flush, because
660
        // we do not call EnqueuePendingFlush().
661
0
        cfd_->imm()->Add(new_mem, &job_context_->memtables_to_free);
662
0
        new_mem->Ref();
663
        // Piggyback FlushJobInfo on the first flushed memtable.
664
0
        db_mutex_->AssertHeld();
665
0
        meta_.fd.file_size = 0;
666
0
        mems_[0]->SetFlushJobInfo(GetFlushJobInfo());
667
0
        db_mutex_->Unlock();
668
0
      } else {
669
0
        s = Status::Aborted(Slice("Mempurge filled more than one memtable."));
670
0
        new_mem_capacity = 1.0;
671
0
        if (new_mem) {
672
0
          job_context_->memtables_to_free.push_back(new_mem);
673
0
        }
674
0
      }
675
0
    } else {
676
      // In this case, the newly allocated new_mem is empty.
677
0
      assert(new_mem != nullptr);
678
0
      job_context_->memtables_to_free.push_back(new_mem);
679
0
    }
680
0
  }
681
682
  // Reacquire the mutex for WriteLevel0 function.
683
0
  db_mutex_->Lock();
684
685
  // If mempurge successful, don't write input tables to level0,
686
  // but write any full output table to level0.
687
0
  if (s.ok()) {
688
0
    TEST_SYNC_POINT("DBImpl::FlushJob:MemPurgeSuccessful");
689
0
  } else {
690
0
    TEST_SYNC_POINT("DBImpl::FlushJob:MemPurgeUnsuccessful");
691
0
  }
692
0
  const uint64_t micros = clock_->NowMicros() - start_micros;
693
0
  const uint64_t cpu_micros = clock_->CPUMicros() - start_cpu_micros;
694
0
  ROCKS_LOG_INFO(db_options_.info_log,
695
0
                 "[%s] [JOB %d] Mempurge lasted %" PRIu64
696
0
                 " microseconds, and %" PRIu64
697
0
                 " cpu "
698
0
                 "microseconds. Status is %s ok. Perc capacity: %f\n",
699
0
                 cfd_->GetName().c_str(), job_context_->job_id, micros,
700
0
                 cpu_micros, s.ok() ? "" : "not", new_mem_capacity);
701
702
0
  return s;
703
0
}
704
705
0
bool FlushJob::MemPurgeDecider(double threshold) {
706
  // Never trigger mempurge if threshold is not a strictly positive value.
707
0
  if (!(threshold > 0.0)) {
708
0
    return false;
709
0
  }
710
0
  if (threshold > (1.0 * mems_.size())) {
711
0
    return true;
712
0
  }
713
  // Payload and useful_payload (in bytes).
714
  // The useful payload ratio of a given MemTable
715
  // is estimated to be useful_payload/payload.
716
0
  uint64_t payload = 0, useful_payload = 0, entry_size = 0;
717
718
  // Local variables used repetitively inside the for-loop
719
  // when iterating over the sampled entries.
720
0
  Slice key_slice, value_slice;
721
0
  ParsedInternalKey res;
722
0
  SnapshotImpl min_snapshot;
723
0
  std::string vget;
724
0
  Status mget_s, parse_s;
725
0
  MergeContext merge_context;
726
0
  SequenceNumber max_covering_tombstone_seq = 0, sqno = 0,
727
0
                 min_seqno_snapshot = 0;
728
0
  bool get_res, can_be_useful_payload, not_in_next_mems;
729
730
  // If estimated_useful_payload is > threshold,
731
  // then flush to storage, else MemPurge.
732
0
  double estimated_useful_payload = 0.0;
733
  // Cochran formula for determining sample size.
734
  // 95% confidence interval, 7% precision.
735
  //    n0 = (1.96*1.96)*0.25/(0.07*0.07) = 196.0
736
0
  double n0 = 196.0;
737
  // TODO: plumb Env::IOActivity, Env::IOPriority
738
0
  ReadOptions ro;
739
0
  ro.total_order_seek = true;
740
741
  // Iterate over each memtable of the set.
742
0
  for (auto mem_iter = std::begin(mems_); mem_iter != std::end(mems_);
743
0
       ++mem_iter) {
744
0
    ReadOnlyMemTable* mt = *mem_iter;
745
746
    // Else sample from the table.
747
0
    uint64_t nentries = mt->NumEntries();
748
    // Corrected Cochran formula for small populations
749
    // (converges to n0 for large populations).
750
0
    uint64_t target_sample_size =
751
0
        static_cast<uint64_t>(ceil(n0 / (1.0 + (n0 / nentries))));
752
0
    std::unordered_set<const char*> sentries = {};
753
    // Populate sample entries set.
754
0
    mt->UniqueRandomSample(target_sample_size, &sentries);
755
756
    // Estimate the garbage ratio by comparing if
757
    // each sample corresponds to a valid entry.
758
0
    for (const char* ss : sentries) {
759
0
      key_slice = GetLengthPrefixedSlice(ss);
760
0
      parse_s = ParseInternalKey(key_slice, &res, true /*log_err_key*/);
761
0
      if (!parse_s.ok()) {
762
0
        ROCKS_LOG_WARN(db_options_.info_log,
763
0
                       "Memtable Decider: ParseInternalKey did not parse "
764
0
                       "key_slice %s successfully.",
765
0
                       key_slice.data());
766
0
      }
767
768
      // Size of the entry is "key size (+ value size if KV entry)"
769
0
      entry_size = key_slice.size();
770
0
      if (res.type == kTypeValue) {
771
0
        value_slice =
772
0
            GetLengthPrefixedSlice(key_slice.data() + key_slice.size());
773
0
        entry_size += value_slice.size();
774
0
      }
775
776
      // Count entry bytes as payload.
777
0
      payload += entry_size;
778
779
0
      LookupKey lkey(res.user_key, kMaxSequenceNumber);
780
781
      // Paranoia: zero out these values just in case.
782
0
      max_covering_tombstone_seq = 0;
783
0
      sqno = 0;
784
785
      // Pick the oldest existing snapshot that is more recent
786
      // than the sequence number of the sampled entry.
787
0
      min_seqno_snapshot = kMaxSequenceNumber;
788
0
      for (SequenceNumber seq_num : job_context_->snapshot_seqs) {
789
0
        if (seq_num > res.sequence && seq_num < min_seqno_snapshot) {
790
0
          min_seqno_snapshot = seq_num;
791
0
        }
792
0
      }
793
0
      min_snapshot.number_ = min_seqno_snapshot;
794
0
      ro.snapshot =
795
0
          min_seqno_snapshot < kMaxSequenceNumber ? &min_snapshot : nullptr;
796
797
      // Estimate if the sample entry is valid or not.
798
0
      get_res = mt->Get(lkey, &vget, /*columns=*/nullptr, /*timestamp=*/nullptr,
799
0
                        &mget_s, &merge_context, &max_covering_tombstone_seq,
800
0
                        &sqno, ro, true /* immutable_memtable */);
801
0
      if (!get_res) {
802
0
        ROCKS_LOG_WARN(
803
0
            db_options_.info_log,
804
0
            "Memtable Get returned false when Get(sampled entry). "
805
0
            "Yet each sample entry should exist somewhere in the memtable, "
806
0
            "unrelated to whether it has been deleted or not.");
807
0
      }
808
809
      // TODO(bjlemaire): evaluate typeMerge.
810
      // This is where the sampled entry is estimated to be
811
      // garbage or not. Note that this is a garbage *estimation*
812
      // because we do not include certain items such as
813
      // CompactionFitlers triggered at flush, or if the same delete
814
      // has been inserted twice or more in the memtable.
815
816
      // Evaluate if the entry can be useful payload
817
      // Situation #1: entry is a KV entry, was found in the memtable mt
818
      //               and the sequence numbers match.
819
0
      can_be_useful_payload = (res.type == kTypeValue) && get_res &&
820
0
                              mget_s.ok() && (sqno == res.sequence);
821
822
      // Situation #2: entry is a delete entry, was found in the memtable mt
823
      //               (because gres==true) and no valid KV entry is found.
824
      //               (note: duplicate delete entries are also taken into
825
      //               account here, because the sequence number 'sqno'
826
      //               in memtable->Get(&sqno) operation is set to be equal
827
      //               to the most recent delete entry as well).
828
0
      can_be_useful_payload |=
829
0
          ((res.type == kTypeDeletion) || (res.type == kTypeSingleDeletion)) &&
830
0
          mget_s.IsNotFound() && get_res && (sqno == res.sequence);
831
832
      // If there is a chance that the entry is useful payload
833
      // Verify that the entry does not appear in the following memtables
834
      // (memtables with greater memtable ID/larger sequence numbers).
835
0
      if (can_be_useful_payload) {
836
0
        not_in_next_mems = true;
837
0
        for (auto next_mem_iter = mem_iter + 1;
838
0
             next_mem_iter != std::end(mems_); next_mem_iter++) {
839
0
          if ((*next_mem_iter)
840
0
                  ->Get(lkey, &vget, /*columns=*/nullptr, /*timestamp=*/nullptr,
841
0
                        &mget_s, &merge_context, &max_covering_tombstone_seq,
842
0
                        &sqno, ro, true /* immutable_memtable */)) {
843
0
            not_in_next_mems = false;
844
0
            break;
845
0
          }
846
0
        }
847
0
        if (not_in_next_mems) {
848
0
          useful_payload += entry_size;
849
0
        }
850
0
      }
851
0
    }
852
0
    if (payload > 0) {
853
      // We use the estimated useful payload ratio to
854
      // evaluate how many of the memtable bytes are useful bytes.
855
0
      estimated_useful_payload +=
856
0
          (mt->ApproximateMemoryUsage()) * (useful_payload * 1.0 / payload);
857
858
0
      ROCKS_LOG_INFO(db_options_.info_log,
859
0
                     "Mempurge sampling [CF %s] - found garbage ratio from "
860
0
                     "sampling: %f. Threshold is %f\n",
861
0
                     cfd_->GetName().c_str(),
862
0
                     (payload - useful_payload) * 1.0 / payload, threshold);
863
0
    } else {
864
0
      ROCKS_LOG_WARN(db_options_.info_log,
865
0
                     "Mempurge sampling: null payload measured, and collected "
866
0
                     "sample size is %zu\n.",
867
0
                     sentries.size());
868
0
    }
869
0
  }
870
  // We convert the total number of useful payload bytes
871
  // into the proportion of memtable necessary to store all these bytes.
872
  // We compare this proportion with the threshold value.
873
0
  return ((estimated_useful_payload / mutable_cf_options_.write_buffer_size) <
874
0
          threshold);
875
0
}
876
877
2.33k
Status FlushJob::WriteLevel0Table() {
878
2.33k
  AutoThreadOperationStageUpdater stage_updater(
879
2.33k
      ThreadStatus::STAGE_FLUSH_WRITE_L0);
880
2.33k
  db_mutex_->AssertHeld();
881
2.33k
  const uint64_t start_micros = clock_->NowMicros();
882
2.33k
  const uint64_t start_cpu_micros = clock_->CPUMicros();
883
2.33k
  Status s;
884
885
2.33k
  meta_.temperature = mutable_cf_options_.default_write_temperature;
886
2.33k
  file_options_.temperature = meta_.temperature;
887
888
2.33k
  const auto* ucmp = cfd_->internal_comparator().user_comparator();
889
2.33k
  assert(ucmp);
890
2.33k
  const size_t ts_sz = ucmp->timestamp_size();
891
2.33k
  const bool logical_strip_timestamp =
892
2.33k
      ts_sz > 0 && !cfd_->ioptions().persist_user_defined_timestamps;
893
894
2.33k
  std::vector<BlobFileAddition> blob_file_additions;
895
2.33k
  std::vector<BlobFileGarbage> blob_file_garbages;
896
  // Only direct-write memtables can carry blob indexes that reference
897
  // pre-existing blob files. Plain flushes write new blob files from inline
898
  // values, so there is no pre-existing blob garbage to meter on the input
899
  // side.
900
2.33k
  std::vector<BlobFileGarbage>* const blob_file_garbages_for_filtering =
901
2.33k
      cfd_->blob_partition_manager() != nullptr ? &blob_file_garbages : nullptr;
902
  // Note that here we treat flush as level 0 compaction in internal stats
903
2.33k
  InternalStats::CompactionStats flush_stats(CompactionReason::kFlush,
904
2.33k
                                             1 /* count**/);
905
2.33k
  {
906
2.33k
    auto write_hint = base_->storage_info()->CalculateSSTWriteHint(
907
2.33k
        /*level=*/0, db_options_.calculate_sst_write_lifetime_hint_set);
908
2.33k
    Env::IOPriority io_priority = GetRateLimiterPriority();
909
2.33k
    db_mutex_->Unlock();
910
2.33k
    if (log_buffer_) {
911
2.33k
      log_buffer_->FlushBufferToLog();
912
2.33k
    }
913
    // memtables and range_del_iters store internal iterators over each data
914
    // memtable and its associated range deletion memtable, respectively, at
915
    // corresponding indexes.
916
2.33k
    std::vector<InternalIterator*> memtables;
917
2.33k
    std::vector<std::unique_ptr<FragmentedRangeTombstoneIterator>>
918
2.33k
        range_del_iters;
919
2.33k
    ReadOptions ro;
920
2.33k
    ro.total_order_seek = true;
921
2.33k
    ro.io_activity = Env::IOActivity::kFlush;
922
2.33k
    Arena arena;
923
2.33k
    uint64_t total_num_input_entries = 0, total_num_deletes = 0;
924
2.33k
    uint64_t total_data_size = 0;
925
2.33k
    size_t total_memory_usage = 0;
926
2.33k
    uint64_t total_num_range_deletes = 0;
927
    // Used for testing:
928
2.33k
    uint64_t mems_size = mems_.size();
929
2.33k
    (void)mems_size;  // avoids unused variable error when
930
                      // TEST_SYNC_POINT_CALLBACK not used.
931
2.33k
    TEST_SYNC_POINT_CALLBACK("FlushJob::WriteLevel0Table:num_memtables",
932
2.33k
                             &mems_size);
933
2.33k
    assert(job_context_);
934
2.33k
    for (ReadOnlyMemTable* m : mems_) {
935
2.33k
      ROCKS_LOG_INFO(db_options_.info_log,
936
2.33k
                     "[%s] [JOB %d] Flushing memtable id %" PRIu64
937
2.33k
                     " with next log file: %" PRIu64 ", marked_for_flush: %d\n",
938
2.33k
                     cfd_->GetName().c_str(), job_context_->job_id, m->GetID(),
939
2.33k
                     m->GetNextLogNumber(), m->IsMarkedForFlush());
940
2.33k
      if (logical_strip_timestamp) {
941
0
        memtables.push_back(m->NewTimestampStrippingIterator(
942
0
            ro, /*seqno_to_time_mapping=*/nullptr, &arena,
943
0
            /*prefix_extractor=*/nullptr, ts_sz));
944
2.33k
      } else {
945
2.33k
        memtables.push_back(
946
2.33k
            m->NewIterator(ro, /*seqno_to_time_mapping=*/nullptr, &arena,
947
2.33k
                           /*prefix_extractor=*/nullptr, /*for_flush=*/true));
948
2.33k
      }
949
2.33k
      auto* range_del_iter =
950
2.33k
          logical_strip_timestamp
951
2.33k
              ? m->NewTimestampStrippingRangeTombstoneIterator(
952
0
                    ro, kMaxSequenceNumber, ts_sz)
953
2.33k
              : m->NewRangeTombstoneIterator(ro, kMaxSequenceNumber,
954
2.33k
                                             true /* immutable_memtable */);
955
2.33k
      if (range_del_iter != nullptr) {
956
0
        range_del_iters.emplace_back(range_del_iter);
957
0
      }
958
2.33k
      total_num_input_entries += m->NumEntries();
959
2.33k
      total_num_deletes += m->NumDeletion();
960
2.33k
      total_data_size += m->GetDataSize();
961
2.33k
      total_memory_usage += m->ApproximateMemoryUsage();
962
2.33k
      total_num_range_deletes += m->NumRangeDeletion();
963
2.33k
    }
964
965
2.33k
    RecordInHistogram(stats_, FLUSH_MEMTABLE_MEMORY_BYTES, total_memory_usage);
966
2.33k
    RecordInHistogram(stats_, FLUSH_MEMTABLE_TOTAL_DATA_SIZE, total_data_size);
967
2.33k
    if (flush_reason_ == FlushReason::kWriteBufferFull) {
968
0
      RecordTick(stats_, FLUSH_REASON_WRITE_BUFFER_FULL);
969
0
      RecordInHistogram(stats_, FLUSH_WRITE_BUFFER_FULL_MEMTABLE_MEMORY_BYTES,
970
0
                        total_memory_usage);
971
2.33k
    } else if (flush_reason_ == FlushReason::kWriteBufferManager) {
972
0
      RecordTick(stats_, FLUSH_REASON_WRITE_BUFFER_MANAGER);
973
0
      RecordInHistogram(stats_,
974
0
                        FLUSH_WRITE_BUFFER_MANAGER_MEMTABLE_MEMORY_BYTES,
975
0
                        total_memory_usage);
976
2.33k
    } else if (flush_reason_ == FlushReason::kMemtableMaxRangeDeletions) {
977
0
      RecordTick(stats_, FLUSH_REASON_MEMTABLE_MAX_RANGE_DELETIONS);
978
0
    }
979
980
2.33k
    event_logger_->Log() << "job" << job_context_->job_id << "event"
981
2.33k
                         << "flush_started" << "num_memtables" << mems_.size()
982
2.33k
                         << "total_num_input_entries" << total_num_input_entries
983
2.33k
                         << "num_deletes" << total_num_deletes
984
2.33k
                         << "total_data_size" << total_data_size
985
2.33k
                         << "memory_usage" << total_memory_usage
986
2.33k
                         << "num_range_deletes" << total_num_range_deletes
987
2.33k
                         << "flush_reason"
988
2.33k
                         << GetFlushReasonString(flush_reason_);
989
990
2.33k
    {
991
2.33k
      ScopedArenaPtr<InternalIterator> iter(
992
2.33k
          NewMergingIterator(&cfd_->internal_comparator(), memtables.data(),
993
2.33k
                             static_cast<int>(memtables.size()), &arena));
994
2.33k
      ROCKS_LOG_INFO(db_options_.info_log,
995
2.33k
                     "[%s] [JOB %d] Level-0 flush table #%" PRIu64 ": started",
996
2.33k
                     cfd_->GetName().c_str(), job_context_->job_id,
997
2.33k
                     meta_.fd.GetNumber());
998
999
2.33k
      TEST_SYNC_POINT_CALLBACK("FlushJob::WriteLevel0Table:output_compression",
1000
2.33k
                               &output_compression_);
1001
2.33k
      int64_t _current_time = 0;
1002
2.33k
      auto status = clock_->GetCurrentTime(&_current_time);
1003
      // Safe to proceed even if GetCurrentTime fails. So, log and proceed.
1004
2.33k
      if (!status.ok()) {
1005
0
        ROCKS_LOG_WARN(
1006
0
            db_options_.info_log,
1007
0
            "Failed to get current time to populate creation_time property. "
1008
0
            "Status: %s",
1009
0
            status.ToString().c_str());
1010
0
      }
1011
2.33k
      const uint64_t current_time = static_cast<uint64_t>(_current_time);
1012
1013
2.33k
      uint64_t oldest_key_time = mems_.front()->ApproximateOldestKeyTime();
1014
1015
      // It's not clear whether oldest_key_time is always available. In case
1016
      // it is not available, use current_time.
1017
2.33k
      uint64_t oldest_ancester_time = std::min(current_time, oldest_key_time);
1018
1019
2.33k
      TEST_SYNC_POINT_CALLBACK(
1020
2.33k
          "FlushJob::WriteLevel0Table:oldest_ancester_time",
1021
2.33k
          &oldest_ancester_time);
1022
2.33k
      meta_.oldest_ancester_time = oldest_ancester_time;
1023
2.33k
      meta_.file_creation_time = current_time;
1024
1025
2.33k
      uint64_t memtable_payload_bytes = 0;
1026
2.33k
      uint64_t memtable_garbage_bytes = 0;
1027
2.33k
      IOStatus io_s;
1028
1029
2.33k
      const std::string* const full_history_ts_low =
1030
2.33k
          (full_history_ts_low_.empty()) ? nullptr : &full_history_ts_low_;
1031
2.33k
      ReadOptions read_options(Env::IOActivity::kFlush);
1032
2.33k
      read_options.rate_limiter_priority = io_priority;
1033
2.33k
      const WriteOptions write_options(io_priority, Env::IOActivity::kFlush);
1034
2.33k
      TableBuilderOptions tboptions(
1035
2.33k
          cfd_->ioptions(), mutable_cf_options_, read_options, write_options,
1036
2.33k
          cfd_->internal_comparator(), cfd_->internal_tbl_prop_coll_factories(),
1037
2.33k
          output_compression_, mutable_cf_options_.compression_opts,
1038
2.33k
          cfd_->GetID(), cfd_->GetName(), 0 /* level */,
1039
2.33k
          current_time /* newest_key_time */, false /* is_bottommost */,
1040
2.33k
          TableFileCreationReason::kFlush, oldest_key_time, current_time,
1041
2.33k
          db_id_, db_session_id_, 0 /* target_file_size */,
1042
2.33k
          meta_.fd.GetNumber(),
1043
2.33k
          preclude_last_level_min_seqno_ == kMaxSequenceNumber
1044
2.33k
              ? preclude_last_level_min_seqno_
1045
2.33k
              : std::min(earliest_snapshot_, preclude_last_level_min_seqno_),
1046
2.33k
          dbname_ /*db_name*/);
1047
2.33k
      s = BuildTable(
1048
2.33k
          dbname_, versions_, db_options_, tboptions, file_options_,
1049
2.33k
          cfd_->table_cache(), iter.get(), std::move(range_del_iters), &meta_,
1050
2.33k
          &blob_file_additions, job_context_->snapshot_seqs, earliest_snapshot_,
1051
2.33k
          job_context_->earliest_write_conflict_snapshot,
1052
2.33k
          job_context_->GetJobSnapshotSequence(),
1053
2.33k
          job_context_->snapshot_checker,
1054
2.33k
          mutable_cf_options_.paranoid_file_checks, cfd_->internal_stats(),
1055
2.33k
          &io_s, io_tracer_, BlobFileCreationReason::kFlush,
1056
2.33k
          seqno_to_time_mapping_.get(), event_logger_, job_context_->job_id,
1057
2.33k
          &table_properties_, write_hint, full_history_ts_low, blob_callback_,
1058
2.33k
          base_, &memtable_payload_bytes, &memtable_garbage_bytes, &flush_stats,
1059
2.33k
          blob_file_garbages_for_filtering, fast_sst_open_);
1060
2.33k
      TEST_SYNC_POINT_CALLBACK("FlushJob::WriteLevel0Table:s", &s);
1061
      // TODO: Cleanup io_status in BuildTable and table builders
1062
2.33k
      assert(!s.ok() || io_s.ok());
1063
2.33k
      io_s.PermitUncheckedError();
1064
2.33k
      if (s.ok() && total_num_input_entries != flush_stats.num_input_records) {
1065
0
        std::string msg = "Expected " +
1066
0
                          std::to_string(total_num_input_entries) +
1067
0
                          " entries in memtables, but read " +
1068
0
                          std::to_string(flush_stats.num_input_records);
1069
0
        ROCKS_LOG_WARN(db_options_.info_log, "[%s] [JOB %d] Level-0 flush %s",
1070
0
                       cfd_->GetName().c_str(), job_context_->job_id,
1071
0
                       msg.c_str());
1072
0
        if (db_options_.flush_verify_memtable_count) {
1073
0
          s = Status::Corruption(msg);
1074
0
        }
1075
0
      }
1076
1077
      // Only verify on table with format collects table properties
1078
2.33k
      if (s.ok() &&
1079
2.33k
          (mutable_cf_options_.table_factory->IsInstanceOf(
1080
2.33k
               TableFactory::kBlockBasedTableName()) ||
1081
0
           mutable_cf_options_.table_factory->IsInstanceOf(
1082
0
               TableFactory::kPlainTableName())) &&
1083
2.33k
          flush_stats.num_output_records != table_properties_.num_entries) {
1084
0
        std::string msg =
1085
0
            "Number of keys in flush output SST files does not match "
1086
0
            "number of keys added to the table. Expected " +
1087
0
            std::to_string(flush_stats.num_output_records) + " but there are " +
1088
0
            std::to_string(table_properties_.num_entries) +
1089
0
            " in output SST files";
1090
0
        ROCKS_LOG_WARN(db_options_.info_log, "[%s] [JOB %d] Level-0 flush %s",
1091
0
                       cfd_->GetName().c_str(), job_context_->job_id,
1092
0
                       msg.c_str());
1093
0
        if (db_options_.flush_verify_memtable_count) {
1094
0
          s = Status::Corruption(msg);
1095
0
        }
1096
0
      }
1097
2.33k
      if (tboptions.reason == TableFileCreationReason::kFlush) {
1098
2.33k
        TEST_SYNC_POINT("DBImpl::FlushJob:Flush");
1099
2.33k
        RecordTick(stats_, MEMTABLE_PAYLOAD_BYTES_AT_FLUSH,
1100
2.33k
                   memtable_payload_bytes);
1101
2.33k
        RecordTick(stats_, MEMTABLE_GARBAGE_BYTES_AT_FLUSH,
1102
2.33k
                   memtable_garbage_bytes);
1103
2.33k
      }
1104
2.33k
      LogFlush(db_options_.info_log);
1105
2.33k
    }
1106
2.33k
    ROCKS_LOG_BUFFER(log_buffer_,
1107
2.33k
                     "[%s] [JOB %d] Level-0 flush table #%" PRIu64 ": %" PRIu64
1108
2.33k
                     " bytes %s"
1109
2.33k
                     " %s"
1110
2.33k
                     " %s",
1111
2.33k
                     cfd_->GetName().c_str(), job_context_->job_id,
1112
2.33k
                     meta_.fd.GetNumber(), meta_.fd.GetFileSize(),
1113
2.33k
                     s.ToString().c_str(),
1114
2.33k
                     s.ok() && meta_.fd.GetFileSize() == 0
1115
2.33k
                         ? "It's an empty SST file from a successful flush so "
1116
2.33k
                           "won't be kept in the DB"
1117
2.33k
                         : "",
1118
2.33k
                     meta_.marked_for_compaction ? " (needs compaction)" : "");
1119
1120
2.33k
    if (s.ok() && output_file_directory_ != nullptr && sync_output_directory_) {
1121
2.33k
      s = output_file_directory_->FsyncWithDirOptions(
1122
2.33k
          IOOptions(), nullptr,
1123
2.33k
          DirFsyncOptions(DirFsyncOptions::FsyncReason::kNewFileSynced));
1124
2.33k
    }
1125
2.33k
    TEST_SYNC_POINT_CALLBACK("FlushJob::WriteLevel0Table", &mems_);
1126
2.33k
    db_mutex_->Lock();
1127
2.33k
  }
1128
2.33k
  base_->Unref();
1129
1130
  // Note that if file_size is zero, the SST has been deleted and should not be
1131
  // added to the manifest. Blob metadata updates may still need to be
1132
  // committed for direct-write files or flush-time filtering.
1133
2.33k
  const bool has_output = meta_.fd.GetFileSize() > 0;
1134
1135
2.33k
  if (s.ok()) {
1136
2.33k
    if (has_output) {
1137
2.33k
      TEST_SYNC_POINT("DBImpl::FlushJob:SSTFileCreated");
1138
      // if we have more than 1 background thread, then we cannot
1139
      // insert files directly into higher levels because some other
1140
      // threads could be concurrently producing compacted files for
1141
      // that key range.
1142
      // Add file to L0
1143
2.33k
      TEST_SYNC_POINT_CALLBACK("FileMetaData::FileMetaData", &meta_);
1144
2.33k
      edit_->AddFile(0 /* level */, meta_);
1145
2.33k
    }
1146
1147
2.33k
    edit_->SetBlobFileAdditions(std::move(blob_file_additions));
1148
2.33k
    for (auto& garbage : blob_file_garbages) {
1149
0
      edit_->AddBlobFileGarbage(std::move(garbage));
1150
0
    }
1151
2.33k
    for (auto& addition : external_blob_file_additions_) {
1152
0
      edit_->AddBlobFile(std::move(addition));
1153
0
    }
1154
2.33k
    for (auto& garbage : external_blob_file_garbages_) {
1155
0
      edit_->AddBlobFileGarbage(std::move(garbage));
1156
0
    }
1157
2.33k
    external_blob_file_additions_.clear();
1158
2.33k
    external_blob_file_garbages_.clear();
1159
2.33k
  }
1160
  // Piggyback FlushJobInfo on the first first flushed memtable.
1161
2.33k
  mems_[0]->SetFlushJobInfo(GetFlushJobInfo());
1162
1163
2.33k
  const uint64_t micros = clock_->NowMicros() - start_micros;
1164
2.33k
  const uint64_t cpu_micros = clock_->CPUMicros() - start_cpu_micros;
1165
2.33k
  flush_stats.micros = micros;
1166
2.33k
  flush_stats.cpu_micros += cpu_micros;
1167
1168
2.33k
  ROCKS_LOG_INFO(db_options_.info_log,
1169
2.33k
                 "[%s] [JOB %d] Flush lasted %" PRIu64
1170
2.33k
                 " microseconds, and %" PRIu64 " cpu microseconds.\n",
1171
2.33k
                 cfd_->GetName().c_str(), job_context_->job_id, micros,
1172
2.33k
                 flush_stats.cpu_micros);
1173
1174
2.33k
  if (has_output) {
1175
2.33k
    flush_stats.bytes_written = meta_.fd.GetFileSize();
1176
2.33k
    flush_stats.num_output_files = 1;
1177
2.33k
  }
1178
1179
2.33k
  const auto& blobs = edit_->GetBlobFileAdditions();
1180
2.33k
  for (const auto& blob : blobs) {
1181
0
    flush_stats.bytes_written_blob += blob.GetTotalBlobBytes();
1182
0
  }
1183
1184
2.33k
  flush_stats.num_output_files_blob = static_cast<int>(blobs.size());
1185
1186
2.33k
  RecordTimeToHistogram(stats_, FLUSH_TIME, flush_stats.micros);
1187
2.33k
  cfd_->internal_stats()->AddCompactionStats(0 /* level */, thread_pri_,
1188
2.33k
                                             flush_stats);
1189
2.33k
  cfd_->internal_stats()->AddCFStats(
1190
2.33k
      InternalStats::BYTES_FLUSHED,
1191
2.33k
      flush_stats.bytes_written + flush_stats.bytes_written_blob);
1192
2.33k
  RecordFlushIOStats();
1193
1194
2.33k
  return s;
1195
2.33k
}
1196
1197
2.33k
Env::IOPriority FlushJob::GetRateLimiterPriority() {
1198
2.33k
  if (versions_ && versions_->GetColumnFamilySet() &&
1199
2.33k
      versions_->GetColumnFamilySet()->write_controller()) {
1200
2.33k
    WriteController* write_controller =
1201
2.33k
        versions_->GetColumnFamilySet()->write_controller();
1202
2.33k
    if (write_controller->IsStopped() || write_controller->NeedsDelay()) {
1203
0
      return Env::IO_USER;
1204
0
    }
1205
2.33k
  }
1206
1207
2.33k
  return Env::IO_HIGH;
1208
2.33k
}
1209
1210
2.33k
std::unique_ptr<FlushJobInfo> FlushJob::GetFlushJobInfo() const {
1211
2.33k
  db_mutex_->AssertHeld();
1212
2.33k
  std::unique_ptr<FlushJobInfo> info(new FlushJobInfo{});
1213
2.33k
  info->cf_id = cfd_->GetID();
1214
2.33k
  info->cf_name = cfd_->GetName();
1215
1216
2.33k
  const uint64_t file_number = meta_.fd.GetNumber();
1217
2.33k
  info->file_path =
1218
2.33k
      MakeTableFileName(cfd_->ioptions().cf_paths[0].path, file_number);
1219
2.33k
  info->file_number = file_number;
1220
2.33k
  info->oldest_blob_file_number = meta_.oldest_blob_file_number;
1221
2.33k
  info->thread_id = db_options_.env->GetThreadID();
1222
2.33k
  info->job_id = job_context_->job_id;
1223
2.33k
  info->smallest_seqno = meta_.fd.smallest_seqno;
1224
2.33k
  info->largest_seqno = meta_.fd.largest_seqno;
1225
2.33k
  info->table_properties = table_properties_;
1226
2.33k
  info->flush_reason = flush_reason_;
1227
2.33k
  info->blob_compression_type = mutable_cf_options_.blob_compression_type;
1228
1229
  // Update BlobFilesInfo.
1230
2.33k
  for (const auto& blob_file : edit_->GetBlobFileAdditions()) {
1231
0
    BlobFileAdditionInfo blob_file_addition_info(
1232
0
        BlobFileName(cfd_->ioptions().cf_paths.front().path,
1233
0
                     blob_file.GetBlobFileNumber()) /*blob_file_path*/,
1234
0
        blob_file.GetBlobFileNumber(), blob_file.GetTotalBlobCount(),
1235
0
        blob_file.GetTotalBlobBytes());
1236
0
    info->blob_file_addition_infos.emplace_back(
1237
0
        std::move(blob_file_addition_info));
1238
0
  }
1239
2.33k
  return info;
1240
2.33k
}
1241
1242
2.33k
void FlushJob::GetEffectiveCutoffUDTForPickedMemTables() {
1243
2.33k
  db_mutex_->AssertHeld();
1244
2.33k
  assert(pick_memtable_called);
1245
2.33k
  const auto* ucmp = cfd_->internal_comparator().user_comparator();
1246
2.33k
  assert(ucmp);
1247
2.33k
  const size_t ts_sz = ucmp->timestamp_size();
1248
2.33k
  if (db_options_.atomic_flush || ts_sz == 0 ||
1249
2.33k
      cfd_->ioptions().persist_user_defined_timestamps) {
1250
2.33k
    return;
1251
2.33k
  }
1252
  // Find the newest user-defined timestamps from all the flushed memtables.
1253
0
  for (const ReadOnlyMemTable* m : mems_) {
1254
0
    Slice table_newest_udt = m->GetNewestUDT();
1255
    // Empty memtables can be legitimately created and flushed, for example
1256
    // by error recovery flush attempts.
1257
0
    if (table_newest_udt.empty()) {
1258
0
      continue;
1259
0
    }
1260
0
    if (cutoff_udt_.empty() ||
1261
0
        ucmp->CompareTimestamp(table_newest_udt, cutoff_udt_) > 0) {
1262
0
      if (!cutoff_udt_.empty()) {
1263
0
        assert(table_newest_udt.size() == cutoff_udt_.size());
1264
0
      }
1265
0
      cutoff_udt_.assign(table_newest_udt.data(), table_newest_udt.size());
1266
0
    }
1267
0
  }
1268
0
}
1269
1270
2.33k
void FlushJob::GetPrecludeLastLevelMinSeqno() {
1271
2.33k
  if (mutable_cf_options_.preclude_last_level_data_seconds == 0) {
1272
2.33k
    return;
1273
2.33k
  }
1274
  // SuperVersion should guarantee this
1275
2.33k
  assert(seqno_to_time_mapping_);
1276
0
  assert(!seqno_to_time_mapping_->Empty());
1277
0
  int64_t current_time = 0;
1278
0
  Status s = db_options_.clock->GetCurrentTime(&current_time);
1279
0
  if (!s.ok()) {
1280
0
    ROCKS_LOG_WARN(db_options_.info_log,
1281
0
                   "Failed to get current time in Flush: Status: %s",
1282
0
                   s.ToString().c_str());
1283
0
  } else {
1284
0
    SequenceNumber preserve_time_min_seqno;
1285
0
    seqno_to_time_mapping_->GetCurrentTieringCutoffSeqnos(
1286
0
        static_cast<uint64_t>(current_time),
1287
0
        mutable_cf_options_.preserve_internal_time_seconds,
1288
0
        mutable_cf_options_.preclude_last_level_data_seconds,
1289
0
        &preserve_time_min_seqno, &preclude_last_level_min_seqno_);
1290
0
  }
1291
0
}
1292
1293
2.33k
Status FlushJob::MaybeIncreaseFullHistoryTsLowToAboveCutoffUDT() {
1294
2.33k
  db_mutex_->AssertHeld();
1295
2.33k
  const auto* ucmp = cfd_->user_comparator();
1296
2.33k
  assert(ucmp);
1297
2.33k
  const std::string& full_history_ts_low = cfd_->GetFullHistoryTsLow();
1298
  // Update full_history_ts_low to right above cutoff udt only if that would
1299
  // increase it.
1300
2.33k
  if (cutoff_udt_.empty() ||
1301
0
      (!full_history_ts_low.empty() &&
1302
2.33k
       ucmp->CompareTimestamp(cutoff_udt_, full_history_ts_low) < 0)) {
1303
2.33k
    return Status::OK();
1304
2.33k
  }
1305
0
  std::string new_full_history_ts_low;
1306
0
  Slice cutoff_udt_slice = cutoff_udt_;
1307
  // TODO(yuzhangyu): Add a member to AdvancedColumnFamilyOptions for an
1308
  //  operation to get the next immediately larger user-defined timestamp to
1309
  //  expand this feature to other user-defined timestamp formats.
1310
0
  GetFullHistoryTsLowFromU64CutoffTs(&cutoff_udt_slice,
1311
0
                                     &new_full_history_ts_low);
1312
0
  VersionEdit edit;
1313
0
  edit.SetColumnFamily(cfd_->GetID());
1314
0
  edit.SetFullHistoryTsLow(new_full_history_ts_low);
1315
0
  return versions_->LogAndApply(cfd_, ReadOptions(Env::IOActivity::kFlush),
1316
0
                                WriteOptions(Env::IOActivity::kFlush), &edit,
1317
0
                                db_mutex_, output_file_directory_);
1318
2.33k
}
1319
1320
}  // namespace ROCKSDB_NAMESPACE