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


##########
be/src/storage/index/inverted/gram/gram_density.cpp:
##########
@@ -0,0 +1,166 @@
+// 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 "storage/index/inverted/gram/gram_density.h"
+
+#include <algorithm>
+
+#include "storage/index/inverted/gram/gram_extractor.h"
+
+namespace doris::segment_v2::gram {
+
+DensitySolver::DensitySolver(size_t min_literal_len, size_t max_gram_len, bool 
lower_case)
+        : _min_literal_len(min_literal_len),
+          _max_gram_len(max_gram_len),
+          _lower_case(lower_case),
+          _histogram(kHashValues, 0) {}
+
+namespace gram_density_detail {
+// The bytes the extractor cuts grams from: ASCII, and never NUL -- a 
candidate gram that
+// contains one is dropped, so a literal window that contains one can never be 
found.
+inline bool indexable_byte(unsigned char c) {
+    return c < 0x80 && c != 0;
+}
+} // namespace gram_density_detail
+
+void DensitySolver::observe(std::string_view value) {
+    if (_min_literal_len < _max_gram_len || _max_gram_len < 2) {
+        return;
+    }
+    // Boundaries with a whole gram's room after them inside a window: a 
boundary at k yields
+    // [k, k+max_gram_len) at worst, so k may not exceed len - max_gram_len. 
That leaves
+    // `per_window` candidate positions in a window of the promised length.
+    const size_t per_window = _min_literal_len - _max_gram_len + 1;
+    // The extractor folds ASCII case before it hashes when the scheme says 
so; the pair
+    // hashes that decide a boundary have to be computed over the same bytes 
here.
+    const auto byte_at = [&](std::string_view run, size_t k) -> uint8_t {
+        auto c = static_cast<unsigned char>(run[k]);
+        if (_lower_case && c >= 'A' && c <= 'Z') {
+            c = static_cast<unsigned char>(c - 'A' + 'a');
+        }
+        return c;
+    };
+
+    size_t i = 0;
+    const size_t n = value.size();
+    while (i < n) {
+        // Only runs of indexable bytes: a window spanning anything else could 
never produce
+        // a gram, so it carries no evidence either way.
+        if (!gram_density_detail::indexable_byte(static_cast<unsigned 
char>(value[i]))) {
+            ++i;
+            continue;
+        }
+        size_t j = i;
+        while (j < n && 
gram_density_detail::indexable_byte(static_cast<unsigned char>(value[j]))) {
+            ++j;
+        }
+        const std::string_view run = value.substr(i, j - i);
+        i = j;
+        if (run.size() < _min_literal_len) {
+            continue;
+        }
+        // Sliding minimum over the run's pair hashes, one window per start 
position. The
+        // monotonic queue keeps this linear in the run rather than quadratic 
in the window,
+        // which is what makes it affordable on the write path.
+        //
+        // A window starting at s spans [s, s + min_literal_len) and its 
candidate boundary
+        // positions are [s, s + per_window); the last real window therefore 
ends the scan at
+        // pair run.size() - max_gram_len, and exactly run.size() - 
min_literal_len + 1 windows
+        // are counted. Running the queue over the remaining pairs would count 
max_gram_len - 2
+        // more windows that extend past the run, which biased the quantile 
toward the tail.
+        const size_t last_pair = run.size() - _max_gram_len;
+        _mono.clear();
+        size_t head = 0;
+        for (size_t p = 0; p <= last_pair; ++p) {
+            const uint16_t h = boundary_hash16(byte_at(run, p), byte_at(run, p 
+ 1));
+            while (_mono.size() > head && boundary_hash16(byte_at(run, 
_mono.back()),
+                                                          byte_at(run, 
_mono.back() + 1)) >= h) {
+                _mono.pop_back();
+            }
+            _mono.push_back(static_cast<uint32_t>(p));

Review Comment:
   [P2] Bound or account for the solver's retained queue. This loop advances 
`head` past expired entries but never removes them, and `_mono.clear()` keeps 
the capacity for the whole sample. A valid 4 MiB ASCII STRING can leave about 
484,000 positions (at least 1.9 MiB) in `_mono`, while `heap_bytes()` reports 
only the 256 KiB histogram. The writer observes the entire row before its 
sample-cap check, so concurrent gram builds can exceed the registered-build 
memory/spill signal by the retained queue bytes. Compact the consumed prefix or 
bound storage to the active window, and report any remaining capacity.



##########
regression-test/suites/inverted_index_p0/gram/test_gram_compaction.groovy:
##########
@@ -0,0 +1,168 @@
+// 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.
+
+// A gram index has to survive compaction, and until now nothing checked that 
it does.
+//
+// Gram indexes are written DOCS_ONLY, and the SNII index-merge fast path 
refuses any source
+// without positions (snii/compaction/eligibility.cpp), so a compaction cannot 
stitch these
+// indexes together -- it rebuilds each one from the merged column data 
instead. That rebuild
+// recomputes everything the segment carries: the gram scheme, the df 
statistics, the high-df
+// digest, and, where it is enabled, which posting lists get dropped. Every 
one of those is
+// derived from the segment's own size, and compaction changes the segment's 
size, so the
+// rebuilt index is not a copy of its inputs and cannot be assumed to behave 
like them.
+//
+// The suite therefore checks the property that has to hold either way: the 
same queries return
+// the same rows before and after, and the index is still there afterwards 
rather than having
+// been quietly dropped. Note the table deliberately does NOT set 
disable_auto_compaction --
+// every existing benchmark and suite for this index does, which is exactly 
why this path had
+// no coverage.
+suite("test_gram_compaction", "p0") {
+    def waitAnalyzerInstalled = { String name ->
+        def deadline = System.currentTimeMillis() + 180_000
+        Exception lastNotFound = null
+        while (System.currentTimeMillis() < deadline) {
+            try {
+                sql """SELECT TOKENIZE('probe', '"analyzer"="${name}"')"""
+                return
+            } catch (Exception e) {
+                if (!e.message.contains("Policy not found")) {
+                    throw e
+                }
+                lastNotFound = e
+                sleep(1000)
+            }
+        }
+        throw new IllegalStateException("analyzer ${name} was not installed on 
BE", lastNotFound)
+    }
+
+    def tableName = "test_gram_compaction"
+
+    // The table goes first: a policy still referenced by a table left behind 
by an earlier run
+    // cannot be dropped.
+    sql "DROP TABLE IF EXISTS ${tableName}"
+    sql "DROP INVERTED INDEX ANALYZER IF EXISTS gram_compact_ana"
+    sql "DROP INVERTED INDEX TOKENIZER IF EXISTS gram_compact_tok"
+    sql """CREATE INVERTED INDEX TOKENIZER gram_compact_tok PROPERTIES (
+        "type"="ngram", "mode"="sparse", "min_gram"="3", "max_gram"="8", 
"density"="0.5")"""
+    sql """CREATE INVERTED INDEX ANALYZER gram_compact_ana
+        PROPERTIES ("tokenizer"="gram_compact_tok")"""
+    waitAnalyzerInstalled("gram_compact_ana")
+
+    sql "DROP TABLE IF EXISTS ${tableName}"
+    sql """CREATE TABLE ${tableName} (
+        `id` bigint NULL,
+        `msg` text NULL,
+        INDEX idx_msg (`msg`) USING INVERTED
+            PROPERTIES('analyzer'='gram_compact_ana', 'support_phrase'='false')
+    ) ENGINE=OLAP DUPLICATE KEY(`id`)
+    DISTRIBUTED BY HASH(`id`) BUCKETS 1
+    PROPERTIES('replication_num'='1', 
'inverted_index_storage_format'='SNII')"""
+
+    sql "SET enable_sql_cache=false"
+    sql "SET enable_condition_cache=false"
+
+    // Several batches, so there are several rowsets for the compaction to 
merge. Row content
+    // spans the selectivity range the gate cares about: a marker on every 
row, one on a few
+    // hundred, and one on a handful.
+    def batches = 5
+    def perBatch = 800
+    for (int b = 0; b < batches; b++) {
+        def values = []
+        for (int i = 0; i < perBatch; i++) {
+            def id = b * perBatch + i
+            def rare = (id % 400 == 0) ? " rare_marker_${id}" : ""
+            def mid = (id % 7 == 0) ? " midfreq_token" : ""
+            values.add("(${id}, 'shared_prefix_common_text 
filler_${id}${mid}${rare}')")
+        }
+        sql "INSERT INTO ${tableName} VALUES ${values.join(',')}"
+    }
+
+    def patterns = [
+        "like_rare"     : "SELECT COUNT(*) FROM ${tableName} WHERE msg LIKE 
'%rare_marker_%'",
+        "like_mid"      : "SELECT COUNT(*) FROM ${tableName} WHERE msg LIKE 
'%midfreq_token%'",
+        "like_common"   : "SELECT COUNT(*) FROM ${tableName} WHERE msg LIKE 
'%shared_prefix_common%'",
+        "regexp_rare"   : "SELECT COUNT(*) FROM ${tableName} WHERE msg REGEXP 
'rare_marker_[0-9]+'",
+        "regexp_mid"    : "SELECT COUNT(*) FROM ${tableName} WHERE msg REGEXP 
'midfreq_[a-z]+'",
+        "regexp_absent" : "SELECT COUNT(*) FROM ${tableName} WHERE msg REGEXP 
'no_such_token_anywhere'",
+        "regexp_alt"    : "SELECT COUNT(*) FROM ${tableName} WHERE msg REGEXP 
'rare_marker_(0|400|800)\\\\b'",
+    ]
+
+    def runAll = { String label ->
+        def out = [:]
+        patterns.each { name, stmt -> out[name] = sql(stmt)[0][0] }
+        logger.info("${label}: ${out}")
+        return out
+    }
+
+    // Tablet statistics reach information_schema asynchronously, so a freshly 
loaded table can
+    // report 0 index bytes for a while: poll rather than read once.
+    def indexBytes = {
+        def deadline = System.currentTimeMillis() + 180_000
+        long bytes = 0L
+        while (System.currentTimeMillis() < deadline) {
+            def rows = sql """SELECT INDEX_LENGTH FROM 
information_schema.tables
+                              WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 
'${tableName}'"""
+            bytes = rows.isEmpty() ? 0L : (rows[0][0] as long)
+            if (bytes > 0) {
+                break
+            }
+            sleep(2000)
+        }
+        return bytes
+    }
+
+    // Ground truth first, with the index taken out of the picture entirely.
+    sql "SET enable_inverted_index_query=false"
+    def scanned = runAll("full scan")
+    sql "SET enable_inverted_index_query=true"
+
+    def before = runAll("indexed, before compaction")
+    scanned.each { name, value ->
+        assertEquals(value, before[name], "index and scan disagree on ${name} 
before compaction")
+    }
+    def bytesBefore = indexBytes()
+    logger.info("index bytes before compaction: ${bytesBefore}")
+    assertTrue(bytesBefore > 0, "the gram index must occupy bytes before 
compaction")
+
+    trigger_and_wait_compaction(tableName, "full")
+
+    def after = runAll("indexed, after compaction")
+    before.each { name, value ->
+        assertEquals(value, after[name], "compaction changed the answer for 
${name}")
+    }
+
+    // The index must have been rebuilt, not dropped: a compaction that 
quietly produced an
+    // index-free rowset would still answer every query above correctly, by 
scanning.
+    def bytesAfter = indexBytes()
+    logger.info("index bytes after compaction: ${bytesAfter}")
+    assertTrue(bytesAfter > 0,

Review Comment:
   [P2] Verify that the compacted gram index is actually used. Every 
post-compaction answer here also passes if gram evaluation is skipped and the 
scalar predicate scans rows. `indexBytes()` stops at the first positive value, 
which may still describe the pre-compaction rowsets because that statistic 
updates asynchronously. A compaction that drops or leaves the new index 
unreadable can therefore pass this test. After confirming the replacement 
rowset is visible, run a selective query with result caches off and require 
`RowsGramIndexFiltered > 0`, as the schema-change suite does after its rewrites.



-- 
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