3232namespace paimon ::test {
3333namespace {
3434
35- constexpr int32_t MAX_LEVEL = 3 ;
35+ constexpr int32_t kMaxLevel = 3 ;
3636
3737KeyValue MakeKeyValue (const RowKind* kind, int64_t sequence_number, int32_t level, int32_t key,
3838 int32_t value, const std::shared_ptr<MemoryPool>& pool) {
@@ -52,7 +52,7 @@ std::unique_ptr<FullChangelogMergeFunctionWrapper> CreateWrapper(
5252 const std::shared_ptr<MemoryPool>& pool,
5353 FieldsComparator::FieldComparatorFunc value_equalizer = {}) {
5454 return std::make_unique<FullChangelogMergeFunctionWrapper>(
55- std::make_unique<DeduplicateMergeFunction>(/* ignore_delete=*/ false ), MAX_LEVEL ,
55+ std::make_unique<DeduplicateMergeFunction>(/* ignore_delete=*/ false ), kMaxLevel ,
5656 CreateValueSerializer (pool), std::move (value_equalizer));
5757}
5858
@@ -71,74 +71,90 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestSingleRecord) {
7171 auto wrapper = CreateWrapper (pool);
7272
7373 wrapper->Reset ();
74- ASSERT_OK (wrapper->Add (
75- MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 1 , /* level=*/ 0 , 1 , 10 , pool)));
74+ ASSERT_OK (
75+ wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 1 , /* level=*/ 0 , /* key=*/ 1 ,
76+ /* value=*/ 10 , pool)));
7677 ASSERT_OK_AND_ASSIGN (std::optional<ChangelogResult> insert_result, wrapper->GetResult ());
7778 ASSERT_TRUE (insert_result);
7879 ASSERT_TRUE (insert_result->result );
7980 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 );
81+ CheckKeyValue (insert_result->changelogs [0 ], RowKind::Insert (), /* sequence_number=*/ 1 ,
82+ /* level=*/ KeyValue::UNKNOWN_LEVEL , /* value=*/ 10 );
83+ CheckKeyValue (*insert_result->result , RowKind::Insert (), /* sequence_number=*/ 1 , /* level=*/ 0 ,
84+ /* value=*/ 10 );
8285
8386 wrapper->Reset ();
84- ASSERT_OK (wrapper->Add (
85- MakeKeyValue (RowKind::Delete (), /* sequence_number=*/ 2 , /* level=*/ 0 , 2 , 20 , pool)));
87+ ASSERT_OK (
88+ wrapper->Add (MakeKeyValue (RowKind::Delete (), /* sequence_number=*/ 2 , /* level=*/ 0 , /* key=*/ 2 ,
89+ /* value=*/ 20 , pool)));
8690 ASSERT_OK_AND_ASSIGN (std::optional<ChangelogResult> delete_result, wrapper->GetResult ());
8791 ASSERT_TRUE (delete_result);
8892 ASSERT_FALSE (delete_result->result );
8993 ASSERT_TRUE (delete_result->changelogs .empty ());
9094
9195 wrapper->Reset ();
92- ASSERT_OK (wrapper->Add (
93- MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 3 , MAX_LEVEL , 3 , 30 , pool)));
96+ ASSERT_OK (wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 3 ,
97+ /* level=*/ kMaxLevel , /* key=*/ 3 ,
98+ /* value=*/ 30 , pool)));
9499 ASSERT_OK_AND_ASSIGN (std::optional<ChangelogResult> top_level_result, wrapper->GetResult ());
95100 ASSERT_TRUE (top_level_result);
96101 ASSERT_TRUE (top_level_result->result );
97102 ASSERT_TRUE (top_level_result->changelogs .empty ());
98- CheckKeyValue (*top_level_result->result , RowKind::Insert (), 3 , MAX_LEVEL , 30 );
103+ CheckKeyValue (*top_level_result->result , RowKind::Insert (), /* sequence_number=*/ 3 ,
104+ /* level=*/ kMaxLevel , /* value=*/ 30 );
99105}
100106
101107TEST (FullChangelogMergeFunctionWrapperTest, TestInsertUpdateAndDelete) {
102108 auto pool = GetDefaultPool ();
103109 auto wrapper = CreateWrapper (pool);
104110
105111 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)));
112+ ASSERT_OK (wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 1 ,
113+ /* level=*/ kMaxLevel , /* key=*/ 1 ,
114+ /* value=*/ 10 , pool)));
115+ ASSERT_OK (
116+ wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 2 , /* level=*/ 0 , /* key=*/ 1 ,
117+ /* value=*/ 20 , pool)));
110118 ASSERT_OK_AND_ASSIGN (std::optional<ChangelogResult> update_result, wrapper->GetResult ());
111119 ASSERT_TRUE (update_result);
112120 ASSERT_TRUE (update_result->result );
113121 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 );
122+ CheckKeyValue (update_result->changelogs [0 ], RowKind::UpdateBefore (), /* sequence_number=*/ 1 ,
123+ /* level=*/ KeyValue::UNKNOWN_LEVEL , /* value=*/ 10 );
124+ CheckKeyValue (update_result->changelogs [1 ], RowKind::UpdateAfter (), /* sequence_number=*/ 2 ,
125+ /* level=*/ KeyValue::UNKNOWN_LEVEL , /* value=*/ 20 );
126+ CheckKeyValue (*update_result->result , RowKind::Insert (), /* sequence_number=*/ 2 , /* level=*/ 0 ,
127+ /* value=*/ 20 );
119128
120129 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)));
130+ ASSERT_OK (wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 3 ,
131+ /* level=*/ kMaxLevel , /* key=*/ 2 ,
132+ /* value=*/ 30 , pool)));
133+ ASSERT_OK (
134+ wrapper->Add (MakeKeyValue (RowKind::Delete (), /* sequence_number=*/ 4 , /* level=*/ 0 , /* key=*/ 2 ,
135+ /* value=*/ 30 , pool)));
125136 ASSERT_OK_AND_ASSIGN (std::optional<ChangelogResult> delete_result, wrapper->GetResult ());
126137 ASSERT_TRUE (delete_result);
127138 ASSERT_FALSE (delete_result->result );
128139 ASSERT_EQ (1 , delete_result->changelogs .size ());
129- CheckKeyValue (delete_result->changelogs [0 ], RowKind::Delete (), 3 , KeyValue::UNKNOWN_LEVEL , 30 );
140+ CheckKeyValue (delete_result->changelogs [0 ], RowKind::Delete (), /* sequence_number=*/ 3 ,
141+ /* level=*/ KeyValue::UNKNOWN_LEVEL , /* value=*/ 30 );
130142
131143 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)));
144+ ASSERT_OK (
145+ wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 5 , /* level=*/ 0 , /* key=*/ 3 ,
146+ /* value=*/ 40 , pool)));
147+ ASSERT_OK (
148+ wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 6 , /* level=*/ 0 , /* key=*/ 3 ,
149+ /* value=*/ 50 , pool)));
136150 ASSERT_OK_AND_ASSIGN (std::optional<ChangelogResult> insert_result, wrapper->GetResult ());
137151 ASSERT_TRUE (insert_result);
138152 ASSERT_TRUE (insert_result->result );
139153 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 );
154+ CheckKeyValue (insert_result->changelogs [0 ], RowKind::Insert (), /* sequence_number=*/ 6 ,
155+ /* level=*/ KeyValue::UNKNOWN_LEVEL , /* value=*/ 50 );
156+ CheckKeyValue (*insert_result->result , RowKind::Insert (), /* sequence_number=*/ 6 , /* level=*/ 0 ,
157+ /* value=*/ 50 );
142158}
143159
144160TEST (FullChangelogMergeFunctionWrapperTest, TestRowDeduplicate) {
@@ -147,9 +163,11 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicate) {
147163 auto wrapper_without_deduplicate = CreateWrapper (pool);
148164 wrapper_without_deduplicate->Reset ();
149165 ASSERT_OK (wrapper_without_deduplicate->Add (
150- MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 1 , MAX_LEVEL , 1 , 10 , pool)));
166+ MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 1 , /* level=*/ kMaxLevel , /* key=*/ 1 ,
167+ /* value=*/ 10 , pool)));
151168 ASSERT_OK (wrapper_without_deduplicate->Add (
152- MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 2 , /* level=*/ 0 , 1 , 10 , pool)));
169+ MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 2 , /* level=*/ 0 , /* key=*/ 1 ,
170+ /* value=*/ 10 , pool)));
153171 ASSERT_OK_AND_ASSIGN (std::optional<ChangelogResult> result_without_deduplicate,
154172 wrapper_without_deduplicate->GetResult ());
155173 ASSERT_TRUE (result_without_deduplicate);
@@ -161,15 +179,18 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicate) {
161179 auto wrapper = CreateWrapper (pool, std::move (value_equalizer));
162180
163181 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)));
182+ ASSERT_OK (wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 1 ,
183+ /* level=*/ kMaxLevel , /* key=*/ 1 ,
184+ /* value=*/ 10 , pool)));
185+ ASSERT_OK (
186+ wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 2 , /* level=*/ 0 , /* key=*/ 1 ,
187+ /* value=*/ 10 , pool)));
168188 ASSERT_OK_AND_ASSIGN (std::optional<ChangelogResult> result, wrapper->GetResult ());
169189 ASSERT_TRUE (result);
170190 ASSERT_TRUE (result->result );
171191 ASSERT_TRUE (result->changelogs .empty ());
172- CheckKeyValue (*result->result , RowKind::Insert (), 2 , 0 , 10 );
192+ CheckKeyValue (*result->result , RowKind::Insert (), /* sequence_number=*/ 2 , /* level=*/ 0 ,
193+ /* value=*/ 10 );
173194}
174195
175196TEST (FullChangelogMergeFunctionWrapperTest, TestRowDeduplicateWithIgnoreFields) {
@@ -181,11 +202,11 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicateWithIgnoreFields)
181202 ASSERT_OK_AND_ASSIGN (std::unique_ptr<RowCompactedSerializer> value_serializer,
182203 RowCompactedSerializer::Create (value_schema, pool));
183204 FullChangelogMergeFunctionWrapper wrapper (
184- std::make_unique<DeduplicateMergeFunction>(/* ignore_delete=*/ false ), MAX_LEVEL ,
205+ std::make_unique<DeduplicateMergeFunction>(/* ignore_delete=*/ false ), kMaxLevel ,
185206 std::move (value_serializer), std::move (value_equalizer));
186207
187208 wrapper.Reset ();
188- ASSERT_OK (wrapper.Add (KeyValue (RowKind::Insert (), /* sequence_number=*/ 1 , MAX_LEVEL ,
209+ ASSERT_OK (wrapper.Add (KeyValue (RowKind::Insert (), /* sequence_number=*/ 1 , kMaxLevel ,
189210 BinaryRowGenerator::GenerateRowPtr ({1 }, pool.get ()),
190211 BinaryRowGenerator::GenerateRowPtr ({10 , 1 }, pool.get ()))));
191212 ASSERT_OK (wrapper.Add (KeyValue (RowKind::Insert (), /* sequence_number=*/ 2 , /* level=*/ 0 ,
@@ -199,7 +220,7 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicateWithIgnoreFields)
199220 ASSERT_EQ (2 , ignored_field_result->result ->value ->GetInt (1 ));
200221
201222 wrapper.Reset ();
202- ASSERT_OK (wrapper.Add (KeyValue (RowKind::Insert (), /* sequence_number=*/ 3 , MAX_LEVEL ,
223+ ASSERT_OK (wrapper.Add (KeyValue (RowKind::Insert (), /* sequence_number=*/ 3 , kMaxLevel ,
203224 BinaryRowGenerator::GenerateRowPtr ({1 }, pool.get ()),
204225 BinaryRowGenerator::GenerateRowPtr ({10 , 1 }, pool.get ()))));
205226 ASSERT_OK (wrapper.Add (KeyValue (RowKind::Insert (), /* sequence_number=*/ 4 , /* level=*/ 0 ,
@@ -218,11 +239,13 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestRejectMultipleTopLevelRecords) {
218239 auto wrapper = CreateWrapper (pool);
219240
220241 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" );
242+ ASSERT_OK (wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 1 ,
243+ /* level=*/ kMaxLevel , /* key=*/ 1 ,
244+ /* value=*/ 10 , pool)));
245+ ASSERT_NOK_WITH_MSG (
246+ wrapper->Add (MakeKeyValue (RowKind::Insert (), /* sequence_number=*/ 2 ,
247+ /* level=*/ kMaxLevel , /* key=*/ 1 , /* value=*/ 20 , pool)),
248+ " Top level key-value already exists" );
226249}
227250
228251} // namespace paimon::test
0 commit comments