|
| 1 | +diff --git a/src/paimon/core/io/merged_key_value_record_reader.cpp b/src/paimon/core/io/merged_key_value_record_reader.cpp |
| 2 | +index 8a8a9cc..5a94236 100644 |
| 3 | +--- a/src/paimon/core/io/merged_key_value_record_reader.cpp |
| 4 | ++++ b/src/paimon/core/io/merged_key_value_record_reader.cpp |
| 5 | +@@ -17,6 +17,7 @@ |
| 6 | + #include "paimon/core/io/merged_key_value_record_reader.h" |
| 7 | + |
| 8 | + #include <cassert> |
| 9 | ++#include <memory> |
| 10 | + #include <optional> |
| 11 | + #include <utility> |
| 12 | + |
| 13 | +@@ -36,56 +37,77 @@ MergedKeyValueRecordReader::MergedKeyValueRecordReader( |
| 14 | + assert(merge_function_wrapper_ != nullptr); |
| 15 | + } |
| 16 | + |
| 17 | +-Result<std::unique_ptr<MergedKeyValueRecordReader::Iterator>> |
| 18 | +-MergedKeyValueRecordReader::Iterator::Create(MergedKeyValueRecordReader* reader) { |
| 19 | +- std::unique_ptr<Iterator> iterator(new Iterator(reader)); |
| 20 | +- PAIMON_RETURN_NOT_OK(iterator->LoadNextKeyValue()); |
| 21 | +- return iterator; |
| 22 | ++Result<bool> MergedKeyValueRecordReader::Iterator::HasNext() const { |
| 23 | ++ if (next_merged_key_value_.has_value()) { |
| 24 | ++ return true; |
| 25 | ++ } |
| 26 | ++ return PrepareNextMergedKeyValue(); |
| 27 | + } |
| 28 | + |
| 29 | + Result<KeyValue> MergedKeyValueRecordReader::Iterator::Next() { |
| 30 | +- assert(next_key_value_.has_value()); |
| 31 | +- reader_->merge_function_wrapper_->Reset(); |
| 32 | +- auto current_key = next_key_value_->key; |
| 33 | +- PAIMON_RETURN_NOT_OK(reader_->merge_function_wrapper_->Add(std::move(*next_key_value_))); |
| 34 | +- next_key_value_.reset(); |
| 35 | ++ if (!next_merged_key_value_.has_value()) { |
| 36 | ++ return Status::Invalid("No more merged key values in current iterator"); |
| 37 | ++ } |
| 38 | + |
| 39 | ++ KeyValue result = std::move(*next_merged_key_value_); |
| 40 | ++ next_merged_key_value_.reset(); |
| 41 | ++ return result; |
| 42 | ++} |
| 43 | ++ |
| 44 | ++Result<bool> MergedKeyValueRecordReader::Iterator::PrepareNextMergedKeyValue() const { |
| 45 | + while (true) { |
| 46 | +- PAIMON_RETURN_NOT_OK(LoadNextKeyValue()); |
| 47 | +- if (!next_key_value_.has_value()) { |
| 48 | +- break; |
| 49 | ++ PAIMON_ASSIGN_OR_RAISE(bool has_group, MergeNextKeyGroup()); |
| 50 | ++ if (!has_group) { |
| 51 | ++ return false; |
| 52 | + } |
| 53 | +- if (reader_->key_comparator_->CompareTo(*current_key, *next_key_value_->key) != 0) { |
| 54 | +- break; |
| 55 | ++ if (next_merged_key_value_.has_value()) { |
| 56 | ++ return true; |
| 57 | + } |
| 58 | +- PAIMON_RETURN_NOT_OK(reader_->merge_function_wrapper_->Add(std::move(*next_key_value_))); |
| 59 | +- next_key_value_.reset(); |
| 60 | + } |
| 61 | ++} |
| 62 | ++ |
| 63 | ++Result<bool> MergedKeyValueRecordReader::Iterator::MergeNextKeyGroup() const { |
| 64 | ++ PAIMON_RETURN_NOT_OK(LoadLookaheadKeyValue()); |
| 65 | ++ if (!lookahead_key_value_.has_value()) { |
| 66 | ++ return false; |
| 67 | ++ } |
| 68 | ++ |
| 69 | ++ reader_->merge_function_wrapper_->Reset(); |
| 70 | ++ auto current_key = lookahead_key_value_->key; |
| 71 | ++ |
| 72 | ++ do { |
| 73 | ++ PAIMON_RETURN_NOT_OK( |
| 74 | ++ reader_->merge_function_wrapper_->Add(std::move(*lookahead_key_value_))); |
| 75 | ++ lookahead_key_value_.reset(); |
| 76 | ++ PAIMON_RETURN_NOT_OK(LoadLookaheadKeyValue()); |
| 77 | ++ } while (lookahead_key_value_.has_value() && |
| 78 | ++ reader_->key_comparator_->CompareTo(*current_key, *lookahead_key_value_->key) == 0); |
| 79 | + |
| 80 | + PAIMON_ASSIGN_OR_RAISE(std::optional<KeyValue> result, |
| 81 | + reader_->merge_function_wrapper_->GetResult()); |
| 82 | +- // TODO(jinli.zjw): support merge function producing no result (e.g. all rows are filtered out) |
| 83 | +- if (result == std::nullopt) { |
| 84 | +- return Status::Invalid("merged key value reader produced empty result"); |
| 85 | +- } |
| 86 | +- return std::move(result).value(); |
| 87 | ++ next_merged_key_value_ = std::move(result); |
| 88 | ++ return true; |
| 89 | + } |
| 90 | + |
| 91 | +-Status MergedKeyValueRecordReader::Iterator::LoadNextKeyValue() { |
| 92 | +- if (next_key_value_.has_value()) { |
| 93 | ++Status MergedKeyValueRecordReader::Iterator::LoadLookaheadKeyValue() const { |
| 94 | ++ if (lookahead_key_value_.has_value() || exhausted_) { |
| 95 | + return Status::OK(); |
| 96 | + } |
| 97 | + |
| 98 | + while (true) { |
| 99 | +- if (current_iterator_ != nullptr && current_iterator_->HasNext()) { |
| 100 | +- PAIMON_ASSIGN_OR_RAISE(KeyValue key_value, current_iterator_->Next()); |
| 101 | +- next_key_value_.emplace(std::move(key_value)); |
| 102 | +- return Status::OK(); |
| 103 | ++ if (current_iterator_ != nullptr) { |
| 104 | ++ PAIMON_ASSIGN_OR_RAISE(bool has_next, current_iterator_->HasNext()); |
| 105 | ++ if (has_next) { |
| 106 | ++ PAIMON_ASSIGN_OR_RAISE(KeyValue key_value, current_iterator_->Next()); |
| 107 | ++ lookahead_key_value_.emplace(std::move(key_value)); |
| 108 | ++ return Status::OK(); |
| 109 | ++ } |
| 110 | + } |
| 111 | + |
| 112 | + current_iterator_.reset(); |
| 113 | + PAIMON_ASSIGN_OR_RAISE(current_iterator_, reader_->reader_->NextBatch()); |
| 114 | + if (current_iterator_ == nullptr) { |
| 115 | ++ exhausted_ = true; |
| 116 | + return Status::OK(); |
| 117 | + } |
| 118 | + } |
| 119 | +@@ -97,8 +119,9 @@ Result<std::unique_ptr<KeyValueRecordReader::Iterator>> MergedKeyValueRecordRead |
| 120 | + } |
| 121 | + visited_ = true; |
| 122 | + |
| 123 | +- PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<Iterator> iterator, Iterator::Create(this)); |
| 124 | +- if (!iterator->HasNext()) { |
| 125 | ++ auto iterator = std::make_unique<Iterator>(this); |
| 126 | ++ PAIMON_ASSIGN_OR_RAISE(bool has_next, iterator->HasNext()); |
| 127 | ++ if (!has_next) { |
| 128 | + return std::unique_ptr<KeyValueRecordReader::Iterator>(); |
| 129 | + } |
| 130 | + return iterator; |
| 131 | +diff --git a/src/paimon/core/io/merged_key_value_record_reader.h b/src/paimon/core/io/merged_key_value_record_reader.h |
| 132 | +index 4d6ef54..e458c1d 100644 |
| 133 | +--- a/src/paimon/core/io/merged_key_value_record_reader.h |
| 134 | ++++ b/src/paimon/core/io/merged_key_value_record_reader.h |
| 135 | +@@ -34,23 +34,25 @@ class MergedKeyValueRecordReader : public KeyValueRecordReader { |
| 136 | + |
| 137 | + class Iterator : public KeyValueRecordReader::Iterator { |
| 138 | + public: |
| 139 | +- static Result<std::unique_ptr<Iterator>> Create(MergedKeyValueRecordReader* reader); |
| 140 | ++ explicit Iterator(MergedKeyValueRecordReader* reader) : reader_(reader) {} |
| 141 | + |
| 142 | +- bool HasNext() const override { |
| 143 | +- return next_key_value_.has_value(); |
| 144 | +- } |
| 145 | ++ Result<bool> HasNext() const override; |
| 146 | + |
| 147 | + Result<KeyValue> Next() override; |
| 148 | + |
| 149 | + private: |
| 150 | +- explicit Iterator(MergedKeyValueRecordReader* reader) : reader_(reader) {} |
| 151 | +- |
| 152 | +- Status LoadNextKeyValue(); |
| 153 | ++ Result<bool> PrepareNextMergedKeyValue() const; |
| 154 | ++ Result<bool> MergeNextKeyGroup() const; |
| 155 | ++ Status LoadLookaheadKeyValue() const; |
| 156 | + |
| 157 | + private: |
| 158 | + MergedKeyValueRecordReader* reader_; |
| 159 | +- std::unique_ptr<KeyValueRecordReader::Iterator> current_iterator_; |
| 160 | +- std::optional<KeyValue> next_key_value_; |
| 161 | ++ mutable std::unique_ptr<KeyValueRecordReader::Iterator> current_iterator_; |
| 162 | ++ // Lookahead raw kv used to detect the boundary between two keys. |
| 163 | ++ mutable std::optional<KeyValue> lookahead_key_value_; |
| 164 | ++ // Merged kv prepared by HasNext() and consumed by Next(). |
| 165 | ++ mutable std::optional<KeyValue> next_merged_key_value_; |
| 166 | ++ mutable bool exhausted_ = false; |
| 167 | + }; |
| 168 | + |
| 169 | + Result<std::unique_ptr<KeyValueRecordReader::Iterator>> NextBatch() override; |
0 commit comments