Skip to content

Commit 18e2454

Browse files
committed
refactor(mergetree): introduce BinaryInMemorySortBuffer and MergedKeyValueRecordReader
Refactor write_buffer to use the new SortBuffer abstraction (BinaryInMemorySortBuffer / BinaryExternalSortBuffer) and introduce MergedKeyValueRecordReader for merged key-value reading. Also adapt key_value_in_memory_record_reader, merge_tree_writer, spill reader/writer and core_options accordingly.
1 parent e7ffea3 commit 18e2454

20 files changed

Lines changed: 837 additions & 241 deletions

src/paimon/CMakeLists.txt

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -213,6 +213,7 @@ set(PAIMON_CORE_SRCS
213213
core/io/key_value_data_file_record_reader.cpp
214214
core/io/key_value_data_file_writer.cpp
215215
core/io/key_value_in_memory_record_reader.cpp
216+
core/io/merged_key_value_record_reader.cpp
216217
core/io/key_value_meta_projection_consumer.cpp
217218
core/io/key_value_projection_consumer.cpp
218219
core/io/key_value_projection_reader.cpp
@@ -246,6 +247,7 @@ set(PAIMON_CORE_SRCS
246247
core/mergetree/compact/lookup_merge_tree_compact_rewriter.cpp
247248
core/mergetree/compact/changelog_merge_tree_rewriter.cpp
248249
core/mergetree/merge_tree_writer.cpp
250+
core/mergetree/binary_in_memory_sort_buffer.cpp
249251
core/mergetree/write_buffer.cpp
250252
core/mergetree/levels.cpp
251253
core/mergetree/lookup_file.cpp
@@ -564,6 +566,7 @@ if(PAIMON_BUILD_TESTS)
564566
core/io/key_value_data_file_record_reader_test.cpp
565567
core/io/key_value_projection_reader_test.cpp
566568
core/io/key_value_in_memory_record_reader_test.cpp
569+
core/io/merged_key_value_record_reader_test.cpp
567570
core/io/complete_row_tracking_fields_reader_test.cpp
568571
core/io/data_file_meta_test.cpp
569572
core/io/file_index_evaluator_test.cpp

src/paimon/core/append/append_only_writer.h

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,10 @@ class AppendOnlyWriter : public BatchWriter {
6868
Status Compact(bool full_compaction) override {
6969
return Flush(/*wait_for_latest_compaction=*/true, full_compaction);
7070
}
71+
72+
Status FlushMemory() override {
73+
return Flush(/*wait_for_latest_compaction=*/false, /*forced_full_compaction=*/false);
74+
}
7175
Result<CommitIncrement> PrepareCommit(bool wait_compaction) override;
7276
Result<bool> CompactNotCompleted() override {
7377
PAIMON_RETURN_NOT_OK(compact_manager_->TriggerCompaction(/*full_compaction=*/false));

src/paimon/core/io/key_value_in_memory_record_reader.cpp

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

1919
#include <cassert>
20-
#include <optional>
2120
#include <utility>
2221

2322
#include "arrow/array/array_base.h"
@@ -28,60 +27,44 @@
2827
#include "arrow/util/checked_cast.h"
2928
#include "fmt/format.h"
3029
#include "paimon/common/data/columnar/columnar_row_ref.h"
31-
#include "paimon/common/data/internal_row.h"
3230
#include "paimon/common/types/row_kind.h"
3331
#include "paimon/common/utils/arrow/arrow_utils.h"
3432
#include "paimon/common/utils/arrow/status_utils.h"
3533
#include "paimon/common/utils/fields_comparator.h"
36-
#include "paimon/core/mergetree/compact/merge_function_wrapper.h"
3734
#include "paimon/status.h"
3835
namespace paimon {
3936
class MemoryPool;
4037

4138
Result<KeyValue> KeyValueInMemoryRecordReader::Iterator::Next() {
42-
reader_->merge_function_wrapper_->Reset();
43-
std::shared_ptr<InternalRow> current_key;
44-
while (cursor_ < reader_->value_struct_array_->length()) {
45-
uint64_t index = reader_->sort_indices_->Value(cursor_);
46-
const RowKind* row_kind = RowKind::Insert();
47-
if (!reader_->row_kinds_.empty()) {
48-
PAIMON_ASSIGN_OR_RAISE(
49-
row_kind, RowKind::FromByteValue(static_cast<int8_t>(reader_->row_kinds_[index])));
50-
}
51-
// key must hold value_struct_array as min/max key may be used after projection
52-
auto key = std::make_unique<ColumnarRowRef>(reader_->key_ctx_, index);
53-
auto value = std::make_unique<ColumnarRowRef>(reader_->value_ctx_, index);
54-
KeyValue kv(row_kind, reader_->last_sequence_num_ + index,
55-
/*level=*/KeyValue::UNKNOWN_LEVEL, std::move(key), std::move(value));
56-
if (current_key == nullptr) {
57-
current_key = kv.key;
58-
} else if (reader_->key_comparator_->CompareTo(*current_key, *kv.key) != 0) {
59-
break;
60-
}
61-
PAIMON_RETURN_NOT_OK(reader_->merge_function_wrapper_->Add(std::move(kv)));
62-
cursor_++;
39+
uint64_t index = reader_->sort_indices_->Value(cursor_++);
40+
const RowKind* row_kind = RowKind::Insert();
41+
if (!reader_->row_kinds_.empty()) {
42+
PAIMON_ASSIGN_OR_RAISE(
43+
row_kind, RowKind::FromByteValue(static_cast<int8_t>(reader_->row_kinds_[index])));
6344
}
64-
PAIMON_ASSIGN_OR_RAISE(std::optional<KeyValue> result,
65-
reader_->merge_function_wrapper_->GetResult());
66-
assert(result != std::nullopt);
67-
return std::move(result).value();
45+
46+
// key must hold value_struct_array as min/max key may be used after projection
47+
auto key = std::make_unique<ColumnarRowRef>(reader_->key_ctx_, index);
48+
auto value = std::make_unique<ColumnarRowRef>(reader_->value_ctx_, index);
49+
return KeyValue(row_kind, reader_->last_sequence_num_ + index,
50+
/*level=*/KeyValue::UNKNOWN_LEVEL, std::move(key), std::move(value));
6851
}
6952

7053
KeyValueInMemoryRecordReader::KeyValueInMemoryRecordReader(
71-
int64_t last_sequence_num, std::shared_ptr<arrow::StructArray>&& struct_array,
72-
std::vector<RecordBatch::RowKind>&& row_kinds, const std::vector<std::string>& primary_keys,
73-
const std::vector<std::string>& user_defined_sequence_fields,
54+
int64_t last_sequence_num, const std::shared_ptr<arrow::StructArray>& struct_array,
55+
const std::vector<RecordBatch::RowKind>& row_kinds,
56+
const std::vector<std::string>& primary_keys,
57+
const std::vector<std::string>& user_defined_sequence_fields, bool sequence_fields_ascending,
7458
const std::shared_ptr<FieldsComparator>& key_comparator,
75-
const std::shared_ptr<MergeFunctionWrapper<KeyValue>>& merge_function_wrapper,
7659
const std::shared_ptr<MemoryPool>& pool)
7760
: last_sequence_num_(last_sequence_num),
7861
primary_keys_(primary_keys),
7962
user_defined_sequence_fields_(user_defined_sequence_fields),
63+
sequence_fields_ascending_(sequence_fields_ascending),
8064
pool_(pool),
81-
value_struct_array_(std::move(struct_array)),
82-
row_kinds_(std::move(row_kinds)),
83-
key_comparator_(key_comparator),
84-
merge_function_wrapper_(merge_function_wrapper) {
65+
value_struct_array_(struct_array),
66+
row_kinds_(row_kinds),
67+
key_comparator_(key_comparator) {
8568
assert(value_struct_array_);
8669
ArrowUtils::TraverseArray(value_struct_array_);
8770
}
@@ -113,6 +96,7 @@ Result<std::unique_ptr<KeyValueRecordReader::Iterator>> KeyValueInMemoryRecordRe
11396
}
11497

11598
void KeyValueInMemoryRecordReader::Close() {
99+
visited_ = true;
116100
value_struct_array_.reset();
117101
row_kinds_.clear();
118102
sort_indices_.reset();
@@ -127,8 +111,11 @@ KeyValueInMemoryRecordReader::SortBatch() const {
127111
for (const auto& name : primary_keys_) {
128112
sort_keys.emplace_back(name, arrow::compute::SortOrder::Ascending);
129113
}
114+
const auto sequence_sort_order = sequence_fields_ascending_
115+
? arrow::compute::SortOrder::Ascending
116+
: arrow::compute::SortOrder::Descending;
130117
for (const auto& name : user_defined_sequence_fields_) {
131-
sort_keys.emplace_back(name, arrow::compute::SortOrder::Ascending);
118+
sort_keys.emplace_back(name, sequence_sort_order);
132119
}
133120
auto sort_options =
134121
arrow::compute::SortOptions(sort_keys, arrow::compute::NullPlacement::AtStart);

src/paimon/core/io/key_value_in_memory_record_reader.h

Lines changed: 9 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -28,7 +28,6 @@
2828
#include "paimon/common/utils/fields_comparator.h"
2929
#include "paimon/core/io/key_value_record_reader.h"
3030
#include "paimon/core/key_value.h"
31-
#include "paimon/core/mergetree/compact/merge_function_wrapper.h"
3231
#include "paimon/record_batch.h"
3332
#include "paimon/result.h"
3433

@@ -41,18 +40,17 @@ struct ColumnarBatchContext;
4140
class FieldsComparator;
4241
class MemoryPool;
4342
class Metrics;
44-
template <typename T>
45-
class MergeFunctionWrapper;
4643

4744
class KeyValueInMemoryRecordReader : public KeyValueRecordReader {
4845
public:
49-
KeyValueInMemoryRecordReader(
50-
int64_t last_sequence_num, std::shared_ptr<arrow::StructArray>&& struct_array,
51-
std::vector<RecordBatch::RowKind>&& row_kinds, const std::vector<std::string>& primary_keys,
52-
const std::vector<std::string>& user_defined_sequence_fields,
53-
const std::shared_ptr<FieldsComparator>& key_comparator,
54-
const std::shared_ptr<MergeFunctionWrapper<KeyValue>>& merge_function_wrapper,
55-
const std::shared_ptr<MemoryPool>& pool);
46+
KeyValueInMemoryRecordReader(int64_t last_sequence_num,
47+
const std::shared_ptr<arrow::StructArray>& struct_array,
48+
const std::vector<RecordBatch::RowKind>& row_kinds,
49+
const std::vector<std::string>& primary_keys,
50+
const std::vector<std::string>& user_defined_sequence_fields,
51+
bool sequence_fields_ascending,
52+
const std::shared_ptr<FieldsComparator>& key_comparator,
53+
const std::shared_ptr<MemoryPool>& pool);
5654

5755
class Iterator : public KeyValueRecordReader::Iterator {
5856
public:
@@ -83,11 +81,11 @@ class KeyValueInMemoryRecordReader : public KeyValueRecordReader {
8381
int64_t last_sequence_num_ = -1;
8482
std::vector<std::string> primary_keys_;
8583
std::vector<std::string> user_defined_sequence_fields_;
84+
bool sequence_fields_ascending_ = true;
8685
std::shared_ptr<MemoryPool> pool_;
8786
std::shared_ptr<arrow::StructArray> value_struct_array_;
8887
std::vector<RecordBatch::RowKind> row_kinds_;
8988
std::shared_ptr<FieldsComparator> key_comparator_;
90-
std::shared_ptr<MergeFunctionWrapper<KeyValue>> merge_function_wrapper_;
9189

9290
std::shared_ptr<arrow::NumericArray<arrow::UInt64Type>> sort_indices_;
9391
std::shared_ptr<ColumnarBatchContext> key_ctx_;

0 commit comments

Comments
 (0)