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


##########
be/src/exprs/function/match.cpp:
##########
@@ -44,6 +49,183 @@ const InvertedIndexAnalyzerCtx* 
get_match_analyzer_ctx(FunctionContext* context)
     return analyzer_ctx;
 }
 
+enum class PhraseMode { EXACT, PREFIX, EDGE };
+
+class StreamingPhraseMatcher {
+public:
+    StreamingPhraseMatcher(const std::vector<segment_v2::TermInfo>& 
query_tokens, PhraseMode mode)
+            : _mode(mode),
+              _query_size(query_tokens.size()),
+              _last_word((_query_size - 1) / 64),
+              _last_bit(uint64_t {1} << ((_query_size - 1) % 64)),
+              _state(_last_word + 1, 0),
+              _pending_masks(_last_word + 1, 0),
+              _first_term(query_tokens.front().get_single_term()),
+              _last_term(query_tokens.back().get_single_term()) {
+        for (size_t pos = 0; pos < _query_size; ++pos) {
+            if ((_mode == PhraseMode::PREFIX && pos == _query_size - 1) ||
+                (_mode == PhraseMode::EDGE && (pos == 0 || pos == _query_size 
- 1))) {
+                continue;
+            }
+            auto& words = _exact_masks[query_tokens[pos].get_single_term()];
+            const size_t word = pos / 64;
+            const uint64_t bit = uint64_t {1} << (pos % 64);
+            if (!words.empty() && words.back().first == word) {
+                words.back().second |= bit;
+            } else {
+                words.emplace_back(word, bit);
+            }
+        }
+    }
+
+    void reset_row() {
+        for (size_t word : _pending_touched) {
+            _pending_masks[word] = 0;
+        }
+        _pending_touched.clear();
+        _pending_position = -1;
+        _previous_position = -1;
+        _active_word_count = 0;
+    }
+
+    bool feed(const std::string& term, int64_t position) {
+        if (_pending_position >= 0 && position != _pending_position) {
+            if (advance()) {
+                return true;
+            }
+        }
+        _pending_position = position;
+        const auto mask_it = _exact_masks.find(term);
+        if (mask_it != _exact_masks.end()) {
+            for (const auto& [word, mask] : mask_it->second) {
+                add_pending_mask(word, mask);

Review Comment:
   [P2] Limit pending masks to reachable words. For a 10,000-term 
`MATCH_PHRASE` query of repeated `alpha` over many one-token `alpha` rows, 
`_exact_masks["alpha"]` has 157 words. Each row starts with no active word, yet 
`feed()` adds all 157 masks and `advance()` clears all 157 while reading only 
word zero; one million rows do about 314 million avoidable mask operations. The 
prior fallback fails immediately after the missing second data token. This 
differs from the earlier absent-term miss because `alpha` is present in every 
mask word. Apply the active-word limit before inserting masks and cover many 
short rows with a long repeated query.



##########
be/src/exprs/function/match.cpp:
##########
@@ -44,6 +49,183 @@ const InvertedIndexAnalyzerCtx* 
get_match_analyzer_ctx(FunctionContext* context)
     return analyzer_ctx;
 }
 
+enum class PhraseMode { EXACT, PREFIX, EDGE };
+
+class StreamingPhraseMatcher {
+public:
+    StreamingPhraseMatcher(const std::vector<segment_v2::TermInfo>& 
query_tokens, PhraseMode mode)
+            : _mode(mode),
+              _query_size(query_tokens.size()),
+              _last_word((_query_size - 1) / 64),
+              _last_bit(uint64_t {1} << ((_query_size - 1) % 64)),
+              _state(_last_word + 1, 0),
+              _pending_masks(_last_word + 1, 0),
+              _first_term(query_tokens.front().get_single_term()),
+              _last_term(query_tokens.back().get_single_term()) {
+        for (size_t pos = 0; pos < _query_size; ++pos) {
+            if ((_mode == PhraseMode::PREFIX && pos == _query_size - 1) ||
+                (_mode == PhraseMode::EDGE && (pos == 0 || pos == _query_size 
- 1))) {
+                continue;
+            }
+            auto& words = _exact_masks[query_tokens[pos].get_single_term()];
+            const size_t word = pos / 64;
+            const uint64_t bit = uint64_t {1} << (pos % 64);
+            if (!words.empty() && words.back().first == word) {
+                words.back().second |= bit;
+            } else {
+                words.emplace_back(word, bit);
+            }
+        }
+    }
+
+    void reset_row() {
+        for (size_t word : _pending_touched) {
+            _pending_masks[word] = 0;
+        }
+        _pending_touched.clear();
+        _pending_position = -1;
+        _previous_position = -1;
+        _active_word_count = 0;
+    }
+
+    bool feed(const std::string& term, int64_t position) {
+        if (_pending_position >= 0 && position != _pending_position) {
+            if (advance()) {
+                return true;
+            }
+        }
+        _pending_position = position;
+        const auto mask_it = _exact_masks.find(term);
+        if (mask_it != _exact_masks.end()) {
+            for (const auto& [word, mask] : mask_it->second) {
+                add_pending_mask(word, mask);
+            }
+        }
+        const bool first_matches = _mode == PhraseMode::EDGE &&
+                                   (_query_size == 1 ? term.find(_first_term) 
!= std::string::npos
+                                                     : 
term.ends_with(_first_term));
+        const bool last_matches =
+                (_mode == PhraseMode::PREFIX || (_mode == PhraseMode::EDGE && 
_query_size > 1)) &&
+                term.starts_with(_last_term);
+
+        if (first_matches) {
+            add_pending_mask(0, 1);
+        }
+        if (last_matches) {
+            add_pending_mask(_last_word, _last_bit);
+        }
+        return false;
+    }
+
+    bool finish() { return _pending_position >= 0 && advance(); }
+
+private:
+    void add_pending_mask(size_t word, uint64_t mask) {
+        if (_pending_masks[word] == 0) {
+            _pending_touched.push_back(word);
+        }
+        _pending_masks[word] |= mask;
+    }
+
+    bool advance() {
+        const int64_t position = _pending_position;
+        _pending_position = -1;
+        if (_previous_position >= 0 && position != _previous_position + 1) {
+            _active_word_count = 0;
+        }
+        _previous_position = position;
+
+        size_t next_active_word_count = 0;
+        uint64_t carry = 1;
+        const size_t active_word_count = _active_word_count;
+        if (!_pending_touched.empty()) {
+            const size_t limit = std::min(_state.size(), active_word_count + 
1);
+            for (size_t word = 0; word < limit; ++word) {
+                const uint64_t previous = word < active_word_count ? 
_state[word] : 0;
+                const uint64_t next = ((previous << 1) | carry) & 
_pending_masks[word];
+                _state[word] = next;
+                if (next != 0) {
+                    next_active_word_count = word + 1;
+                }
+                carry = previous >> 63;
+            }
+        }
+        for (size_t word : _pending_touched) {
+            _pending_masks[word] = 0;
+        }
+        _pending_touched.clear();
+        _active_word_count = next_active_word_count;
+        return _active_word_count > _last_word && (_state[_last_word] & 
_last_bit) != 0;
+    }
+
+    PhraseMode _mode;
+    size_t _query_size;
+    size_t _last_word;
+    uint64_t _last_bit;
+    std::vector<uint64_t> _state;
+    std::vector<uint64_t> _pending_masks;
+    std::vector<size_t> _pending_touched;
+    std::string _first_term;
+    std::string _last_term;
+    std::unordered_map<std::string, std::vector<std::pair<size_t, uint64_t>>> 
_exact_masks;
+    int64_t _pending_position = -1;
+    int64_t _previous_position = -1;
+    size_t _active_word_count = 0;
+};
+
+template <typename Callback>
+bool for_each_data_element_tokens(const FunctionMatchBase& function, const 
std::string& column_name,
+                                  const InvertedIndexAnalyzerCtx* analyzer_ctx,
+                                  const ColumnString* string_col, size_t row,
+                                  const ColumnArray::Offsets64* array_offsets,
+                                  const ColumnUInt8::Container* 
array_element_null_map,
+                                  Callback&& callback) {
+    const size_t begin = array_offsets ? (row == 0 ? 0 : (*array_offsets)[row 
- 1]) : row;
+    const size_t end = array_offsets ? (*array_offsets)[row] : row + 1;
+    int32_t unused_array_offset = 0;
+    for (size_t element = begin; element < end; ++element) {
+        if (array_element_null_map && (*array_element_null_map)[element]) {
+            continue;
+        }
+        auto tokens = function.analyse_data_token(column_name, analyzer_ctx, 
string_col, element,
+                                                  nullptr, 
unused_array_offset);
+        if (tokens.empty()) {
+            continue;
+        }
+        if (callback(tokens)) {
+            return true;
+        }
+    }
+    return false;
+}
+
+bool match_phrase_data_tokens(const FunctionMatchBase& function, const 
std::string& column_name,
+                              const InvertedIndexAnalyzerCtx* analyzer_ctx,
+                              const ColumnString* string_col, size_t row,
+                              const ColumnArray::Offsets64* array_offsets,
+                              const ColumnUInt8::Container* 
array_element_null_map,
+                              StreamingPhraseMatcher& matcher) {
+    matcher.reset_row();
+    int64_t position_base = 0;
+    const bool matched = for_each_data_element_tokens(
+            function, column_name, analyzer_ctx, string_col, row, 
array_offsets,
+            array_element_null_map, [&](const 
std::vector<segment_v2::TermInfo>& tokens) {
+                const int32_t first_position = tokens.front().position;
+                int32_t last_position = first_position;
+                for (const auto& token : tokens) {
+                    const int64_t position = position_base + 
static_cast<int64_t>(token.position) -

Review Comment:
   [P1] Preserve leading analyzer position gaps across array elements. With a 
custom `char_group` whitespace tokenizer and default `word_delimiter` filter, 
`["hello", "--- world"]` emits `hello@1` and `world@2`: the filter discards 
`---` but carries its position increment. Subtracting `first_position` here 
moves `world` directly after `hello`, so fallback `MATCH_PHRASE 'hello world'` 
returns true; the indexed writer preserves the gap and its exact phrase matcher 
rejects the row. This differs from the earlier same-position synonym thread: 
the lost increment is on the first token of a later field. Preserve that 
increment and cover indexed/fallback parity.



##########
be/src/exprs/function/match.cpp:
##########
@@ -44,6 +49,183 @@ const InvertedIndexAnalyzerCtx* 
get_match_analyzer_ctx(FunctionContext* context)
     return analyzer_ctx;
 }
 
+enum class PhraseMode { EXACT, PREFIX, EDGE };
+
+class StreamingPhraseMatcher {
+public:
+    StreamingPhraseMatcher(const std::vector<segment_v2::TermInfo>& 
query_tokens, PhraseMode mode)
+            : _mode(mode),
+              _query_size(query_tokens.size()),
+              _last_word((_query_size - 1) / 64),
+              _last_bit(uint64_t {1} << ((_query_size - 1) % 64)),
+              _state(_last_word + 1, 0),
+              _pending_masks(_last_word + 1, 0),
+              _first_term(query_tokens.front().get_single_term()),
+              _last_term(query_tokens.back().get_single_term()) {
+        for (size_t pos = 0; pos < _query_size; ++pos) {
+            if ((_mode == PhraseMode::PREFIX && pos == _query_size - 1) ||
+                (_mode == PhraseMode::EDGE && (pos == 0 || pos == _query_size 
- 1))) {
+                continue;
+            }
+            auto& words = _exact_masks[query_tokens[pos].get_single_term()];
+            const size_t word = pos / 64;
+            const uint64_t bit = uint64_t {1} << (pos % 64);
+            if (!words.empty() && words.back().first == word) {
+                words.back().second |= bit;
+            } else {
+                words.emplace_back(word, bit);
+            }
+        }
+    }
+
+    void reset_row() {
+        for (size_t word : _pending_touched) {
+            _pending_masks[word] = 0;
+        }
+        _pending_touched.clear();
+        _pending_position = -1;
+        _previous_position = -1;
+        _active_word_count = 0;
+    }
+
+    bool feed(const std::string& term, int64_t position) {
+        if (_pending_position >= 0 && position != _pending_position) {
+            if (advance()) {
+                return true;
+            }
+        }
+        _pending_position = position;
+        const auto mask_it = _exact_masks.find(term);
+        if (mask_it != _exact_masks.end()) {
+            for (const auto& [word, mask] : mask_it->second) {
+                add_pending_mask(word, mask);
+            }
+        }
+        const bool first_matches = _mode == PhraseMode::EDGE &&
+                                   (_query_size == 1 ? term.find(_first_term) 
!= std::string::npos
+                                                     : 
term.ends_with(_first_term));
+        const bool last_matches =
+                (_mode == PhraseMode::PREFIX || (_mode == PhraseMode::EDGE && 
_query_size > 1)) &&
+                term.starts_with(_last_term);
+
+        if (first_matches) {
+            add_pending_mask(0, 1);
+        }
+        if (last_matches) {
+            add_pending_mask(_last_word, _last_bit);
+        }
+        return false;
+    }
+
+    bool finish() { return _pending_position >= 0 && advance(); }
+
+private:
+    void add_pending_mask(size_t word, uint64_t mask) {
+        if (_pending_masks[word] == 0) {
+            _pending_touched.push_back(word);
+        }
+        _pending_masks[word] |= mask;
+    }
+
+    bool advance() {
+        const int64_t position = _pending_position;
+        _pending_position = -1;
+        if (_previous_position >= 0 && position != _previous_position + 1) {
+            _active_word_count = 0;
+        }
+        _previous_position = position;
+
+        size_t next_active_word_count = 0;
+        uint64_t carry = 1;
+        const size_t active_word_count = _active_word_count;
+        if (!_pending_touched.empty()) {
+            const size_t limit = std::min(_state.size(), active_word_count + 
1);

Review Comment:
   [P2] Skip empty words below a lone live phrase prefix. For `MATCH_PHRASE` 
query `beta` followed by 100,000 `alpha` terms and then `gamma`, against a row 
containing `beta` followed by 100,000 `alpha` terms, only one prefix bit stays 
live. `_active_word_count` records its highest word plus one, so this loop 
revisits every lower empty word on every token (about 78 million word visits); 
the prior fallback checked the unique `beta` start and read the row linearly. 
Limiting mask insertion does not remove this scan. Track active words sparsely 
or skip zero runs and cover a long unique-start miss.



##########
be/src/exprs/function/match.cpp:
##########
@@ -44,6 +49,183 @@ const InvertedIndexAnalyzerCtx* 
get_match_analyzer_ctx(FunctionContext* context)
     return analyzer_ctx;
 }
 
+enum class PhraseMode { EXACT, PREFIX, EDGE };
+
+class StreamingPhraseMatcher {
+public:
+    StreamingPhraseMatcher(const std::vector<segment_v2::TermInfo>& 
query_tokens, PhraseMode mode)
+            : _mode(mode),
+              _query_size(query_tokens.size()),
+              _last_word((_query_size - 1) / 64),
+              _last_bit(uint64_t {1} << ((_query_size - 1) % 64)),
+              _state(_last_word + 1, 0),
+              _pending_masks(_last_word + 1, 0),
+              _first_term(query_tokens.front().get_single_term()),
+              _last_term(query_tokens.back().get_single_term()) {
+        for (size_t pos = 0; pos < _query_size; ++pos) {
+            if ((_mode == PhraseMode::PREFIX && pos == _query_size - 1) ||
+                (_mode == PhraseMode::EDGE && (pos == 0 || pos == _query_size 
- 1))) {
+                continue;
+            }
+            auto& words = _exact_masks[query_tokens[pos].get_single_term()];
+            const size_t word = pos / 64;
+            const uint64_t bit = uint64_t {1} << (pos % 64);
+            if (!words.empty() && words.back().first == word) {
+                words.back().second |= bit;
+            } else {
+                words.emplace_back(word, bit);
+            }
+        }
+    }
+
+    void reset_row() {
+        for (size_t word : _pending_touched) {
+            _pending_masks[word] = 0;
+        }
+        _pending_touched.clear();
+        _pending_position = -1;
+        _previous_position = -1;
+        _active_word_count = 0;
+    }
+
+    bool feed(const std::string& term, int64_t position) {
+        if (_pending_position >= 0 && position != _pending_position) {
+            if (advance()) {
+                return true;
+            }
+        }
+        _pending_position = position;
+        const auto mask_it = _exact_masks.find(term);
+        if (mask_it != _exact_masks.end()) {
+            for (const auto& [word, mask] : mask_it->second) {
+                add_pending_mask(word, mask);
+            }
+        }
+        const bool first_matches = _mode == PhraseMode::EDGE &&
+                                   (_query_size == 1 ? term.find(_first_term) 
!= std::string::npos
+                                                     : 
term.ends_with(_first_term));
+        const bool last_matches =
+                (_mode == PhraseMode::PREFIX || (_mode == PhraseMode::EDGE && 
_query_size > 1)) &&
+                term.starts_with(_last_term);
+
+        if (first_matches) {
+            add_pending_mask(0, 1);
+        }
+        if (last_matches) {
+            add_pending_mask(_last_word, _last_bit);
+        }
+        return false;
+    }
+
+    bool finish() { return _pending_position >= 0 && advance(); }
+
+private:
+    void add_pending_mask(size_t word, uint64_t mask) {
+        if (_pending_masks[word] == 0) {
+            _pending_touched.push_back(word);
+        }
+        _pending_masks[word] |= mask;
+    }
+
+    bool advance() {
+        const int64_t position = _pending_position;
+        _pending_position = -1;
+        if (_previous_position >= 0 && position != _previous_position + 1) {
+            _active_word_count = 0;
+        }
+        _previous_position = position;
+
+        size_t next_active_word_count = 0;
+        uint64_t carry = 1;
+        const size_t active_word_count = _active_word_count;
+        if (!_pending_touched.empty()) {
+            const size_t limit = std::min(_state.size(), active_word_count + 
1);
+            for (size_t word = 0; word < limit; ++word) {
+                const uint64_t previous = word < active_word_count ? 
_state[word] : 0;
+                const uint64_t next = ((previous << 1) | carry) & 
_pending_masks[word];
+                _state[word] = next;
+                if (next != 0) {
+                    next_active_word_count = word + 1;
+                }
+                carry = previous >> 63;
+            }
+        }
+        for (size_t word : _pending_touched) {
+            _pending_masks[word] = 0;
+        }
+        _pending_touched.clear();
+        _active_word_count = next_active_word_count;
+        return _active_word_count > _last_word && (_state[_last_word] & 
_last_bit) != 0;
+    }
+
+    PhraseMode _mode;
+    size_t _query_size;
+    size_t _last_word;
+    uint64_t _last_bit;
+    std::vector<uint64_t> _state;
+    std::vector<uint64_t> _pending_masks;
+    std::vector<size_t> _pending_touched;
+    std::string _first_term;
+    std::string _last_term;
+    std::unordered_map<std::string, std::vector<std::pair<size_t, uint64_t>>> 
_exact_masks;
+    int64_t _pending_position = -1;
+    int64_t _previous_position = -1;
+    size_t _active_word_count = 0;
+};
+
+template <typename Callback>
+bool for_each_data_element_tokens(const FunctionMatchBase& function, const 
std::string& column_name,
+                                  const InvertedIndexAnalyzerCtx* analyzer_ctx,
+                                  const ColumnString* string_col, size_t row,
+                                  const ColumnArray::Offsets64* array_offsets,
+                                  const ColumnUInt8::Container* 
array_element_null_map,
+                                  Callback&& callback) {
+    const size_t begin = array_offsets ? (row == 0 ? 0 : (*array_offsets)[row 
- 1]) : row;
+    const size_t end = array_offsets ? (*array_offsets)[row] : row + 1;
+    int32_t unused_array_offset = 0;
+    for (size_t element = begin; element < end; ++element) {
+        if (array_element_null_map && (*array_element_null_map)[element]) {
+            continue;
+        }
+        auto tokens = function.analyse_data_token(column_name, analyzer_ctx, 
string_col, element,

Review Comment:
   [P2] Reuse token storage for keyword array elements. On a `MATCH_ANY` miss 
over a million-element keyword `ARRAY<STRING>`, this call constructs a fresh 
`std::vector<TermInfo>` with one entry per element and destroys it before the 
next element, causing a million allocations and frees even though no analysis 
is needed. The old keyword array path grew one vector; the earlier retention 
thread concerned its memory cost, not this new allocator churn. Reuse a 
one-token scratch buffer or pass the keyword term directly to the callback 
while keeping bounded row memory.



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