Skip to content

Commit ff02c94

Browse files
committed
fix
1 parent 1176255 commit ff02c94

4 files changed

Lines changed: 10 additions & 10 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_->Flush(&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: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
/*
2-
* Copyright 2024-present Alibaba Inc.
2+
* Copyright 2026-present Alibaba Inc.
33
*
44
* Licensed under the Apache License, Version 2.0 (the "License");
55
* you may not use this file except in compliance with the License.
@@ -67,16 +67,16 @@ Status WriteBuffer::Write(std::unique_ptr<RecordBatch>&& moved_batch) {
6767
}
6868

6969
Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> WriteBuffer::Flush(
70-
int64_t& last_sequence_number) {
70+
int64_t* last_sequence_number) {
7171
std::vector<std::unique_ptr<KeyValueRecordReader>> readers;
7272
if (batch_vec_.empty()) {
7373
return readers;
7474
}
7575

7676
readers.reserve(batch_vec_.size());
7777
for (size_t i = 0; i < batch_vec_.size(); ++i) {
78-
int64_t sequence_number = last_sequence_number;
79-
last_sequence_number += batch_vec_[i]->length();
78+
int64_t sequence_number = *last_sequence_number;
79+
*last_sequence_number += batch_vec_[i]->length();
8080
auto in_memory_reader = std::make_unique<KeyValueInMemoryRecordReader>(
8181
sequence_number, std::move(batch_vec_[i]), std::move(row_kinds_vec_[i]),
8282
trimmed_primary_keys_, user_defined_sequence_fields_, key_comparator_,

src/paimon/core/mergetree/write_buffer.h

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
/*
2-
* Copyright 2024-present Alibaba Inc.
2+
* Copyright 2026-present Alibaba Inc.
33
*
44
* Licensed under the Apache License, Version 2.0 (the "License");
55
* you may not use this file except in compliance with the License.
@@ -61,7 +61,7 @@ class WriteBuffer {
6161
/// Flush all buffered batches into KeyValueInMemoryRecordReaders and clear the buffer.
6262
/// @param[in,out] last_sequence_number current sequence number, updated after flush
6363
/// @return list of KeyValueRecordReaders built from buffered data
64-
Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> Flush(int64_t& last_sequence_number);
64+
Result<std::vector<std::unique_ptr<KeyValueRecordReader>>> Flush(int64_t* last_sequence_number);
6565

6666
/// Return current memory usage in bytes.
6767
int64_t GetMemoryUsage() const {

src/paimon/core/mergetree/write_buffer_test.cpp

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
/*
2-
* Copyright 2024-present Alibaba Inc.
2+
* Copyright 2026-present Alibaba Inc.
33
*
44
* Licensed under the Apache License, Version 2.0 (the "License");
55
* you may not use this file except in compliance with the License.
@@ -100,7 +100,7 @@ TEST_F(WriteBufferTest, TestFlushResetsStateAndAdvancesSequenceNumber) {
100100
ASSERT_GT(write_buffer.GetMemoryUsage(), 0);
101101

102102
int64_t last_sequence_number = 10;
103-
ASSERT_OK_AND_ASSIGN(auto readers, write_buffer.Flush(last_sequence_number));
103+
ASSERT_OK_AND_ASSIGN(auto readers, write_buffer.Flush(&last_sequence_number));
104104

105105
ASSERT_EQ(readers.size(), 2);
106106
ASSERT_TRUE(write_buffer.IsEmpty());
@@ -150,7 +150,7 @@ TEST_F(WriteBufferTest, TestFlushPreservesRowKinds) {
150150
ASSERT_OK(write_buffer.Write(CreateBatch(array, row_kinds)));
151151

152152
int64_t last_sequence_number = 0;
153-
ASSERT_OK_AND_ASSIGN(auto readers, write_buffer.Flush(last_sequence_number));
153+
ASSERT_OK_AND_ASSIGN(auto readers, write_buffer.Flush(&last_sequence_number));
154154
ASSERT_EQ(readers.size(), 1);
155155
ASSERT_EQ(last_sequence_number, 4);
156156

0 commit comments

Comments
 (0)