Skip to content

Commit fec911d

Browse files
committed
feat(changelog): support full-compaction mode changelog producer
1 parent 53f9c86 commit fec911d

15 files changed

Lines changed: 845 additions & 37 deletions

include/paimon/defs.h

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -390,7 +390,6 @@ struct PAIMON_EXPORT Options {
390390
/// keeps the details of data changes, it can be read directly during stream reads. This can be
391391
/// applied to tables with primary keys. Values can be "none", "input", "lookup",
392392
/// "full-compaction". Default value is "none".
393-
/// @note C++ Paimon currently supports "none", "input", and "lookup".
394393
static const char CHANGELOG_PRODUCER[];
395394

396395
/// "changelog-producer.row-deduplicate" - Whether to generate update-before and update-after

src/paimon/CMakeLists.txt

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -329,6 +329,7 @@ set(PAIMON_CORE_SRCS
329329
core/mergetree/compact/merge_tree_compact_manager_factory.cpp
330330
core/mergetree/compact/merge_tree_compact_rewriter.cpp
331331
core/mergetree/compact/merge_tree_compact_task.cpp
332+
core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.cpp
332333
core/mergetree/compact/partial_update_merge_function.cpp
333334
core/mergetree/compact/sort_merge_reader_with_loser_tree.cpp
334335
core/mergetree/compact/sort_merge_reader_with_min_heap.cpp
@@ -815,6 +816,7 @@ if(PAIMON_BUILD_TESTS)
815816
core/mergetree/compact/deduplicate_merge_function_test.cpp
816817
core/mergetree/compact/first_row_merge_function_test.cpp
817818
core/mergetree/compact/first_row_merge_function_wrapper_test.cpp
819+
core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp
818820
core/mergetree/compact/internal_row_equalizer_test.cpp
819821
core/mergetree/compact/interval_partition_test.cpp
820822
core/mergetree/compact/lookup_changelog_merge_function_wrapper_test.cpp
Lines changed: 142 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,142 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
#pragma once
21+
22+
#include <memory>
23+
#include <optional>
24+
#include <utility>
25+
26+
#include "paimon/common/data/serializer/row_compacted_serializer.h"
27+
#include "paimon/common/utils/fields_comparator.h"
28+
#include "paimon/core/key_value.h"
29+
#include "paimon/core/mergetree/compact/changelog_result.h"
30+
#include "paimon/core/mergetree/compact/merge_function.h"
31+
#include "paimon/core/mergetree/compact/merge_function_wrapper.h"
32+
#include "paimon/result.h"
33+
#include "paimon/status.h"
34+
35+
namespace paimon {
36+
37+
/// Wrapper for `MergeFunction`s which produces changelog during a full compaction.
38+
class FullChangelogMergeFunctionWrapper : public MergeFunctionWrapper<ChangelogResult> {
39+
public:
40+
FullChangelogMergeFunctionWrapper(std::unique_ptr<MergeFunction>&& merge_function,
41+
int32_t max_level,
42+
std::unique_ptr<RowCompactedSerializer>&& value_serializer,
43+
FieldsComparator::FieldComparatorFunc value_equalizer)
44+
: merge_function_(std::move(merge_function)),
45+
max_level_(max_level),
46+
value_serializer_(std::move(value_serializer)),
47+
value_equalizer_(std::move(value_equalizer)) {}
48+
49+
void Reset() override {
50+
merge_function_->Reset();
51+
top_level_kv_ = std::nullopt;
52+
initial_kv_ = std::nullopt;
53+
is_initialized_ = false;
54+
}
55+
56+
Status Add(KeyValue&& kv) override {
57+
if (!initial_kv_) {
58+
initial_kv_ = std::move(kv);
59+
return Status::OK();
60+
}
61+
62+
if (!is_initialized_) {
63+
if (initial_kv_->level == max_level_) {
64+
PAIMON_RETURN_NOT_OK(RememberTopLevel(*initial_kv_));
65+
}
66+
PAIMON_RETURN_NOT_OK(merge_function_->Add(std::move(initial_kv_).value()));
67+
is_initialized_ = true;
68+
}
69+
70+
if (kv.level == max_level_) {
71+
PAIMON_RETURN_NOT_OK(RememberTopLevel(kv));
72+
}
73+
return merge_function_->Add(std::move(kv));
74+
}
75+
76+
Result<std::optional<ChangelogResult>> GetResult() override {
77+
std::optional<KeyValue> merged;
78+
if (is_initialized_) {
79+
PAIMON_ASSIGN_OR_RAISE(merged, merge_function_->GetResult());
80+
} else {
81+
merged = std::move(initial_kv_);
82+
}
83+
84+
ChangelogResult result;
85+
if (is_initialized_) {
86+
if (!top_level_kv_) {
87+
if (merged && merged->value_kind->IsAdd()) {
88+
PAIMON_ASSIGN_OR_RAISE(KeyValue insert,
89+
CloneKeyValue(*merged, RowKind::Insert()));
90+
result.changelogs.emplace_back(std::move(insert));
91+
}
92+
} else if (!merged || !merged->value_kind->IsAdd()) {
93+
top_level_kv_->value_kind = RowKind::Delete();
94+
result.changelogs.emplace_back(std::move(top_level_kv_).value());
95+
} else if (!value_equalizer_ ||
96+
value_equalizer_(*top_level_kv_->value, *merged->value) != 0) {
97+
top_level_kv_->value_kind = RowKind::UpdateBefore();
98+
result.changelogs.emplace_back(std::move(top_level_kv_).value());
99+
PAIMON_ASSIGN_OR_RAISE(KeyValue update_after,
100+
CloneKeyValue(*merged, RowKind::UpdateAfter()));
101+
result.changelogs.emplace_back(std::move(update_after));
102+
}
103+
} else if (merged && merged->level != max_level_ && merged->value_kind->IsAdd()) {
104+
PAIMON_ASSIGN_OR_RAISE(KeyValue insert, CloneKeyValue(*merged, RowKind::Insert()));
105+
result.changelogs.emplace_back(std::move(insert));
106+
}
107+
108+
if (merged && merged->value_kind->IsAdd()) {
109+
result.result = std::move(merged);
110+
}
111+
Reset();
112+
return std::optional<ChangelogResult>(std::move(result));
113+
}
114+
115+
private:
116+
Status RememberTopLevel(const KeyValue& kv) {
117+
if (top_level_kv_) {
118+
return Status::Invalid("Top level key-value already exists. This is unexpected.");
119+
}
120+
PAIMON_ASSIGN_OR_RAISE(top_level_kv_, CloneKeyValue(kv, kv.value_kind));
121+
return Status::OK();
122+
}
123+
124+
Result<KeyValue> CloneKeyValue(const KeyValue& from, const RowKind* value_kind) const {
125+
PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<Bytes> bytes,
126+
value_serializer_->SerializeToBytes(*from.value));
127+
PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<InternalRow> value,
128+
value_serializer_->Deserialize(bytes));
129+
return KeyValue(value_kind, from.sequence_number, KeyValue::UNKNOWN_LEVEL, from.key,
130+
std::move(value));
131+
}
132+
133+
std::unique_ptr<MergeFunction> merge_function_;
134+
int32_t max_level_;
135+
std::unique_ptr<RowCompactedSerializer> value_serializer_;
136+
FieldsComparator::FieldComparatorFunc value_equalizer_;
137+
std::optional<KeyValue> top_level_kv_;
138+
std::optional<KeyValue> initial_kv_;
139+
bool is_initialized_ = false;
140+
};
141+
142+
} // namespace paimon
Lines changed: 228 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,228 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
#include "paimon/core/mergetree/compact/full_changelog_merge_function_wrapper.h"
21+
22+
#include <memory>
23+
#include <utility>
24+
25+
#include "gtest/gtest.h"
26+
#include "paimon/core/mergetree/compact/deduplicate_merge_function.h"
27+
#include "paimon/core/mergetree/compact/internal_row_equalizer.h"
28+
#include "paimon/memory/memory_pool.h"
29+
#include "paimon/testing/utils/binary_row_generator.h"
30+
#include "paimon/testing/utils/testharness.h"
31+
32+
namespace paimon::test {
33+
namespace {
34+
35+
constexpr int32_t MAX_LEVEL = 3;
36+
37+
KeyValue MakeKeyValue(const RowKind* kind, int64_t sequence_number, int32_t level, int32_t key,
38+
int32_t value, const std::shared_ptr<MemoryPool>& pool) {
39+
return KeyValue(kind, sequence_number, level,
40+
BinaryRowGenerator::GenerateRowPtr({key}, pool.get()),
41+
BinaryRowGenerator::GenerateRowPtr({value}, pool.get()));
42+
}
43+
44+
std::unique_ptr<RowCompactedSerializer> CreateValueSerializer(
45+
const std::shared_ptr<MemoryPool>& pool) {
46+
return RowCompactedSerializer::Create(arrow::schema({arrow::field("value", arrow::int32())}),
47+
pool)
48+
.value();
49+
}
50+
51+
std::unique_ptr<FullChangelogMergeFunctionWrapper> CreateWrapper(
52+
const std::shared_ptr<MemoryPool>& pool,
53+
FieldsComparator::FieldComparatorFunc value_equalizer = {}) {
54+
return std::make_unique<FullChangelogMergeFunctionWrapper>(
55+
std::make_unique<DeduplicateMergeFunction>(/*ignore_delete=*/false), MAX_LEVEL,
56+
CreateValueSerializer(pool), std::move(value_equalizer));
57+
}
58+
59+
void CheckKeyValue(const KeyValue& actual, const RowKind* kind, int64_t sequence_number,
60+
int32_t level, int32_t value) {
61+
ASSERT_EQ(kind, actual.value_kind);
62+
ASSERT_EQ(sequence_number, actual.sequence_number);
63+
ASSERT_EQ(level, actual.level);
64+
ASSERT_EQ(value, actual.value->GetInt(0));
65+
}
66+
67+
} // namespace
68+
69+
TEST(FullChangelogMergeFunctionWrapperTest, TestSingleRecord) {
70+
auto pool = GetDefaultPool();
71+
auto wrapper = CreateWrapper(pool);
72+
73+
wrapper->Reset();
74+
ASSERT_OK(wrapper->Add(
75+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, /*level=*/0, 1, 10, pool)));
76+
ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> insert_result, wrapper->GetResult());
77+
ASSERT_TRUE(insert_result);
78+
ASSERT_TRUE(insert_result->result);
79+
ASSERT_EQ(1, insert_result->changelogs.size());
80+
CheckKeyValue(insert_result->changelogs[0], RowKind::Insert(), 1, KeyValue::UNKNOWN_LEVEL, 10);
81+
CheckKeyValue(*insert_result->result, RowKind::Insert(), 1, 0, 10);
82+
83+
wrapper->Reset();
84+
ASSERT_OK(wrapper->Add(
85+
MakeKeyValue(RowKind::Delete(), /*sequence_number=*/2, /*level=*/0, 2, 20, pool)));
86+
ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> delete_result, wrapper->GetResult());
87+
ASSERT_TRUE(delete_result);
88+
ASSERT_FALSE(delete_result->result);
89+
ASSERT_TRUE(delete_result->changelogs.empty());
90+
91+
wrapper->Reset();
92+
ASSERT_OK(wrapper->Add(
93+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/3, MAX_LEVEL, 3, 30, pool)));
94+
ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> top_level_result, wrapper->GetResult());
95+
ASSERT_TRUE(top_level_result);
96+
ASSERT_TRUE(top_level_result->result);
97+
ASSERT_TRUE(top_level_result->changelogs.empty());
98+
CheckKeyValue(*top_level_result->result, RowKind::Insert(), 3, MAX_LEVEL, 30);
99+
}
100+
101+
TEST(FullChangelogMergeFunctionWrapperTest, TestInsertUpdateAndDelete) {
102+
auto pool = GetDefaultPool();
103+
auto wrapper = CreateWrapper(pool);
104+
105+
wrapper->Reset();
106+
ASSERT_OK(wrapper->Add(
107+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool)));
108+
ASSERT_OK(wrapper->Add(
109+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, 1, 20, pool)));
110+
ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> update_result, wrapper->GetResult());
111+
ASSERT_TRUE(update_result);
112+
ASSERT_TRUE(update_result->result);
113+
ASSERT_EQ(2, update_result->changelogs.size());
114+
CheckKeyValue(update_result->changelogs[0], RowKind::UpdateBefore(), 1, KeyValue::UNKNOWN_LEVEL,
115+
10);
116+
CheckKeyValue(update_result->changelogs[1], RowKind::UpdateAfter(), 2, KeyValue::UNKNOWN_LEVEL,
117+
20);
118+
CheckKeyValue(*update_result->result, RowKind::Insert(), 2, 0, 20);
119+
120+
wrapper->Reset();
121+
ASSERT_OK(wrapper->Add(
122+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/3, MAX_LEVEL, 2, 30, pool)));
123+
ASSERT_OK(wrapper->Add(
124+
MakeKeyValue(RowKind::Delete(), /*sequence_number=*/4, /*level=*/0, 2, 30, pool)));
125+
ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> delete_result, wrapper->GetResult());
126+
ASSERT_TRUE(delete_result);
127+
ASSERT_FALSE(delete_result->result);
128+
ASSERT_EQ(1, delete_result->changelogs.size());
129+
CheckKeyValue(delete_result->changelogs[0], RowKind::Delete(), 3, KeyValue::UNKNOWN_LEVEL, 30);
130+
131+
wrapper->Reset();
132+
ASSERT_OK(wrapper->Add(
133+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/5, /*level=*/0, 3, 40, pool)));
134+
ASSERT_OK(wrapper->Add(
135+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/6, /*level=*/0, 3, 50, pool)));
136+
ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> insert_result, wrapper->GetResult());
137+
ASSERT_TRUE(insert_result);
138+
ASSERT_TRUE(insert_result->result);
139+
ASSERT_EQ(1, insert_result->changelogs.size());
140+
CheckKeyValue(insert_result->changelogs[0], RowKind::Insert(), 6, KeyValue::UNKNOWN_LEVEL, 50);
141+
CheckKeyValue(*insert_result->result, RowKind::Insert(), 6, 0, 50);
142+
}
143+
144+
TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicate) {
145+
auto pool = GetDefaultPool();
146+
147+
auto wrapper_without_deduplicate = CreateWrapper(pool);
148+
wrapper_without_deduplicate->Reset();
149+
ASSERT_OK(wrapper_without_deduplicate->Add(
150+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool)));
151+
ASSERT_OK(wrapper_without_deduplicate->Add(
152+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, 1, 10, pool)));
153+
ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> result_without_deduplicate,
154+
wrapper_without_deduplicate->GetResult());
155+
ASSERT_TRUE(result_without_deduplicate);
156+
ASSERT_EQ(2, result_without_deduplicate->changelogs.size());
157+
158+
auto value_schema = arrow::schema({arrow::field("value", arrow::int32())});
159+
ASSERT_OK_AND_ASSIGN(FieldsComparator::FieldComparatorFunc value_equalizer,
160+
InternalRowEqualizer::Create(value_schema, /*ignore_fields=*/{}));
161+
auto wrapper = CreateWrapper(pool, std::move(value_equalizer));
162+
163+
wrapper->Reset();
164+
ASSERT_OK(wrapper->Add(
165+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool)));
166+
ASSERT_OK(wrapper->Add(
167+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, 1, 10, pool)));
168+
ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> result, wrapper->GetResult());
169+
ASSERT_TRUE(result);
170+
ASSERT_TRUE(result->result);
171+
ASSERT_TRUE(result->changelogs.empty());
172+
CheckKeyValue(*result->result, RowKind::Insert(), 2, 0, 10);
173+
}
174+
175+
TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicateWithIgnoreFields) {
176+
auto pool = GetDefaultPool();
177+
auto value_schema = arrow::schema(
178+
{arrow::field("value", arrow::int32()), arrow::field("ignored", arrow::int32())});
179+
ASSERT_OK_AND_ASSIGN(FieldsComparator::FieldComparatorFunc value_equalizer,
180+
InternalRowEqualizer::Create(value_schema, {"ignored"}));
181+
ASSERT_OK_AND_ASSIGN(std::unique_ptr<RowCompactedSerializer> value_serializer,
182+
RowCompactedSerializer::Create(value_schema, pool));
183+
FullChangelogMergeFunctionWrapper wrapper(
184+
std::make_unique<DeduplicateMergeFunction>(/*ignore_delete=*/false), MAX_LEVEL,
185+
std::move(value_serializer), std::move(value_equalizer));
186+
187+
wrapper.Reset();
188+
ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL,
189+
BinaryRowGenerator::GenerateRowPtr({1}, pool.get()),
190+
BinaryRowGenerator::GenerateRowPtr({10, 1}, pool.get()))));
191+
ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0,
192+
BinaryRowGenerator::GenerateRowPtr({1}, pool.get()),
193+
BinaryRowGenerator::GenerateRowPtr({10, 2}, pool.get()))));
194+
ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> ignored_field_result, wrapper.GetResult());
195+
ASSERT_TRUE(ignored_field_result);
196+
ASSERT_TRUE(ignored_field_result->result);
197+
ASSERT_TRUE(ignored_field_result->changelogs.empty());
198+
ASSERT_EQ(10, ignored_field_result->result->value->GetInt(0));
199+
ASSERT_EQ(2, ignored_field_result->result->value->GetInt(1));
200+
201+
wrapper.Reset();
202+
ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/3, MAX_LEVEL,
203+
BinaryRowGenerator::GenerateRowPtr({1}, pool.get()),
204+
BinaryRowGenerator::GenerateRowPtr({10, 1}, pool.get()))));
205+
ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/4, /*level=*/0,
206+
BinaryRowGenerator::GenerateRowPtr({1}, pool.get()),
207+
BinaryRowGenerator::GenerateRowPtr({11, 2}, pool.get()))));
208+
ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> value_field_result, wrapper.GetResult());
209+
ASSERT_TRUE(value_field_result);
210+
ASSERT_TRUE(value_field_result->result);
211+
ASSERT_EQ(2, value_field_result->changelogs.size());
212+
ASSERT_EQ(RowKind::UpdateBefore(), value_field_result->changelogs[0].value_kind);
213+
ASSERT_EQ(RowKind::UpdateAfter(), value_field_result->changelogs[1].value_kind);
214+
}
215+
216+
TEST(FullChangelogMergeFunctionWrapperTest, TestRejectMultipleTopLevelRecords) {
217+
auto pool = GetDefaultPool();
218+
auto wrapper = CreateWrapper(pool);
219+
220+
wrapper->Reset();
221+
ASSERT_OK(wrapper->Add(
222+
MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool)));
223+
ASSERT_NOK_WITH_MSG(wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2,
224+
MAX_LEVEL, 1, 20, pool)),
225+
"Top level key-value already exists");
226+
}
227+
228+
} // namespace paimon::test

0 commit comments

Comments
 (0)