Skip to content

Commit 3d3dee6

Browse files
committed
feat(core): support null merge results in key-value readers
1 parent 2aeaca6 commit 3d3dee6

19 files changed

Lines changed: 235 additions & 75 deletions

src/paimon/core/io/key_value_data_file_record_reader.cpp

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,7 @@ KeyValueDataFileRecordReader::KeyValueDataFileRecordReader(
4848
value_schema_(value_schema),
4949
value_names_(value_schema_->field_names()) {}
5050

51-
bool KeyValueDataFileRecordReader::Iterator::HasNext() const {
51+
Result<bool> KeyValueDataFileRecordReader::Iterator::HasNext() const {
5252
int64_t array_length = reader_->row_kind_array_->length();
5353
const auto& selection_bitmap = reader_->selection_bitmap_;
5454
if (selection_bitmap.Cardinality() == array_length) {
@@ -67,7 +67,6 @@ bool KeyValueDataFileRecordReader::Iterator::HasNext() const {
6767
}
6868

6969
Result<KeyValue> KeyValueDataFileRecordReader::Iterator::Next() {
70-
assert(HasNext());
7170
// key is only used in merge sort; key context does not hold parent struct array
7271
auto key = std::make_unique<ColumnarRowRef>(reader_->key_ctx_, cursor_);
7372
// value is used in merge sort and projection (maybe async and multi-thread), so value context

src/paimon/core/io/key_value_data_file_record_reader.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -56,7 +56,7 @@ class KeyValueDataFileRecordReader : public KeyValueRecordReader {
5656
public:
5757
Iterator(KeyValueDataFileRecordReader* reader, int64_t previous_batch_first_row_number)
5858
: previous_batch_first_row_number_(previous_batch_first_row_number), reader_(reader) {}
59-
bool HasNext() const override;
59+
Result<bool> HasNext() const override;
6060
Result<KeyValue> Next() override;
6161
Result<std::pair<int64_t, KeyValue>> NextWithFilePos();
6262

src/paimon/core/io/key_value_data_file_record_reader_test.cpp

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -285,7 +285,11 @@ TEST_F(KeyValueDataFileRecordReaderTest, TestWithSelectedBitmapWithFilePos) {
285285
auto typed_iter = dynamic_cast<KeyValueDataFileRecordReader::Iterator*>(iter);
286286
ASSERT_TRUE(typed_iter);
287287
size_t pos_iter = 0;
288-
while (iter->HasNext()) {
288+
while (true) {
289+
ASSERT_OK_AND_ASSIGN(bool has_next, iter->HasNext());
290+
if (!has_next) {
291+
break;
292+
}
289293
ASSERT_OK_AND_ASSIGN(auto kv_and_pos, typed_iter->NextWithFilePos());
290294
const auto& [pos, kv] = kv_and_pos;
291295
ASSERT_EQ(pos, expected_pos_vector[pos_iter++]);

src/paimon/core/io/key_value_in_memory_record_reader.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -55,7 +55,7 @@ class KeyValueInMemoryRecordReader : public KeyValueRecordReader {
5555
class Iterator : public KeyValueRecordReader::Iterator {
5656
public:
5757
explicit Iterator(KeyValueInMemoryRecordReader* reader) : reader_(reader) {}
58-
bool HasNext() const override {
58+
Result<bool> HasNext() const override {
5959
return cursor_ < reader_->value_struct_array_->length();
6060
}
6161
Result<KeyValue> Next() override;

src/paimon/core/io/key_value_record_reader.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020

2121
#include "paimon/core/key_value.h"
2222
#include "paimon/metrics.h"
23+
#include "paimon/result.h"
2324
namespace paimon {
2425
class KeyValueRecordReader {
2526
public:
@@ -28,7 +29,7 @@ class KeyValueRecordReader {
2829
class Iterator {
2930
public:
3031
virtual ~Iterator() = default;
31-
virtual bool HasNext() const = 0;
32+
virtual Result<bool> HasNext() const = 0;
3233
virtual Result<KeyValue> Next() = 0;
3334
};
3435

src/paimon/core/io/merged_key_value_record_reader.cpp

Lines changed: 55 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
#include "paimon/core/io/merged_key_value_record_reader.h"
1818

1919
#include <cassert>
20+
#include <memory>
2021
#include <optional>
2122
#include <utility>
2223

@@ -36,56 +37,79 @@ MergedKeyValueRecordReader::MergedKeyValueRecordReader(
3637
assert(merge_function_wrapper_ != nullptr);
3738
}
3839

39-
Result<std::unique_ptr<MergedKeyValueRecordReader::Iterator>>
40-
MergedKeyValueRecordReader::Iterator::Create(MergedKeyValueRecordReader* reader) {
41-
std::unique_ptr<Iterator> iterator(new Iterator(reader));
42-
PAIMON_RETURN_NOT_OK(iterator->LoadNextKeyValue());
43-
return iterator;
40+
Result<bool> MergedKeyValueRecordReader::Iterator::HasNext() const {
41+
if (merged_key_value_.has_value()) {
42+
return true;
43+
}
44+
return PrepareNextMergedKeyValue();
4445
}
4546

4647
Result<KeyValue> MergedKeyValueRecordReader::Iterator::Next() {
47-
assert(next_key_value_.has_value());
48-
reader_->merge_function_wrapper_->Reset();
49-
auto current_key = next_key_value_->key;
50-
PAIMON_RETURN_NOT_OK(reader_->merge_function_wrapper_->Add(std::move(*next_key_value_)));
51-
next_key_value_.reset();
48+
if (!merged_key_value_.has_value()) {
49+
return Status::Invalid("No more merged key values in current iterator");
50+
}
5251

52+
KeyValue result = std::move(*merged_key_value_);
53+
merged_key_value_.reset();
54+
return result;
55+
}
56+
57+
Result<bool> MergedKeyValueRecordReader::Iterator::PrepareNextMergedKeyValue() const {
5358
while (true) {
54-
PAIMON_RETURN_NOT_OK(LoadNextKeyValue());
55-
if (!next_key_value_.has_value()) {
56-
break;
59+
PAIMON_ASSIGN_OR_RAISE(bool has_next, MergeNextKey());
60+
if (!has_next) {
61+
return false;
5762
}
58-
if (reader_->key_comparator_->CompareTo(*current_key, *next_key_value_->key) != 0) {
59-
break;
63+
// if merged_key_value_ is empty(maybe all filtered out), continue to fetch next key and
64+
// merge until we get a non-empty merged result or no more keys
65+
if (merged_key_value_.has_value()) {
66+
return true;
6067
}
61-
PAIMON_RETURN_NOT_OK(reader_->merge_function_wrapper_->Add(std::move(*next_key_value_)));
62-
next_key_value_.reset();
6368
}
69+
}
70+
71+
Result<bool> MergedKeyValueRecordReader::Iterator::MergeNextKey() const {
72+
PAIMON_RETURN_NOT_OK(LoadNextRawKeyValue());
73+
if (!next_raw_key_value_.has_value()) {
74+
return false;
75+
}
76+
77+
reader_->merge_function_wrapper_->Reset();
78+
auto current_key = next_raw_key_value_->key;
79+
80+
do {
81+
PAIMON_RETURN_NOT_OK(
82+
reader_->merge_function_wrapper_->Add(std::move(*next_raw_key_value_)));
83+
next_raw_key_value_.reset();
84+
PAIMON_RETURN_NOT_OK(LoadNextRawKeyValue());
85+
} while (next_raw_key_value_.has_value() &&
86+
reader_->key_comparator_->CompareTo(*current_key, *next_raw_key_value_->key) == 0);
6487

6588
PAIMON_ASSIGN_OR_RAISE(std::optional<KeyValue> result,
6689
reader_->merge_function_wrapper_->GetResult());
67-
// TODO(jinli.zjw): support merge function producing no result (e.g. all rows are filtered out)
68-
if (result == std::nullopt) {
69-
return Status::Invalid("merged key value reader produced empty result");
70-
}
71-
return std::move(result).value();
90+
merged_key_value_ = std::move(result);
91+
return true;
7292
}
7393

74-
Status MergedKeyValueRecordReader::Iterator::LoadNextKeyValue() {
75-
if (next_key_value_.has_value()) {
94+
Status MergedKeyValueRecordReader::Iterator::LoadNextRawKeyValue() const {
95+
if (next_raw_key_value_.has_value() || eof_) {
7696
return Status::OK();
7797
}
7898

7999
while (true) {
80-
if (current_iterator_ != nullptr && current_iterator_->HasNext()) {
81-
PAIMON_ASSIGN_OR_RAISE(KeyValue key_value, current_iterator_->Next());
82-
next_key_value_.emplace(std::move(key_value));
83-
return Status::OK();
100+
if (current_iterator_ != nullptr) {
101+
PAIMON_ASSIGN_OR_RAISE(bool has_next, current_iterator_->HasNext());
102+
if (has_next) {
103+
PAIMON_ASSIGN_OR_RAISE(KeyValue key_value, current_iterator_->Next());
104+
next_raw_key_value_.emplace(std::move(key_value));
105+
return Status::OK();
106+
}
84107
}
85108

86109
current_iterator_.reset();
87110
PAIMON_ASSIGN_OR_RAISE(current_iterator_, reader_->reader_->NextBatch());
88111
if (current_iterator_ == nullptr) {
112+
eof_ = true;
89113
return Status::OK();
90114
}
91115
}
@@ -97,8 +121,9 @@ Result<std::unique_ptr<KeyValueRecordReader::Iterator>> MergedKeyValueRecordRead
97121
}
98122
visited_ = true;
99123

100-
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<Iterator> iterator, Iterator::Create(this));
101-
if (!iterator->HasNext()) {
124+
auto iterator = std::make_unique<Iterator>(this);
125+
PAIMON_ASSIGN_OR_RAISE(bool has_next, iterator->HasNext());
126+
if (!has_next) {
102127
return std::unique_ptr<KeyValueRecordReader::Iterator>();
103128
}
104129
return iterator;

src/paimon/core/io/merged_key_value_record_reader.h

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -34,23 +34,25 @@ class MergedKeyValueRecordReader : public KeyValueRecordReader {
3434

3535
class Iterator : public KeyValueRecordReader::Iterator {
3636
public:
37-
static Result<std::unique_ptr<Iterator>> Create(MergedKeyValueRecordReader* reader);
37+
explicit Iterator(MergedKeyValueRecordReader* reader) : reader_(reader) {}
3838

39-
bool HasNext() const override {
40-
return next_key_value_.has_value();
41-
}
39+
Result<bool> HasNext() const override;
4240

4341
Result<KeyValue> Next() override;
4442

4543
private:
46-
explicit Iterator(MergedKeyValueRecordReader* reader) : reader_(reader) {}
47-
48-
Status LoadNextKeyValue();
44+
Result<bool> PrepareNextMergedKeyValue() const;
45+
Result<bool> MergeNextKey() const;
46+
Status LoadNextRawKeyValue() const;
4947

5048
private:
5149
MergedKeyValueRecordReader* reader_;
52-
std::unique_ptr<KeyValueRecordReader::Iterator> current_iterator_;
53-
std::optional<KeyValue> next_key_value_;
50+
mutable std::unique_ptr<KeyValueRecordReader::Iterator> current_iterator_;
51+
// Lookahead raw kv used to detect the boundary between two keys.
52+
mutable std::optional<KeyValue> next_raw_key_value_;
53+
// Merged kv prepared by HasNext() and consumed by Next().
54+
mutable std::optional<KeyValue> merged_key_value_;
55+
mutable bool eof_ = false;
5456
};
5557

5658
Result<std::unique_ptr<KeyValueRecordReader::Iterator>> NextBatch() override;

src/paimon/core/io/merged_key_value_record_reader_test.cpp

Lines changed: 52 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
#include "arrow/array/array_nested.h"
2424
#include "arrow/ipc/json_simple.h"
2525
#include "gtest/gtest.h"
26+
#include "paimon/common/table/special_fields.h"
2627
#include "paimon/common/types/data_field.h"
2728
#include "paimon/common/utils/fields_comparator.h"
2829
#include "paimon/core/mergetree/compact/deduplicate_merge_function.h"
@@ -55,14 +56,11 @@ TEST_F(MergedKeyValueRecordReaderTest, TestMergeAcrossUnderlyingBatches) {
5556
DataField(3, arrow::field("v1", arrow::int32())),
5657
DataField(4, arrow::field("v2", arrow::int32()))};
5758

58-
auto key_schema = arrow::schema({fields[0].ArrowField(), fields[1].ArrowField()});
59-
auto value_schema =
60-
arrow::schema({fields[0].ArrowField(), fields[1].ArrowField(), fields[2].ArrowField(),
61-
fields[3].ArrowField(), fields[4].ArrowField()});
62-
std::shared_ptr<arrow::DataType> src_type = arrow::struct_(
63-
{arrow::field("_SEQUENCE_NUMBER", arrow::int64()),
64-
arrow::field("_VALUE_KIND", arrow::int8()), fields[0].ArrowField(), fields[1].ArrowField(),
65-
fields[2].ArrowField(), fields[3].ArrowField(), fields[4].ArrowField()});
59+
auto value_schema = DataField::ConvertDataFieldsToArrowSchema(fields);
60+
auto arrow_fields = value_schema->fields();
61+
auto key_schema = arrow::schema({arrow_fields[0], arrow_fields[1]});
62+
std::shared_ptr<arrow::DataType> src_type =
63+
arrow::struct_(SpecialFields::CompleteSequenceAndValueKindField(value_schema)->fields());
6664

6765
auto src_array = std::dynamic_pointer_cast<arrow::StructArray>(
6866
arrow::ipc::internal::json::ArrayFromJSON(src_type, R"([
@@ -97,4 +95,50 @@ TEST_F(MergedKeyValueRecordReaderTest, TestMergeAcrossUnderlyingBatches) {
9795
}
9896
}
9997

98+
TEST_F(MergedKeyValueRecordReaderTest, TestSkipMergedNulloptResultInHasNext) {
99+
auto mfunc = std::make_unique<DeduplicateMergeFunction>(/*ignore_delete=*/true);
100+
auto merge_function_wrapper = std::make_shared<ReducerMergeFunctionWrapper>(std::move(mfunc));
101+
102+
std::vector<DataField> fields = {DataField(0, arrow::field("k0", arrow::int32())),
103+
DataField(1, arrow::field("v0", arrow::int32()))};
104+
105+
auto value_schema = DataField::ConvertDataFieldsToArrowSchema(fields);
106+
auto arrow_fields = value_schema->fields();
107+
auto key_schema = arrow::schema({arrow_fields[0]});
108+
std::shared_ptr<arrow::DataType> src_type =
109+
arrow::struct_(SpecialFields::CompleteSequenceAndValueKindField(value_schema)->fields());
110+
111+
auto src_array = std::dynamic_pointer_cast<arrow::StructArray>(
112+
arrow::ipc::internal::json::ArrayFromJSON(src_type, R"([
113+
[0, 3, 1, 10],
114+
[3, 3, 1, 11],
115+
[2, 3, 2, 200],
116+
[4, 3, 2, 240],
117+
[1, 0, 3, 300],
118+
[5, 3, 3, 30]
119+
])")
120+
.ValueOrDie());
121+
122+
ASSERT_OK_AND_ASSIGN(std::shared_ptr<FieldsComparator> key_comparator,
123+
FieldsComparator::Create({fields[0]}, /*is_ascending_order=*/true));
124+
125+
auto expected = KeyValueChecker::GenerateKeyValues(
126+
/*seq_vec=*/{1}, /*key_vec=*/{{3}}, /*value_vec=*/{{3, 300}}, pool_);
127+
128+
for (auto batch_size : {1, 2, 3}) {
129+
auto file_batch_reader =
130+
std::make_unique<MockFileBatchReader>(src_array, src_type, /*batch_size=*/batch_size);
131+
auto raw_reader = std::make_unique<MockKeyValueDataFileRecordReader>(
132+
std::move(file_batch_reader), key_schema, value_schema, /*level=*/0, pool_);
133+
auto merged_reader = std::make_unique<MergedKeyValueRecordReader>(
134+
std::move(raw_reader), key_comparator, merge_function_wrapper);
135+
136+
ASSERT_OK_AND_ASSIGN(
137+
auto results,
138+
(ReadResultCollector::CollectKeyValueResult<
139+
MergedKeyValueRecordReader, KeyValueRecordReader::Iterator>(merged_reader.get())));
140+
KeyValueChecker::CheckResult(expected, results, /*key_arity=*/1, /*value_arity=*/2);
141+
}
142+
}
143+
100144
} // namespace paimon::test

src/paimon/core/mergetree/compact/loser_tree.h

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -116,15 +116,23 @@ class LoserTree {
116116
Status AdvanceIfAvailable() {
117117
first_same_key_index = -1;
118118
state = State::WINNER_WITH_NEW_KEY;
119-
if (iterator == nullptr || !iterator->HasNext()) {
119+
bool has_next = false;
120+
if (iterator != nullptr) {
121+
PAIMON_ASSIGN_OR_RAISE(has_next, iterator->HasNext());
122+
}
123+
if (iterator == nullptr || !has_next) {
120124
while (!end_of_input) {
121125
PAIMON_ASSIGN_OR_RAISE(iterator, reader->NextBatch());
122126
if (!iterator) {
123127
// read eof
124128
reader->Close();
125129
end_of_input = true;
126130
kv = std::nullopt;
127-
} else if (iterator->HasNext()) {
131+
} else {
132+
PAIMON_ASSIGN_OR_RAISE(has_next, iterator->HasNext());
133+
if (!has_next) {
134+
continue;
135+
}
128136
PAIMON_ASSIGN_OR_RAISE(kv, iterator->Next());
129137
break;
130138
}

0 commit comments

Comments
 (0)