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]
