/
smychkov
/
SStorage
Обзор
Документация
Войти
/
smychkov
/
SStorage
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/lsm/compaction.cpp
184 строки
6 KB
Андрей Смычков
docs: перевод оставшихся комментариев на русский + актуализация AGENTS
25 апр 2026, 10:08
25 апр 2026, 10:08
f2b9fd9
Код
Авторство
О чём код?
#include "compaction.hpp" #include <algorithm> #include <memory> namespace sstorage { //============================================================================ // MergingIterator::Stream //============================================================================ bool MergingIterator::Stream::init() { if (isStreaming()) { // Загружаем первый непустой блок (пропускаем пустые, если такие есть) while (nextBlockIdx < reader->blockCount()) { currentBlock.clear(); if (!reader->loadBlockAt(nextBlockIdx, currentBlock)) return false; ++nextBlockIdx; if (!currentBlock.empty()) { pos = 0; return true; } } return false; // пустой SSTable } else { pos = 0; return !inMemory.empty(); } } bool MergingIterator::Stream::advance() { ++pos; if (isStreaming()) { if (pos < currentBlock.size()) return true; // Текущий блок закончился — загружаем следующий while (nextBlockIdx < reader->blockCount()) { currentBlock.clear(); if (!reader->loadBlockAt(nextBlockIdx, currentBlock)) return false; ++nextBlockIdx; if (!currentBlock.empty()) { pos = 0; return true; } } return false; } else { return pos < inMemory.size(); } } const Record& MergingIterator::Stream::current() const { return isStreaming() ? currentBlock[pos] : inMemory[pos]; } //============================================================================ // MergingIterator — объединяет несколько отсортированных потоков записей //============================================================================ MergingIterator::MergingIterator(std::vector<std::vector<Record>> inMemStreams) { streams_.resize(inMemStreams.size()); for (size_t i = 0; i < inMemStreams.size(); ++i) { streams_[i].inMemory = std::move(inMemStreams[i]); if (streams_[i].init()) { pushFromStream(i); } } std::make_heap(heap_.begin(), heap_.end(), std::greater<HeapEntry>()); } MergingIterator::MergingIterator(std::vector<SSTablePtr> readers) { streams_.resize(readers.size()); for (size_t i = 0; i < readers.size(); ++i) { streams_[i].reader = readers[i]; if (streams_[i].init()) { pushFromStream(i); } } std::make_heap(heap_.begin(), heap_.end(), std::greater<HeapEntry>()); } void MergingIterator::pushFromStream(size_t streamIdx) { const auto& r = streams_[streamIdx].current(); heap_.push_back({r.key(), r.seqNo(), streamIdx}); std::push_heap(heap_.begin(), heap_.end(), std::greater<HeapEntry>()); } const Record& MergingIterator::current() const { return streams_[heap_.front().streamIdx].current(); } void MergingIterator::next() { if (heap_.empty()) return; std::pop_heap(heap_.begin(), heap_.end(), std::greater<HeapEntry>()); HeapEntry top = heap_.back(); heap_.pop_back(); // Продвигаем этот поток и, если остались записи, заново кладём в heap if (streams_[top.streamIdx].advance()) { pushFromStream(top.streamIdx); } } //============================================================================ // runCompaction — основная процедура слияния SSTable-ов //============================================================================ bool runCompaction( const std::vector<SSTablePtr>& inputs, int targetLevel, const CompactionOptions& opts, std::function<std::string(int, uint64_t)> filePathBuilder, std::function<uint64_t()> seqNoAllocator, std::function<bool(const std::string& key)> shouldDropTombstone, std::vector<std::string>& outNewFiles) { if (inputs.empty()) return true; // 1. Streaming k-way merge: читаем SSTable поблочно, а не целиком. // Это ��ащищает от OOM при compaction больших уровней (гигабайты данных). MergingIterator it(inputs); std::unique_ptr<SSTableWriter> writer; std::string currentPath; uint64_t bytesWritten = 0; std::string prevKey; bool havePrevKey = false; auto closeCurrent = [&]() -> bool { if (!writer) return true; if (!writer->finish()) return false; outNewFiles.push_back(currentPath); writer.reset(); bytesWritten = 0; return true; }; auto openNew = [&]() -> bool { uint64_t seq = seqNoAllocator(); currentPath = filePathBuilder(targetLevel, seq); writer = std::make_unique<SSTableWriter>( currentPath, opts.blockSize, opts.compression, 10000, opts.bloomBitsPerKey); return writer->open(); }; while (it.valid()) { const Record& r = it.current(); // Dedupe: для одного key берём только первую (свежую) версию if (havePrevKey && r.key() == prevKey) { it.next(); continue; } // Tombstone elimination на последнем уровне (или если ниже нет) if (r.isTombstone() && shouldDropTombstone(r.key())) { prevKey = r.key(); havePrevKey = true; it.next(); continue; } // Открыть writer при необходимости if (!writer) { if (!openNew()) return false; } if (!writer->add(r)) return false; bytesWritten += r.serializedSize(); // Если размер превысил таргет — закрываем и открываем новый if (bytesWritten >= opts.targetFileSize) { if (!closeCurrent()) return false; } prevKey = r.key(); havePrevKey = true; it.next(); } // Финализируем последний writer return closeCurrent(); } }