github-actions[bot] commented on code in PR #67657:
URL: https://github.com/apache/doris/pull/67657#discussion_r4122320368


##########
be/src/storage/index/snii/writer/bounded_run_merge.cpp:
##########
@@ -0,0 +1,497 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include <sys/resource.h>
+#include <unistd.h>
+
+#include <algorithm>
+#include <array>
+#include <atomic>
+#include <climits>
+#include <cstdio>
+#include <cstring>
+#include <limits>
+#include <memory>
+#include <queue>
+#include <utility>
+
+#include "storage/index/snii/encoding/crc32c.h"
+#include "storage/index/snii/writer/encoded_spill_run.h"
+#include "storage/index/snii/writer/spill_run_codec.h"
+#include "storage/index/snii/writer/temp_dir.h"
+
+namespace doris::snii::writer {
+namespace {
+
+Status invalid_run(const char* reason) {
+    return Status::Error<ErrorCode::INVERTED_INDEX_FILE_CORRUPTED, 
false>("spill merge: {}",
+                                                                          
reason);
+}
+
+std::atomic<uint64_t> g_run_compactions {0};
+
+struct RunRange {
+    std::string path;
+    uint64_t begin = 0;
+    uint64_t end = UINT64_MAX;
+};
+
+// Caller-owned fixture files, or a bounded directory of sealed ranges in one
+// ingestion spool. The latter has no resident path/offset array proportional 
to
+// the number of spills.
+struct RunInputs {
+    const std::vector<std::string>* files = nullptr;
+    const std::string* spool = nullptr;
+    PostingByteBuffer* ends = nullptr;
+    size_t count = 0;
+    size_t fan_in_limit = 0;
+};
+
+using Readers = std::vector<std::unique_ptr<StreamingRunReader>>;
+
+std::string merge_run_path() {
+    static std::atomic<uint64_t> sequence {0};
+    return resolve_temp_dir() + "/snii_merge_" + std::to_string(::getpid()) + 
"_" +
+           std::to_string(sequence.fetch_add(1, std::memory_order_relaxed)) + 
".run";
+}
+
+// One-token lookahead survives a fill boundary, including a document split
+// across any number of input fragments. Positions are appended incrementally
+// into the consumer's replayable buffer, never accumulated per document here.
+class RunTermSource final : public TermPostingSource {
+public:
+    RunTermSource(Readers* readers, const std::vector<size_t>* matching, bool 
positions)
+            : readers_(readers), matching_(matching), positions_(positions) {}
+
+    Status fill(uint32_t target_docs, TermPostingBuffer* out, bool* exhausted) 
override {
+        if (target_docs == 0 || out == nullptr || exhausted == nullptr || 
!out->empty()) {
+            return Status::Error<ErrorCode::INVALID_ARGUMENT, false>(
+                    "spill merge: invalid source fill arguments");
+        }
+        if (!pending_) {
+            RETURN_IF_ERROR(next());
+        }
+        while (pending_ && out->document_count() < target_docs) {
+            const uint32_t docid = docid_;
+            MutableTermPostingSpan document;
+            RETURN_IF_ERROR(out->grow_uninitialized(1, true, 0, &document));
+            document.docids[0] = docid;
+            uint64_t frequency = 0;
+            do {
+                if (positions_) {
+                    RETURN_IF_ERROR(out->append_position(position_));
+                }
+                if (++frequency > std::numeric_limits<uint32_t>::max()) {
+                    return invalid_run("document frequency overflows uint32");
+                }
+                RETURN_IF_ERROR(next());
+            } while (pending_ && docid_ == docid);
+            // Appending positions can replace the buffer's vectors. Acquire 
the
+            // frequency slot again rather than retaining an invalidated span.
+            out->set_last_frequency(static_cast<uint32_t>(frequency));
+            if (pending_ && docid_ < docid) {
+                return invalid_run("docids go backwards across run fragments");
+            }
+        }
+        *exhausted = !pending_ && run_ == matching_->size();
+        exhausted_ = *exhausted;
+        return Status::OK();
+    }
+
+    bool exhausted() const { return exhausted_; }
+
+private:
+    Status next() {
+        pending_ = false;
+        while (run_ < matching_->size()) {
+            bool end = false;
+            
RETURN_IF_ERROR((*readers_)[(*matching_)[run_]]->next_token(&docid_, 
&position_, &end));
+            if (!end) {
+                pending_ = true;
+                break;
+            }
+            ++run_;
+        }
+        return Status::OK();
+    }
+
+    Readers* readers_;
+    const std::vector<size_t>* matching_;
+    const bool positions_;
+    size_t run_ = 0;
+    uint32_t docid_ = 0;
+    uint32_t position_ = 0;
+    bool pending_ = false;
+    bool exhausted_ = false;
+};
+
+struct HeapItem {
+    uint32_t id;
+    size_t run;
+};
+struct Greater {
+    const std::vector<uint32_t>* rank;
+    bool operator()(const HeapItem& left, const HeapItem& right) const {
+        const uint32_t a = (*rank)[left.id];
+        const uint32_t b = (*rank)[right.id];
+        return a == b ? left.run > right.run : a > b;
+    }
+};
+
+using RunHeap = std::priority_queue<HeapItem, std::vector<HeapItem>, Greater>;
+
+Status take_matching_runs(RunHeap* heap, const Readers& readers, 
EncodedRunTerm* term,
+                          std::vector<size_t>* matching) {
+    matching->clear();
+    while (!heap->empty() && heap->top().id == term->term_id) {
+        const size_t run = heap->top().run;
+        heap->pop();
+        const auto& input = readers[run]->current();
+        if (input.has_positions != term->has_positions) {
+            return invalid_run("posting shape differs");
+        }
+        if (input.document_groups > UINT64_MAX - term->document_groups ||
+            input.tokens > UINT64_MAX - term->tokens) {
+            return invalid_run("term counts overflow");
+        }
+        term->document_groups += input.document_groups;
+        term->tokens += input.tokens;
+        matching->push_back(run);
+    }
+    return Status::OK();
+}
+
+template <typename Consumer>
+Status merge_group(const std::vector<RunRange>& paths, const 
std::vector<uint32_t>& rank,
+                   bool positions, MemoryReporter* reporter, bool allow_legacy,
+                   const Consumer& consume) {
+    auto metadata = reporter == nullptr ? MemoryReporter::Reservation()
+                                        : 
reporter->make_postings_reservation();
+    if (reporter != nullptr) {
+        // Reader objects own fixed, separately charged payload buffers. 
Reserve
+        // their object storage plus both heap and matching-index capacities.
+        RETURN_IF_ERROR(metadata.set_bytes(paths.size() * 
(StreamingRunReader::object_bytes() +
+                                                           sizeof(HeapItem) + 
2 * sizeof(size_t))));
+    }
+    Readers readers;
+    readers.reserve(paths.size());
+    std::vector<HeapItem> heap_storage;
+    heap_storage.reserve(paths.size());
+    RunHeap heap(Greater {&rank}, std::move(heap_storage));
+    for (size_t run = 0; run < paths.size(); ++run) {
+        auto reader = std::make_unique<StreamingRunReader>(reporter);
+        RETURN_IF_ERROR(reader->open(paths[run].path, positions, allow_legacy, 
paths[run].begin,
+                                     paths[run].end));
+        if (!reader->exhausted()) {
+            if (reader->current().term_id >= rank.size()) {
+                return invalid_run("term id out of range");
+            }
+            heap.push({reader->current().term_id, run});
+        }
+        readers.push_back(std::move(reader));
+    }
+    std::vector<size_t> matching;
+    matching.reserve(paths.size());
+    while (!heap.empty()) {
+        const uint32_t id = heap.top().id;
+        EncodedRunTerm term {.term_id = id,
+                             .has_positions = 
readers[heap.top().run]->current().has_positions};
+        RETURN_IF_ERROR(take_matching_runs(&heap, readers, &term, &matching));
+        RETURN_IF_ERROR(consume(term, &readers, matching));
+        for (size_t run : matching) {
+            RETURN_IF_ERROR(readers[run]->advance());
+            if (!readers[run]->exhausted()) {
+                const uint32_t next_id = readers[run]->current().term_id;
+                if (next_id >= rank.size()) {
+                    return invalid_run("term id out of range");
+                }
+                if (rank[next_id] <= rank[id]) {
+                    return invalid_run("terms are not strictly ordered");
+                }
+                heap.push({next_id, run});
+            }
+        }
+    }
+    return Status::OK();
+}
+
+Status compact_group(const std::vector<RunRange>& paths, const 
std::vector<uint32_t>& rank,
+                     bool positions, const std::string& output, 
MemoryReporter* reporter,
+                     bool allow_legacy) {
+    EncodedRunWriter writer(reporter);
+    RETURN_IF_ERROR(writer.open(output));
+    RETURN_IF_ERROR(merge_group(
+            paths, rank, positions, reporter, allow_legacy,
+            [&](const EncodedRunTerm& term, Readers* readers, const 
std::vector<size_t>& matching) {
+                RETURN_IF_ERROR(writer.begin_term(term));
+                for (size_t run : matching) {
+                    
RETURN_IF_ERROR((*readers)[run]->copy_fragments_to(&writer));
+                }
+                return writer.end_term();
+            }));
+    RETURN_IF_ERROR(writer.close());
+    g_run_compactions.fetch_add(1, std::memory_order_relaxed);
+    return Status::OK();
+}
+
+class ReducedRuns {
+public:
+    ~ReducedRuns() {
+        // Manifests also retain names from failed/partially completed passes.
+        // Cleanup reads directly into stack storage and cannot need a budget
+        // reservation while unwinding a memory-limit failure.
+        std::array<char, PATH_MAX> path {};
+        for (auto& manifest : manifests_) {
+            uint64_t offset = 0;
+            while (offset < manifest->size()) {
+                uint32_t length = 0;
+                if (!manifest->read_at(offset, 
{reinterpret_cast<uint8_t*>(&length), 4}).ok()) {

Review Comment:
   On a full temporary volume, a manifest flush can fail after earlier merge 
groups have created `.run` files. `PostingByteBuffer::read_at()` retries that 
failed flush before reading even completed manifest records, so this destructor 
breaks and never unlinks those outputs. Repeated failed merges can keep 
consuming the scarce disk space. Keep cleanup paths readable without a 
successful manifest flush and remove the current partial output on failure.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to