airborne12 commented on code in PR #68661:
URL: https://github.com/apache/doris/pull/68661#discussion_r4202216780


##########
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:
   Fixed in `be0a2e4b74f17daba7142e7546f62f82517351f6`. 
`FunctionMatchTest.array_phrase_preserves_leading_analyzer_gap` reproduced the 
false positive on the previous head (1 instead of 0) using the reported custom 
analyzer. The fallback now retains the leading position increment of each 
analyzed element, while keyword elements advance one position. The test and all 
41 MATCH tests pass; the BE Release build passes.



##########
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:
   Fixed in `be0a2e4b74f17daba7142e7546f62f82517351f6`. The matcher now builds 
masks only for words reachable from the current active state, intersecting the 
term mask with that small candidate set. A 1,000-term repeated query over 2,000 
one-token rows is covered by `long_phrase_with_many_short_rows`. All 41 MATCH 
tests and the BE Release build pass.



##########
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:
   Fixed in `be0a2e4b74f17daba7142e7546f62f82517351f6`. Active phrase words are 
tracked explicitly, so an isolated live prefix no longer causes scans through 
empty lower words. `long_phrase_with_unique_start_miss` covers a unique start 
followed by a long repeated tail and a missing final term. All 41 MATCH tests 
and the BE Release build pass.



##########
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:
   Fixed in `be0a2e4b74f17daba7142e7546f62f82517351f6`. The keyword path now 
reuses one `TermInfo` and passes a one-element span to the callback; it no 
longer allocates a token vector per array element. 
`keyword_array_many_elements` covers a large miss followed by a later hit. All 
41 MATCH tests and the BE Release build pass.



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