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"
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"
3835namespace paimon {
3936class MemoryPool ;
4037
4138Result<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
7053KeyValueInMemoryRecordReader::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
11598void 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);
0 commit comments