Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions src/paimon/common/data/generic_row.h
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,15 @@ class GenericRow : public InternalRow {
return row;
}

void ResetFields() {
for (auto& field : fields_) {
field = VariantType{};
}
kind_ = RowKind::Insert();
row_holder_.clear();
bytes_holder_.reset();
}

private:
/// The array to store the actual internal format values.
std::vector<VariantType> fields_;
Expand Down
31 changes: 31 additions & 0 deletions src/paimon/common/data/generic_row_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -104,4 +104,35 @@ TEST(GenericRowTest, TestSimple) {
"00:00:00.100000020,123.45678998765432145678,array,row,map,null)",
row.ToString());
}

TEST(GenericRowTest, TestResetFields) {
auto pool = GetDefaultPool();
GenericRow row(3);
row.SetField(0, static_cast<int32_t>(1));
row.SetField(1, BinaryString::FromString("old", pool.get()));
row.SetField(2, Bytes::AllocateBytes("bytes", pool.get()));
row.SetRowKind(RowKind::Delete());

ASSERT_FALSE(row.IsNullAt(0));
ASSERT_FALSE(row.IsNullAt(1));
ASSERT_FALSE(row.IsNullAt(2));
ASSERT_EQ(row.GetRowKind().value(), RowKind::Delete());

row.ResetFields();

ASSERT_EQ(row.GetFieldCount(), 3);
ASSERT_TRUE(row.IsNullAt(0));
ASSERT_TRUE(row.IsNullAt(1));
ASSERT_TRUE(row.IsNullAt(2));
ASSERT_EQ(row.GetRowKind().value(), RowKind::Insert());

row.SetField(0, static_cast<int32_t>(2));
row.SetField(1, BinaryString::FromString("new", pool.get()));
row.SetRowKind(RowKind::UpdateAfter());

ASSERT_EQ(row.GetInt(0), static_cast<int32_t>(2));
ASSERT_EQ(row.GetString(1), BinaryString::FromString("new", pool.get()));
ASSERT_TRUE(row.IsNullAt(2));
ASSERT_EQ(row.GetRowKind().value(), RowKind::UpdateAfter());
}
} // namespace paimon::test
4 changes: 2 additions & 2 deletions src/paimon/core/io/key_value_data_file_record_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -50,11 +50,11 @@ KeyValueDataFileRecordReader::KeyValueDataFileRecordReader(

Result<bool> KeyValueDataFileRecordReader::Iterator::HasNext() const {
int64_t array_length = reader_->row_kind_array_->length();
const auto& selection_bitmap = reader_->selection_bitmap_;
if (selection_bitmap.Cardinality() == array_length) {
if (selection_cardinality_ == array_length) {
// all rows are selected in bitmap
return cursor_ < array_length;
}
const auto& selection_bitmap = reader_->selection_bitmap_;
auto iter = selection_bitmap.EqualOrLarger(cursor_);
if (iter == selection_bitmap.End()) {
// no row are selected
Expand Down
5 changes: 4 additions & 1 deletion src/paimon/core/io/key_value_data_file_record_reader.h
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,9 @@ class KeyValueDataFileRecordReader : public KeyValueRecordReader {
class Iterator : public KeyValueRecordReader::Iterator {
public:
Iterator(KeyValueDataFileRecordReader* reader, int64_t previous_batch_first_row_number)
: previous_batch_first_row_number_(previous_batch_first_row_number), reader_(reader) {}
Comment thread
dalingmeng marked this conversation as resolved.
: previous_batch_first_row_number_(previous_batch_first_row_number),
reader_(reader),
selection_cardinality_(reader->selection_bitmap_.Cardinality()) {}
Result<bool> HasNext() const override;
Result<KeyValue> Next() override;
Result<std::pair<int64_t, KeyValue>> NextWithFilePos();
Expand All @@ -64,6 +66,7 @@ class KeyValueDataFileRecordReader : public KeyValueRecordReader {
int64_t previous_batch_first_row_number_;
mutable int64_t cursor_ = 0;
KeyValueDataFileRecordReader* reader_ = nullptr;
int64_t selection_cardinality_ = 0;
};

Result<std::unique_ptr<KeyValueRecordReader::Iterator>> NextBatch() override;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,11 @@ Status AggregateMergeFunction::Add(KeyValue&& kv) {
// mark the current row for deletion and initialize the row with input values.
if (remove_record_on_delete_ && kv.value_kind == RowKind::Delete()) {
current_delete_row_ = true;
row_ = std::make_unique<GenericRow>(getters_.size());
if (row_) {
row_->ResetFields();
} else {
row_ = std::make_unique<GenericRow>(getters_.size());
}
for (size_t i = 0; i < getters_.size(); i++) {
row_->SetField(i, getters_[i](*(kv.value)));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,11 @@ class AggregateMergeFunction : public MergeFunction {
void Reset() override {
latest_kv_ = std::nullopt;
current_delete_row_ = false;
row_ = std::make_unique<GenericRow>(getters_.size());
if (row_) {
row_->ResetFields();
} else {
row_ = std::make_unique<GenericRow>(getters_.size());
}
for (const auto& agg : aggregators_) {
agg->Reset();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -269,7 +269,11 @@ void PartialUpdateMergeFunction::Reset() {
current_key_.reset();
meet_insert_ = false;
not_null_column_filled_ = false;
row_ = std::make_unique<GenericRow>(getters_.size());
if (row_) {
row_->ResetFields();
} else {
row_ = std::make_unique<GenericRow>(getters_.size());
}
last_seq_num_ = 0;
for (auto& [_, agg] : field_aggregators_) {
assert(agg);
Expand Down Expand Up @@ -304,7 +308,11 @@ Status PartialUpdateMergeFunction::Add(KeyValue&& moved_kv) {
if (remove_record_on_delete_) {
if (kv.value_kind == RowKind::Delete()) {
current_delete_row_ = true;
row_ = std::make_unique<GenericRow>(getters_.size());
if (row_) {
row_->ResetFields();
} else {
row_ = std::make_unique<GenericRow>(getters_.size());
}
InitRowAndHoldData(std::move(kv.value));
} else if (!not_null_column_filled_) {
InitRowAndHoldData(std::move(kv.value));
Expand Down
Loading