1 // Copyright (c) 2018-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 #pragma once 7 8 #include <algorithm> 9 #include <iterator> 10 #include <list> 11 #include <map> 12 #include <set> 13 #include <string> 14 #include <vector> 15 16 #include "db/compaction/compaction_iteration_stats.h" 17 #include "db/dbformat.h" 18 #include "db/pinned_iterators_manager.h" 19 #include "db/range_del_aggregator.h" 20 #include "db/range_tombstone_fragmenter.h" 21 #include "db/version_edit.h" 22 #include "rocksdb/comparator.h" 23 #include "rocksdb/types.h" 24 #include "table/internal_iterator.h" 25 #include "table/scoped_arena_iterator.h" 26 #include "table/table_builder.h" 27 #include "util/heap.h" 28 #include "util/kv_map.h" 29 30 namespace ROCKSDB_NAMESPACE { 31 32 class TruncatedRangeDelIterator { 33 public: 34 TruncatedRangeDelIterator( 35 std::unique_ptr<FragmentedRangeTombstoneIterator> iter, 36 const InternalKeyComparator* icmp, const InternalKey* smallest, 37 const InternalKey* largest); 38 39 bool Valid() const; 40 41 void Next(); 42 void Prev(); 43 44 void InternalNext(); 45 46 // Seeks to the tombstone with the highest viisble sequence number that covers 47 // target (a user key). If no such tombstone exists, the position will be at 48 // the earliest tombstone that ends after target. 49 void Seek(const Slice& target); 50 51 // Seeks to the tombstone with the highest viisble sequence number that covers 52 // target (a user key). If no such tombstone exists, the position will be at 53 // the latest tombstone that starts before target. 54 void SeekForPrev(const Slice& target); 55 56 void SeekToFirst(); 57 void SeekToLast(); 58 start_key()59 ParsedInternalKey start_key() const { 60 return (smallest_ == nullptr || 61 icmp_->Compare(*smallest_, iter_->parsed_start_key()) <= 0) 62 ? iter_->parsed_start_key() 63 : *smallest_; 64 } 65 end_key()66 ParsedInternalKey end_key() const { 67 return (largest_ == nullptr || 68 icmp_->Compare(iter_->parsed_end_key(), *largest_) <= 0) 69 ? iter_->parsed_end_key() 70 : *largest_; 71 } 72 seq()73 SequenceNumber seq() const { return iter_->seq(); } 74 75 std::map<SequenceNumber, std::unique_ptr<TruncatedRangeDelIterator>> 76 SplitBySnapshot(const std::vector<SequenceNumber>& snapshots); 77 upper_bound()78 SequenceNumber upper_bound() const { return iter_->upper_bound(); } 79 lower_bound()80 SequenceNumber lower_bound() const { return iter_->lower_bound(); } 81 82 private: 83 std::unique_ptr<FragmentedRangeTombstoneIterator> iter_; 84 const InternalKeyComparator* icmp_; 85 const ParsedInternalKey* smallest_ = nullptr; 86 const ParsedInternalKey* largest_ = nullptr; 87 std::list<ParsedInternalKey> pinned_bounds_; 88 89 const InternalKey* smallest_ikey_; 90 const InternalKey* largest_ikey_; 91 }; 92 93 struct SeqMaxComparator { operatorSeqMaxComparator94 bool operator()(const TruncatedRangeDelIterator* a, 95 const TruncatedRangeDelIterator* b) const { 96 return a->seq() > b->seq(); 97 } 98 }; 99 100 struct StartKeyMinComparator { StartKeyMinComparatorStartKeyMinComparator101 explicit StartKeyMinComparator(const InternalKeyComparator* c) : icmp(c) {} 102 operatorStartKeyMinComparator103 bool operator()(const TruncatedRangeDelIterator* a, 104 const TruncatedRangeDelIterator* b) const { 105 return icmp->Compare(a->start_key(), b->start_key()) > 0; 106 } 107 108 const InternalKeyComparator* icmp; 109 }; 110 111 class ForwardRangeDelIterator { 112 public: 113 explicit ForwardRangeDelIterator(const InternalKeyComparator* icmp); 114 115 bool ShouldDelete(const ParsedInternalKey& parsed); 116 void Invalidate(); 117 AddNewIter(TruncatedRangeDelIterator * iter,const ParsedInternalKey & parsed)118 void AddNewIter(TruncatedRangeDelIterator* iter, 119 const ParsedInternalKey& parsed) { 120 iter->Seek(parsed.user_key); 121 PushIter(iter, parsed); 122 assert(active_iters_.size() == active_seqnums_.size()); 123 } 124 UnusedIdx()125 size_t UnusedIdx() const { return unused_idx_; } IncUnusedIdx()126 void IncUnusedIdx() { unused_idx_++; } 127 128 private: 129 using ActiveSeqSet = 130 std::multiset<TruncatedRangeDelIterator*, SeqMaxComparator>; 131 132 struct EndKeyMinComparator { EndKeyMinComparatorEndKeyMinComparator133 explicit EndKeyMinComparator(const InternalKeyComparator* c) : icmp(c) {} 134 operatorEndKeyMinComparator135 bool operator()(const ActiveSeqSet::const_iterator& a, 136 const ActiveSeqSet::const_iterator& b) const { 137 return icmp->Compare((*a)->end_key(), (*b)->end_key()) > 0; 138 } 139 140 const InternalKeyComparator* icmp; 141 }; 142 PushIter(TruncatedRangeDelIterator * iter,const ParsedInternalKey & parsed)143 void PushIter(TruncatedRangeDelIterator* iter, 144 const ParsedInternalKey& parsed) { 145 if (!iter->Valid()) { 146 // The iterator has been fully consumed, so we don't need to add it to 147 // either of the heaps. 148 return; 149 } 150 int cmp = icmp_->Compare(parsed, iter->start_key()); 151 if (cmp < 0) { 152 PushInactiveIter(iter); 153 } else { 154 PushActiveIter(iter); 155 } 156 } 157 PushActiveIter(TruncatedRangeDelIterator * iter)158 void PushActiveIter(TruncatedRangeDelIterator* iter) { 159 auto seq_pos = active_seqnums_.insert(iter); 160 active_iters_.push(seq_pos); 161 } 162 PopActiveIter()163 TruncatedRangeDelIterator* PopActiveIter() { 164 auto active_top = active_iters_.top(); 165 auto iter = *active_top; 166 active_iters_.pop(); 167 active_seqnums_.erase(active_top); 168 return iter; 169 } 170 PushInactiveIter(TruncatedRangeDelIterator * iter)171 void PushInactiveIter(TruncatedRangeDelIterator* iter) { 172 inactive_iters_.push(iter); 173 } 174 PopInactiveIter()175 TruncatedRangeDelIterator* PopInactiveIter() { 176 auto* iter = inactive_iters_.top(); 177 inactive_iters_.pop(); 178 return iter; 179 } 180 181 const InternalKeyComparator* icmp_; 182 size_t unused_idx_; 183 ActiveSeqSet active_seqnums_; 184 BinaryHeap<ActiveSeqSet::const_iterator, EndKeyMinComparator> active_iters_; 185 BinaryHeap<TruncatedRangeDelIterator*, StartKeyMinComparator> inactive_iters_; 186 }; 187 188 class ReverseRangeDelIterator { 189 public: 190 explicit ReverseRangeDelIterator(const InternalKeyComparator* icmp); 191 192 bool ShouldDelete(const ParsedInternalKey& parsed); 193 void Invalidate(); 194 AddNewIter(TruncatedRangeDelIterator * iter,const ParsedInternalKey & parsed)195 void AddNewIter(TruncatedRangeDelIterator* iter, 196 const ParsedInternalKey& parsed) { 197 iter->SeekForPrev(parsed.user_key); 198 PushIter(iter, parsed); 199 assert(active_iters_.size() == active_seqnums_.size()); 200 } 201 UnusedIdx()202 size_t UnusedIdx() const { return unused_idx_; } IncUnusedIdx()203 void IncUnusedIdx() { unused_idx_++; } 204 205 private: 206 using ActiveSeqSet = 207 std::multiset<TruncatedRangeDelIterator*, SeqMaxComparator>; 208 209 struct EndKeyMaxComparator { EndKeyMaxComparatorEndKeyMaxComparator210 explicit EndKeyMaxComparator(const InternalKeyComparator* c) : icmp(c) {} 211 operatorEndKeyMaxComparator212 bool operator()(const TruncatedRangeDelIterator* a, 213 const TruncatedRangeDelIterator* b) const { 214 return icmp->Compare(a->end_key(), b->end_key()) < 0; 215 } 216 217 const InternalKeyComparator* icmp; 218 }; 219 struct StartKeyMaxComparator { StartKeyMaxComparatorStartKeyMaxComparator220 explicit StartKeyMaxComparator(const InternalKeyComparator* c) : icmp(c) {} 221 operatorStartKeyMaxComparator222 bool operator()(const ActiveSeqSet::const_iterator& a, 223 const ActiveSeqSet::const_iterator& b) const { 224 return icmp->Compare((*a)->start_key(), (*b)->start_key()) < 0; 225 } 226 227 const InternalKeyComparator* icmp; 228 }; 229 PushIter(TruncatedRangeDelIterator * iter,const ParsedInternalKey & parsed)230 void PushIter(TruncatedRangeDelIterator* iter, 231 const ParsedInternalKey& parsed) { 232 if (!iter->Valid()) { 233 // The iterator has been fully consumed, so we don't need to add it to 234 // either of the heaps. 235 } else if (icmp_->Compare(iter->end_key(), parsed) <= 0) { 236 PushInactiveIter(iter); 237 } else { 238 PushActiveIter(iter); 239 } 240 } 241 PushActiveIter(TruncatedRangeDelIterator * iter)242 void PushActiveIter(TruncatedRangeDelIterator* iter) { 243 auto seq_pos = active_seqnums_.insert(iter); 244 active_iters_.push(seq_pos); 245 } 246 PopActiveIter()247 TruncatedRangeDelIterator* PopActiveIter() { 248 auto active_top = active_iters_.top(); 249 auto iter = *active_top; 250 active_iters_.pop(); 251 active_seqnums_.erase(active_top); 252 return iter; 253 } 254 PushInactiveIter(TruncatedRangeDelIterator * iter)255 void PushInactiveIter(TruncatedRangeDelIterator* iter) { 256 inactive_iters_.push(iter); 257 } 258 PopInactiveIter()259 TruncatedRangeDelIterator* PopInactiveIter() { 260 auto* iter = inactive_iters_.top(); 261 inactive_iters_.pop(); 262 return iter; 263 } 264 265 const InternalKeyComparator* icmp_; 266 size_t unused_idx_; 267 ActiveSeqSet active_seqnums_; 268 BinaryHeap<ActiveSeqSet::const_iterator, StartKeyMaxComparator> active_iters_; 269 BinaryHeap<TruncatedRangeDelIterator*, EndKeyMaxComparator> inactive_iters_; 270 }; 271 272 enum class RangeDelPositioningMode { kForwardTraversal, kBackwardTraversal }; 273 class RangeDelAggregator { 274 public: RangeDelAggregator(const InternalKeyComparator * icmp)275 explicit RangeDelAggregator(const InternalKeyComparator* icmp) 276 : icmp_(icmp) {} ~RangeDelAggregator()277 virtual ~RangeDelAggregator() {} 278 279 virtual void AddTombstones( 280 std::unique_ptr<FragmentedRangeTombstoneIterator> input_iter, 281 const InternalKey* smallest = nullptr, 282 const InternalKey* largest = nullptr) = 0; 283 ShouldDelete(const Slice & key,RangeDelPositioningMode mode)284 bool ShouldDelete(const Slice& key, RangeDelPositioningMode mode) { 285 ParsedInternalKey parsed; 286 if (!ParseInternalKey(key, &parsed)) { 287 return false; 288 } 289 return ShouldDelete(parsed, mode); 290 } 291 virtual bool ShouldDelete(const ParsedInternalKey& parsed, 292 RangeDelPositioningMode mode) = 0; 293 294 virtual void InvalidateRangeDelMapPositions() = 0; 295 296 virtual bool IsEmpty() const = 0; 297 AddFile(uint64_t file_number)298 bool AddFile(uint64_t file_number) { 299 return files_seen_.insert(file_number).second; 300 } 301 302 protected: 303 class StripeRep { 304 public: StripeRep(const InternalKeyComparator * icmp,SequenceNumber upper_bound,SequenceNumber lower_bound)305 StripeRep(const InternalKeyComparator* icmp, SequenceNumber upper_bound, 306 SequenceNumber lower_bound) 307 : icmp_(icmp), 308 forward_iter_(icmp), 309 reverse_iter_(icmp), 310 upper_bound_(upper_bound), 311 lower_bound_(lower_bound) {} 312 AddTombstones(std::unique_ptr<TruncatedRangeDelIterator> input_iter)313 void AddTombstones(std::unique_ptr<TruncatedRangeDelIterator> input_iter) { 314 iters_.push_back(std::move(input_iter)); 315 } 316 IsEmpty()317 bool IsEmpty() const { return iters_.empty(); } 318 319 bool ShouldDelete(const ParsedInternalKey& parsed, 320 RangeDelPositioningMode mode); 321 Invalidate()322 void Invalidate() { 323 if (!IsEmpty()) { 324 InvalidateForwardIter(); 325 InvalidateReverseIter(); 326 } 327 } 328 329 bool IsRangeOverlapped(const Slice& start, const Slice& end); 330 331 private: InStripe(SequenceNumber seq)332 bool InStripe(SequenceNumber seq) const { 333 return lower_bound_ <= seq && seq <= upper_bound_; 334 } 335 InvalidateForwardIter()336 void InvalidateForwardIter() { forward_iter_.Invalidate(); } 337 InvalidateReverseIter()338 void InvalidateReverseIter() { reverse_iter_.Invalidate(); } 339 340 const InternalKeyComparator* icmp_; 341 std::vector<std::unique_ptr<TruncatedRangeDelIterator>> iters_; 342 ForwardRangeDelIterator forward_iter_; 343 ReverseRangeDelIterator reverse_iter_; 344 SequenceNumber upper_bound_; 345 SequenceNumber lower_bound_; 346 }; 347 348 const InternalKeyComparator* icmp_; 349 350 private: 351 std::set<uint64_t> files_seen_; 352 }; 353 354 class ReadRangeDelAggregator final : public RangeDelAggregator { 355 public: ReadRangeDelAggregator(const InternalKeyComparator * icmp,SequenceNumber upper_bound)356 ReadRangeDelAggregator(const InternalKeyComparator* icmp, 357 SequenceNumber upper_bound) 358 : RangeDelAggregator(icmp), 359 rep_(icmp, upper_bound, 0 /* lower_bound */) {} ~ReadRangeDelAggregator()360 ~ReadRangeDelAggregator() override {} 361 362 using RangeDelAggregator::ShouldDelete; 363 void AddTombstones( 364 std::unique_ptr<FragmentedRangeTombstoneIterator> input_iter, 365 const InternalKey* smallest = nullptr, 366 const InternalKey* largest = nullptr) override; 367 ShouldDelete(const ParsedInternalKey & parsed,RangeDelPositioningMode mode)368 bool ShouldDelete(const ParsedInternalKey& parsed, 369 RangeDelPositioningMode mode) final override { 370 if (rep_.IsEmpty()) { 371 return false; 372 } 373 return ShouldDeleteImpl(parsed, mode); 374 } 375 376 bool IsRangeOverlapped(const Slice& start, const Slice& end); 377 InvalidateRangeDelMapPositions()378 void InvalidateRangeDelMapPositions() override { rep_.Invalidate(); } 379 IsEmpty()380 bool IsEmpty() const override { return rep_.IsEmpty(); } 381 382 private: 383 StripeRep rep_; 384 385 bool ShouldDeleteImpl(const ParsedInternalKey& parsed, 386 RangeDelPositioningMode mode); 387 }; 388 389 class CompactionRangeDelAggregator : public RangeDelAggregator { 390 public: CompactionRangeDelAggregator(const InternalKeyComparator * icmp,const std::vector<SequenceNumber> & snapshots)391 CompactionRangeDelAggregator(const InternalKeyComparator* icmp, 392 const std::vector<SequenceNumber>& snapshots) 393 : RangeDelAggregator(icmp), snapshots_(&snapshots) {} ~CompactionRangeDelAggregator()394 ~CompactionRangeDelAggregator() override {} 395 396 void AddTombstones( 397 std::unique_ptr<FragmentedRangeTombstoneIterator> input_iter, 398 const InternalKey* smallest = nullptr, 399 const InternalKey* largest = nullptr) override; 400 401 using RangeDelAggregator::ShouldDelete; 402 bool ShouldDelete(const ParsedInternalKey& parsed, 403 RangeDelPositioningMode mode) override; 404 405 bool IsRangeOverlapped(const Slice& start, const Slice& end); 406 InvalidateRangeDelMapPositions()407 void InvalidateRangeDelMapPositions() override { 408 for (auto& rep : reps_) { 409 rep.second.Invalidate(); 410 } 411 } 412 IsEmpty()413 bool IsEmpty() const override { 414 for (const auto& rep : reps_) { 415 if (!rep.second.IsEmpty()) { 416 return false; 417 } 418 } 419 return true; 420 } 421 422 // Creates an iterator over all the range tombstones in the aggregator, for 423 // use in compaction. Nullptr arguments indicate that the iterator range is 424 // unbounded. 425 // NOTE: the boundaries are used for optimization purposes to reduce the 426 // number of tombstones that are passed to the fragmenter; they do not 427 // guarantee that the resulting iterator only contains range tombstones that 428 // cover keys in the provided range. If required, these bounds must be 429 // enforced during iteration. 430 std::unique_ptr<FragmentedRangeTombstoneIterator> NewIterator( 431 const Slice* lower_bound = nullptr, const Slice* upper_bound = nullptr, 432 bool upper_bound_inclusive = false); 433 434 private: 435 std::vector<std::unique_ptr<TruncatedRangeDelIterator>> parent_iters_; 436 std::map<SequenceNumber, StripeRep> reps_; 437 438 const std::vector<SequenceNumber>* snapshots_; 439 }; 440 441 } // namespace ROCKSDB_NAMESPACE 442