Skip to content

Commit a47d52d

Browse files
committed
fix review
1 parent bdf52d4 commit a47d52d

4 files changed

Lines changed: 45 additions & 46 deletions

File tree

src/paimon/core/mergetree/merge_tree_writer.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -226,7 +226,7 @@ Status MergeTreeWriter::Flush(bool wait_for_latest_compaction, bool forced_full_
226226
}
227227
// 1. flush write buffer to get in-memory readers
228228
PAIMON_ASSIGN_OR_RAISE(std::vector<std::unique_ptr<KeyValueRecordReader>> readers,
229-
write_buffer_->Flush(&last_sequence_number_));
229+
write_buffer_->DrainToReaders(&last_sequence_number_));
230230
// 2. prepare loser tree sort merge reader
231231
auto sort_merge_reader = std::make_unique<SortMergeReaderWithLoserTree>(
232232
std::move(readers), key_comparator_, user_defined_seq_comparator_,

src/paimon/core/mergetree/write_buffer.cpp

Lines changed: 7 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -39,12 +39,12 @@ WriteBuffer::WriteBuffer(
3939
const std::shared_ptr<FieldsComparator>& key_comparator,
4040
const std::shared_ptr<MergeFunctionWrapper<KeyValue>>& merge_function_wrapper,
4141
const std::shared_ptr<MemoryPool>& pool)
42-
: value_type_(value_type),
42+
: pool_(pool),
43+
value_type_(value_type),
4344
trimmed_primary_keys_(trimmed_primary_keys),
4445
user_defined_sequence_fields_(user_defined_sequence_fields),
4546
key_comparator_(key_comparator),
46-
merge_function_wrapper_(merge_function_wrapper),
47-
pool_(pool) {}
47+
merge_function_wrapper_(merge_function_wrapper) {}
4848

4949
Status WriteBuffer::Write(std::unique_ptr<RecordBatch>&& moved_batch) {
5050
if (ArrowArrayIsReleased(moved_batch->GetData())) {
@@ -66,7 +66,7 @@ Status WriteBuffer::Write(std::unique_ptr<RecordBatch>&& moved_batch) {
6666
return Status::OK();
6767
}
6868

69-
Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> WriteBuffer::Flush(
69+
Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> WriteBuffer::DrainToReaders(
7070
int64_t* last_sequence_number) {
7171
std::vector<std::unique_ptr<KeyValueRecordReader>> readers;
7272
if (batch_vec_.empty()) {
@@ -84,10 +84,7 @@ Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> WriteBuffer::Flush(
8484
readers.push_back(std::move(in_memory_reader));
8585
}
8686

87-
batch_vec_.clear();
88-
row_kinds_vec_.clear();
89-
current_memory_in_bytes_ = 0;
90-
87+
Clear();
9188
return readers;
9289
}
9390

@@ -97,6 +94,8 @@ void WriteBuffer::Clear() {
9794
current_memory_in_bytes_ = 0;
9895
}
9996

97+
// TODO(jinli.zjw): Consider making the memory estimation more accurate.
98+
// https://github.com/alibaba/paimon-cpp/pull/206#discussion_r3021325389
10099
Result<int64_t> WriteBuffer::EstimateMemoryUse(const std::shared_ptr<arrow::Array>& array) {
101100
arrow::Type::type type = array->type()->id();
102101
int64_t null_bits_size_in_bytes = (array->length() + 7) / 8;

src/paimon/core/mergetree/write_buffer.h

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -57,10 +57,11 @@ class WriteBuffer {
5757
/// Does NOT check memory thresholds or trigger flush.
5858
Status Write(std::unique_ptr<RecordBatch>&& batch);
5959

60-
/// Flush all buffered batches into KeyValueInMemoryRecordReaders and clear the buffer.
61-
/// @param[in,out] last_sequence_number current sequence number, updated after flush
60+
/// Drain all buffered batches into KeyValueInMemoryRecordReaders and clear the buffer.
61+
/// @param[in,out] last_sequence_number current sequence number, updated after draining
6262
/// @return list of KeyValueRecordReaders built from buffered data
63-
Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> Flush(int64_t* last_sequence_number);
63+
Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> DrainToReaders(
64+
int64_t* last_sequence_number);
6465

6566
/// Return current memory usage in bytes.
6667
int64_t GetMemoryUsage() const {
@@ -80,12 +81,12 @@ class WriteBuffer {
8081
static Result<int64_t> EstimateMemoryUse(const std::shared_ptr<arrow::Array>& array);
8182

8283
// Immutable configuration
84+
const std::shared_ptr<MemoryPool> pool_;
8385
const std::shared_ptr<arrow::DataType> value_type_;
8486
const std::vector<std::string> trimmed_primary_keys_;
8587
const std::vector<std::string> user_defined_sequence_fields_;
8688
const std::shared_ptr<FieldsComparator> key_comparator_;
8789
const std::shared_ptr<MergeFunctionWrapper<KeyValue>> merge_function_wrapper_;
88-
const std::shared_ptr<MemoryPool> pool_;
8990

9091
// Mutable buffer state
9192
std::vector<std::shared_ptr<arrow::StructArray>> batch_vec_;

src/paimon/core/mergetree/write_buffer_test.cpp

Lines changed: 32 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,11 @@ class MergeFunctionWrapper;
3939
} // namespace paimon
4040

4141
namespace paimon::test {
42+
struct ReaderResult {
43+
std::vector<int64_t> sequence_numbers;
44+
std::vector<int8_t> row_kind_values;
45+
};
46+
4247
class WriteBufferTest : public ::testing::Test {
4348
public:
4449
void SetUp() override {
@@ -69,6 +74,18 @@ class WriteBufferTest : public ::testing::Test {
6974
return batch;
7075
}
7176

77+
Result<ReaderResult> ReadReaderResult(KeyValueRecordReader* reader) const {
78+
PAIMON_ASSIGN_OR_RAISE(auto iterator, reader->NextBatch());
79+
80+
ReaderResult result;
81+
while (iterator->HasNext()) {
82+
PAIMON_ASSIGN_OR_RAISE(KeyValue key_value, iterator->Next());
83+
result.sequence_numbers.push_back(key_value.sequence_number);
84+
result.row_kind_values.push_back(key_value.value_kind->ToByteValue());
85+
}
86+
return result;
87+
}
88+
7289
protected:
7390
std::shared_ptr<MemoryPool> pool_;
7491
std::vector<DataField> value_fields_;
@@ -100,32 +117,23 @@ TEST_F(WriteBufferTest, TestFlushResetsStateAndAdvancesSequenceNumber) {
100117
ASSERT_GT(write_buffer.GetMemoryUsage(), 0);
101118

102119
int64_t last_sequence_number = 10;
103-
ASSERT_OK_AND_ASSIGN(auto readers, write_buffer.Flush(&last_sequence_number));
120+
ASSERT_OK_AND_ASSIGN(auto readers, write_buffer.DrainToReaders(&last_sequence_number));
104121

105122
ASSERT_EQ(readers.size(), 2);
106123
ASSERT_TRUE(write_buffer.IsEmpty());
107124
ASSERT_EQ(write_buffer.GetMemoryUsage(), 0);
108125
ASSERT_EQ(last_sequence_number, 13);
109126

110-
ASSERT_OK_AND_ASSIGN(auto first_iterator, readers[0]->NextBatch());
111-
ASSERT_TRUE(first_iterator);
112-
ASSERT_TRUE(first_iterator->HasNext());
113-
ASSERT_OK_AND_ASSIGN(KeyValue first, first_iterator->Next());
114-
ASSERT_EQ(first.sequence_number, 10);
115-
ASSERT_EQ(first.value_kind->ToByteValue(), RowKind::Insert()->ToByteValue());
116-
ASSERT_TRUE(first_iterator->HasNext());
117-
ASSERT_OK_AND_ASSIGN(KeyValue second, first_iterator->Next());
118-
ASSERT_EQ(second.sequence_number, 11);
119-
ASSERT_EQ(second.value_kind->ToByteValue(), RowKind::Insert()->ToByteValue());
120-
ASSERT_FALSE(first_iterator->HasNext());
121-
122-
ASSERT_OK_AND_ASSIGN(auto second_iterator, readers[1]->NextBatch());
123-
ASSERT_TRUE(second_iterator);
124-
ASSERT_TRUE(second_iterator->HasNext());
125-
ASSERT_OK_AND_ASSIGN(KeyValue third, second_iterator->Next());
126-
ASSERT_EQ(third.sequence_number, 12);
127-
ASSERT_EQ(third.value_kind->ToByteValue(), RowKind::Insert()->ToByteValue());
128-
ASSERT_FALSE(second_iterator->HasNext());
127+
ASSERT_OK_AND_ASSIGN(auto first_result, ReadReaderResult(readers[0].get()));
128+
ASSERT_EQ(first_result.sequence_numbers, (std::vector<int64_t>{10, 11}));
129+
ASSERT_EQ(
130+
first_result.row_kind_values,
131+
(std::vector<int8_t>{RowKind::Insert()->ToByteValue(), RowKind::Insert()->ToByteValue()}));
132+
133+
ASSERT_OK_AND_ASSIGN(auto second_result, ReadReaderResult(readers[1].get()));
134+
ASSERT_EQ(second_result.sequence_numbers, (std::vector<int64_t>{12}));
135+
ASSERT_EQ(second_result.row_kind_values,
136+
(std::vector<int8_t>{RowKind::Insert()->ToByteValue()}));
129137
}
130138

131139
TEST_F(WriteBufferTest, TestFlushPreservesRowKinds) {
@@ -150,26 +158,17 @@ TEST_F(WriteBufferTest, TestFlushPreservesRowKinds) {
150158
ASSERT_OK(write_buffer.Write(CreateBatch(array, row_kinds)));
151159

152160
int64_t last_sequence_number = 0;
153-
ASSERT_OK_AND_ASSIGN(auto readers, write_buffer.Flush(&last_sequence_number));
161+
ASSERT_OK_AND_ASSIGN(auto readers, write_buffer.DrainToReaders(&last_sequence_number));
154162
ASSERT_EQ(readers.size(), 1);
155163
ASSERT_EQ(last_sequence_number, 4);
156164

157-
ASSERT_OK_AND_ASSIGN(auto iterator, readers[0]->NextBatch());
158-
ASSERT_TRUE(iterator);
159-
160-
std::vector<int8_t> actual_row_kind_values;
161-
std::vector<int64_t> actual_sequence_numbers;
162-
while (iterator->HasNext()) {
163-
ASSERT_OK_AND_ASSIGN(KeyValue key_value, iterator->Next());
164-
actual_row_kind_values.push_back(key_value.value_kind->ToByteValue());
165-
actual_sequence_numbers.push_back(key_value.sequence_number);
166-
}
165+
ASSERT_OK_AND_ASSIGN(auto reader_result, ReadReaderResult(readers[0].get()));
167166

168-
ASSERT_EQ(actual_row_kind_values,
167+
ASSERT_EQ(reader_result.row_kind_values,
169168
(std::vector<int8_t>{
170169
RowKind::Insert()->ToByteValue(), RowKind::UpdateBefore()->ToByteValue(),
171170
RowKind::UpdateAfter()->ToByteValue(), RowKind::Delete()->ToByteValue()}));
172-
ASSERT_EQ(actual_sequence_numbers, (std::vector<int64_t>{0, 1, 2, 3}));
171+
ASSERT_EQ(reader_result.sequence_numbers, (std::vector<int64_t>{0, 1, 2, 3}));
173172
}
174173

175174
TEST_F(WriteBufferTest, TestEstimateMemoryUse) {

0 commit comments

Comments
 (0)