From bb7933f9981623ea2b5021320dc38f69afeb0fbf Mon Sep 17 00:00:00 2001 From: "lisizhuo.lsz" Date: Mon, 31 Aug 2026 17:15:26 +0800 Subject: [PATCH 1/5] feat(changelog): support full-compaction mode changelog producer --- include/paimon/defs.h | 1 - src/paimon/CMakeLists.txt | 2 + .../full_changelog_merge_function_wrapper.h | 142 ++++++++++ ..._changelog_merge_function_wrapper_test.cpp | 228 ++++++++++++++++ ..._changelog_merge_tree_compact_rewriter.cpp | 131 ++++++++++ ...ll_changelog_merge_tree_compact_rewriter.h | 67 +++++ .../lookup_merge_tree_compact_rewriter.cpp | 4 +- .../merge_tree_compact_manager_factory.cpp | 10 +- ...erge_tree_compact_manager_factory_test.cpp | 13 +- .../compact/merge_tree_compact_rewriter.cpp | 6 +- .../compact/merge_tree_compact_rewriter.h | 5 +- src/paimon/core/schema/schema_validation.cpp | 5 - .../core/schema/schema_validation_test.cpp | 14 +- .../table/source/data_table_stream_scan.cpp | 9 +- test/inte/write_and_read_inte_test.cpp | 243 ++++++++++++++++++ 15 files changed, 844 insertions(+), 36 deletions(-) create mode 100644 src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper.h create mode 100644 src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp create mode 100644 src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.cpp create mode 100644 src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.h diff --git a/include/paimon/defs.h b/include/paimon/defs.h index e05faefd8..0069ded74 100644 --- a/include/paimon/defs.h +++ b/include/paimon/defs.h @@ -390,7 +390,6 @@ struct PAIMON_EXPORT Options { /// keeps the details of data changes, it can be read directly during stream reads. This can be /// applied to tables with primary keys. Values can be "none", "input", "lookup", /// "full-compaction". Default value is "none". - /// @note C++ Paimon currently supports "none", "input", and "lookup". static const char CHANGELOG_PRODUCER[]; /// "changelog-producer.row-deduplicate" - Whether to generate update-before and update-after diff --git a/src/paimon/CMakeLists.txt b/src/paimon/CMakeLists.txt index 051eba324..54b5aafb4 100644 --- a/src/paimon/CMakeLists.txt +++ b/src/paimon/CMakeLists.txt @@ -329,6 +329,7 @@ set(PAIMON_CORE_SRCS core/mergetree/compact/merge_tree_compact_manager_factory.cpp core/mergetree/compact/merge_tree_compact_rewriter.cpp core/mergetree/compact/merge_tree_compact_task.cpp + core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.cpp core/mergetree/compact/partial_update_merge_function.cpp core/mergetree/compact/sort_merge_reader_with_loser_tree.cpp core/mergetree/compact/sort_merge_reader_with_min_heap.cpp @@ -820,6 +821,7 @@ if(PAIMON_BUILD_TESTS) core/mergetree/compact/deduplicate_merge_function_test.cpp core/mergetree/compact/first_row_merge_function_test.cpp core/mergetree/compact/first_row_merge_function_wrapper_test.cpp + core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp core/mergetree/compact/internal_row_equalizer_test.cpp core/mergetree/compact/interval_partition_test.cpp core/mergetree/compact/lookup_changelog_merge_function_wrapper_test.cpp diff --git a/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper.h b/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper.h new file mode 100644 index 000000000..f5698b582 --- /dev/null +++ b/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper.h @@ -0,0 +1,142 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#pragma once + +#include +#include +#include + +#include "paimon/common/data/serializer/row_compacted_serializer.h" +#include "paimon/common/utils/fields_comparator.h" +#include "paimon/core/key_value.h" +#include "paimon/core/mergetree/compact/changelog_result.h" +#include "paimon/core/mergetree/compact/merge_function.h" +#include "paimon/core/mergetree/compact/merge_function_wrapper.h" +#include "paimon/result.h" +#include "paimon/status.h" + +namespace paimon { + +/// Wrapper for `MergeFunction`s which produces changelog during a full compaction. +class FullChangelogMergeFunctionWrapper : public MergeFunctionWrapper { + public: + FullChangelogMergeFunctionWrapper(std::unique_ptr&& merge_function, + int32_t max_level, + std::unique_ptr&& value_serializer, + FieldsComparator::FieldComparatorFunc value_equalizer) + : merge_function_(std::move(merge_function)), + max_level_(max_level), + value_serializer_(std::move(value_serializer)), + value_equalizer_(std::move(value_equalizer)) {} + + void Reset() override { + merge_function_->Reset(); + top_level_kv_ = std::nullopt; + initial_kv_ = std::nullopt; + is_initialized_ = false; + } + + Status Add(KeyValue&& kv) override { + if (!initial_kv_) { + initial_kv_ = std::move(kv); + return Status::OK(); + } + + if (!is_initialized_) { + if (initial_kv_->level == max_level_) { + PAIMON_RETURN_NOT_OK(RememberTopLevel(*initial_kv_)); + } + PAIMON_RETURN_NOT_OK(merge_function_->Add(std::move(initial_kv_).value())); + is_initialized_ = true; + } + + if (kv.level == max_level_) { + PAIMON_RETURN_NOT_OK(RememberTopLevel(kv)); + } + return merge_function_->Add(std::move(kv)); + } + + Result> GetResult() override { + std::optional merged; + if (is_initialized_) { + PAIMON_ASSIGN_OR_RAISE(merged, merge_function_->GetResult()); + } else { + merged = std::move(initial_kv_); + } + + ChangelogResult result; + if (is_initialized_) { + if (!top_level_kv_) { + if (merged && merged->value_kind->IsAdd()) { + PAIMON_ASSIGN_OR_RAISE(KeyValue insert, + CloneKeyValue(*merged, RowKind::Insert())); + result.changelogs.emplace_back(std::move(insert)); + } + } else if (!merged || !merged->value_kind->IsAdd()) { + top_level_kv_->value_kind = RowKind::Delete(); + result.changelogs.emplace_back(std::move(top_level_kv_).value()); + } else if (!value_equalizer_ || + value_equalizer_(*top_level_kv_->value, *merged->value) != 0) { + top_level_kv_->value_kind = RowKind::UpdateBefore(); + result.changelogs.emplace_back(std::move(top_level_kv_).value()); + PAIMON_ASSIGN_OR_RAISE(KeyValue update_after, + CloneKeyValue(*merged, RowKind::UpdateAfter())); + result.changelogs.emplace_back(std::move(update_after)); + } + } else if (merged && merged->level != max_level_ && merged->value_kind->IsAdd()) { + PAIMON_ASSIGN_OR_RAISE(KeyValue insert, CloneKeyValue(*merged, RowKind::Insert())); + result.changelogs.emplace_back(std::move(insert)); + } + + if (merged && merged->value_kind->IsAdd()) { + result.result = std::move(merged); + } + Reset(); + return std::optional(std::move(result)); + } + + private: + Status RememberTopLevel(const KeyValue& kv) { + if (top_level_kv_) { + return Status::Invalid("Top level key-value already exists. This is unexpected."); + } + PAIMON_ASSIGN_OR_RAISE(top_level_kv_, CloneKeyValue(kv, kv.value_kind)); + return Status::OK(); + } + + Result CloneKeyValue(const KeyValue& from, const RowKind* value_kind) const { + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr bytes, + value_serializer_->SerializeToBytes(*from.value)); + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr value, + value_serializer_->Deserialize(bytes)); + return KeyValue(value_kind, from.sequence_number, KeyValue::UNKNOWN_LEVEL, from.key, + std::move(value)); + } + + std::unique_ptr merge_function_; + int32_t max_level_; + std::unique_ptr value_serializer_; + FieldsComparator::FieldComparatorFunc value_equalizer_; + std::optional top_level_kv_; + std::optional initial_kv_; + bool is_initialized_ = false; +}; + +} // namespace paimon diff --git a/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp b/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp new file mode 100644 index 000000000..dd61bee85 --- /dev/null +++ b/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp @@ -0,0 +1,228 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "paimon/core/mergetree/compact/full_changelog_merge_function_wrapper.h" + +#include +#include + +#include "gtest/gtest.h" +#include "paimon/core/mergetree/compact/deduplicate_merge_function.h" +#include "paimon/core/mergetree/compact/internal_row_equalizer.h" +#include "paimon/memory/memory_pool.h" +#include "paimon/testing/utils/binary_row_generator.h" +#include "paimon/testing/utils/testharness.h" + +namespace paimon::test { +namespace { + +constexpr int32_t MAX_LEVEL = 3; + +KeyValue MakeKeyValue(const RowKind* kind, int64_t sequence_number, int32_t level, int32_t key, + int32_t value, const std::shared_ptr& pool) { + return KeyValue(kind, sequence_number, level, + BinaryRowGenerator::GenerateRowPtr({key}, pool.get()), + BinaryRowGenerator::GenerateRowPtr({value}, pool.get())); +} + +std::unique_ptr CreateValueSerializer( + const std::shared_ptr& pool) { + return RowCompactedSerializer::Create(arrow::schema({arrow::field("value", arrow::int32())}), + pool) + .value(); +} + +std::unique_ptr CreateWrapper( + const std::shared_ptr& pool, + FieldsComparator::FieldComparatorFunc value_equalizer = {}) { + return std::make_unique( + std::make_unique(/*ignore_delete=*/false), MAX_LEVEL, + CreateValueSerializer(pool), std::move(value_equalizer)); +} + +void CheckKeyValue(const KeyValue& actual, const RowKind* kind, int64_t sequence_number, + int32_t level, int32_t value) { + ASSERT_EQ(kind, actual.value_kind); + ASSERT_EQ(sequence_number, actual.sequence_number); + ASSERT_EQ(level, actual.level); + ASSERT_EQ(value, actual.value->GetInt(0)); +} + +} // namespace + +TEST(FullChangelogMergeFunctionWrapperTest, TestSingleRecord) { + auto pool = GetDefaultPool(); + auto wrapper = CreateWrapper(pool); + + wrapper->Reset(); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, /*level=*/0, 1, 10, pool))); + ASSERT_OK_AND_ASSIGN(std::optional insert_result, wrapper->GetResult()); + ASSERT_TRUE(insert_result); + ASSERT_TRUE(insert_result->result); + ASSERT_EQ(1, insert_result->changelogs.size()); + CheckKeyValue(insert_result->changelogs[0], RowKind::Insert(), 1, KeyValue::UNKNOWN_LEVEL, 10); + CheckKeyValue(*insert_result->result, RowKind::Insert(), 1, 0, 10); + + wrapper->Reset(); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Delete(), /*sequence_number=*/2, /*level=*/0, 2, 20, pool))); + ASSERT_OK_AND_ASSIGN(std::optional delete_result, wrapper->GetResult()); + ASSERT_TRUE(delete_result); + ASSERT_FALSE(delete_result->result); + ASSERT_TRUE(delete_result->changelogs.empty()); + + wrapper->Reset(); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/3, MAX_LEVEL, 3, 30, pool))); + ASSERT_OK_AND_ASSIGN(std::optional top_level_result, wrapper->GetResult()); + ASSERT_TRUE(top_level_result); + ASSERT_TRUE(top_level_result->result); + ASSERT_TRUE(top_level_result->changelogs.empty()); + CheckKeyValue(*top_level_result->result, RowKind::Insert(), 3, MAX_LEVEL, 30); +} + +TEST(FullChangelogMergeFunctionWrapperTest, TestInsertUpdateAndDelete) { + auto pool = GetDefaultPool(); + auto wrapper = CreateWrapper(pool); + + wrapper->Reset(); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool))); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, 1, 20, pool))); + ASSERT_OK_AND_ASSIGN(std::optional update_result, wrapper->GetResult()); + ASSERT_TRUE(update_result); + ASSERT_TRUE(update_result->result); + ASSERT_EQ(2, update_result->changelogs.size()); + CheckKeyValue(update_result->changelogs[0], RowKind::UpdateBefore(), 1, KeyValue::UNKNOWN_LEVEL, + 10); + CheckKeyValue(update_result->changelogs[1], RowKind::UpdateAfter(), 2, KeyValue::UNKNOWN_LEVEL, + 20); + CheckKeyValue(*update_result->result, RowKind::Insert(), 2, 0, 20); + + wrapper->Reset(); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/3, MAX_LEVEL, 2, 30, pool))); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Delete(), /*sequence_number=*/4, /*level=*/0, 2, 30, pool))); + ASSERT_OK_AND_ASSIGN(std::optional delete_result, wrapper->GetResult()); + ASSERT_TRUE(delete_result); + ASSERT_FALSE(delete_result->result); + ASSERT_EQ(1, delete_result->changelogs.size()); + CheckKeyValue(delete_result->changelogs[0], RowKind::Delete(), 3, KeyValue::UNKNOWN_LEVEL, 30); + + wrapper->Reset(); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/5, /*level=*/0, 3, 40, pool))); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/6, /*level=*/0, 3, 50, pool))); + ASSERT_OK_AND_ASSIGN(std::optional insert_result, wrapper->GetResult()); + ASSERT_TRUE(insert_result); + ASSERT_TRUE(insert_result->result); + ASSERT_EQ(1, insert_result->changelogs.size()); + CheckKeyValue(insert_result->changelogs[0], RowKind::Insert(), 6, KeyValue::UNKNOWN_LEVEL, 50); + CheckKeyValue(*insert_result->result, RowKind::Insert(), 6, 0, 50); +} + +TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicate) { + auto pool = GetDefaultPool(); + + auto wrapper_without_deduplicate = CreateWrapper(pool); + wrapper_without_deduplicate->Reset(); + ASSERT_OK(wrapper_without_deduplicate->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool))); + ASSERT_OK(wrapper_without_deduplicate->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, 1, 10, pool))); + ASSERT_OK_AND_ASSIGN(std::optional result_without_deduplicate, + wrapper_without_deduplicate->GetResult()); + ASSERT_TRUE(result_without_deduplicate); + ASSERT_EQ(2, result_without_deduplicate->changelogs.size()); + + auto value_schema = arrow::schema({arrow::field("value", arrow::int32())}); + ASSERT_OK_AND_ASSIGN(FieldsComparator::FieldComparatorFunc value_equalizer, + InternalRowEqualizer::Create(value_schema, /*ignore_fields=*/{})); + auto wrapper = CreateWrapper(pool, std::move(value_equalizer)); + + wrapper->Reset(); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool))); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, 1, 10, pool))); + ASSERT_OK_AND_ASSIGN(std::optional result, wrapper->GetResult()); + ASSERT_TRUE(result); + ASSERT_TRUE(result->result); + ASSERT_TRUE(result->changelogs.empty()); + CheckKeyValue(*result->result, RowKind::Insert(), 2, 0, 10); +} + +TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicateWithIgnoreFields) { + auto pool = GetDefaultPool(); + auto value_schema = arrow::schema( + {arrow::field("value", arrow::int32()), arrow::field("ignored", arrow::int32())}); + ASSERT_OK_AND_ASSIGN(FieldsComparator::FieldComparatorFunc value_equalizer, + InternalRowEqualizer::Create(value_schema, {"ignored"})); + ASSERT_OK_AND_ASSIGN(std::unique_ptr value_serializer, + RowCompactedSerializer::Create(value_schema, pool)); + FullChangelogMergeFunctionWrapper wrapper( + std::make_unique(/*ignore_delete=*/false), MAX_LEVEL, + std::move(value_serializer), std::move(value_equalizer)); + + wrapper.Reset(); + ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, + BinaryRowGenerator::GenerateRowPtr({1}, pool.get()), + BinaryRowGenerator::GenerateRowPtr({10, 1}, pool.get())))); + ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, + BinaryRowGenerator::GenerateRowPtr({1}, pool.get()), + BinaryRowGenerator::GenerateRowPtr({10, 2}, pool.get())))); + ASSERT_OK_AND_ASSIGN(std::optional ignored_field_result, wrapper.GetResult()); + ASSERT_TRUE(ignored_field_result); + ASSERT_TRUE(ignored_field_result->result); + ASSERT_TRUE(ignored_field_result->changelogs.empty()); + ASSERT_EQ(10, ignored_field_result->result->value->GetInt(0)); + ASSERT_EQ(2, ignored_field_result->result->value->GetInt(1)); + + wrapper.Reset(); + ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/3, MAX_LEVEL, + BinaryRowGenerator::GenerateRowPtr({1}, pool.get()), + BinaryRowGenerator::GenerateRowPtr({10, 1}, pool.get())))); + ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/4, /*level=*/0, + BinaryRowGenerator::GenerateRowPtr({1}, pool.get()), + BinaryRowGenerator::GenerateRowPtr({11, 2}, pool.get())))); + ASSERT_OK_AND_ASSIGN(std::optional value_field_result, wrapper.GetResult()); + ASSERT_TRUE(value_field_result); + ASSERT_TRUE(value_field_result->result); + ASSERT_EQ(2, value_field_result->changelogs.size()); + ASSERT_EQ(RowKind::UpdateBefore(), value_field_result->changelogs[0].value_kind); + ASSERT_EQ(RowKind::UpdateAfter(), value_field_result->changelogs[1].value_kind); +} + +TEST(FullChangelogMergeFunctionWrapperTest, TestRejectMultipleTopLevelRecords) { + auto pool = GetDefaultPool(); + auto wrapper = CreateWrapper(pool); + + wrapper->Reset(); + ASSERT_OK(wrapper->Add( + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool))); + ASSERT_NOK_WITH_MSG(wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, + MAX_LEVEL, 1, 20, pool)), + "Top level key-value already exists"); +} + +} // namespace paimon::test diff --git a/src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.cpp b/src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.cpp new file mode 100644 index 000000000..a85c1a0fb --- /dev/null +++ b/src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.cpp @@ -0,0 +1,131 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.h" + +#include + +#include "paimon/common/data/serializer/row_compacted_serializer.h" +#include "paimon/common/table/special_fields.h" +#include "paimon/core/mergetree/compact/full_changelog_merge_function_wrapper.h" +#include "paimon/core/mergetree/compact/internal_row_equalizer.h" +#include "paimon/core/operation/internal_read_context.h" +#include "paimon/core/utils/primary_key_table_utils.h" +#include "paimon/read_context.h" + +namespace paimon { + +FullChangelogMergeTreeCompactRewriter::FullChangelogMergeTreeCompactRewriter( + int32_t max_level, const BinaryRow& partition, int32_t bucket, int64_t schema_id, + const std::vector& trimmed_primary_keys, const CoreOptions& options, + const std::shared_ptr& data_schema, + const std::shared_ptr& write_schema, DeletionVector::Factory dv_factory, + const std::shared_ptr& path_factory_cache, + std::unique_ptr&& merge_file_split_read, + MergeFunctionWrapperFactory merge_function_wrapper_factory, + ChangelogMergeFunctionWrapperFactory changelog_merge_function_wrapper_factory, + const std::shared_ptr& cancellation_controller, + const std::shared_ptr& pool) + : ChangelogMergeTreeRewriter(max_level, /*force_drop_delete=*/false, partition, bucket, + schema_id, trimmed_primary_keys, options, data_schema, + write_schema, std::move(dv_factory), path_factory_cache, + std::move(merge_file_split_read), + std::move(merge_function_wrapper_factory), + std::move(changelog_merge_function_wrapper_factory), + /*produce_changelog=*/true, cancellation_controller, pool) {} + +Result> +FullChangelogMergeTreeCompactRewriter::Create( + int32_t max_level, int32_t bucket, const BinaryRow& partition, + const std::shared_ptr& table_schema, DeletionVector::Factory dv_factory, + const std::shared_ptr& path_factory_cache, + const CoreOptions& options, + const std::shared_ptr& cancellation_controller, + const std::shared_ptr& pool) { + PAIMON_ASSIGN_OR_RAISE(std::vector trimmed_primary_keys, + table_schema->TrimmedPrimaryKeys()); + std::shared_ptr data_schema = + DataField::ConvertDataFieldsToArrowSchema(table_schema->Fields()); + std::shared_ptr write_schema = + SpecialFields::CompleteSequenceAndValueKindField(data_schema); + + ReadContextBuilder read_context_builder(path_factory_cache->RootPath()); + read_context_builder.SetOptions(options.ToMap()) + .WithFileSystem(options.GetFileSystem()) + .EnablePrefetch(true) + .SetPrefetchMaxParallelNum(1) + .SetPrefetchBatchCount(3) + .WithMemoryPool(pool); + PAIMON_ASSIGN_OR_RAISE(std::shared_ptr read_context, + read_context_builder.Finish()); + PAIMON_ASSIGN_OR_RAISE( + std::shared_ptr internal_context, + InternalReadContext::Create(read_context, table_schema, options.ToMap())); + PAIMON_ASSIGN_OR_RAISE( + std::shared_ptr path_factory, + path_factory_cache->GetOrCreatePathFactory(options.GetFileFormat()->Identifier())); + PAIMON_ASSIGN_OR_RAISE( + std::unique_ptr merge_file_split_read, + MergeFileSplitRead::Create(path_factory, internal_context, pool, CreateDefaultExecutor())); + + MergeFunctionWrapperFactory merge_function_wrapper_factory = + []() -> Result>> { + return std::shared_ptr>(); + }; + + FieldsComparator::FieldComparatorFunc value_equalizer; + if (options.ChangelogRowDeduplicate()) { + PAIMON_ASSIGN_OR_RAISE(value_equalizer, + InternalRowEqualizer::Create( + data_schema, options.GetChangelogRowDeduplicateIgnoreFields())); + } + ChangelogMergeFunctionWrapperFactory changelog_merge_function_wrapper_factory = + [data_schema, trimmed_primary_keys, options, max_level, value_equalizer, + pool](int32_t /*output_level*/) + -> Result>> { + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr merge_function, + PrimaryKeyTableUtils::CreateMergeFunction( + data_schema, trimmed_primary_keys, options, pool)); + PAIMON_ASSIGN_OR_RAISE(std::unique_ptr value_serializer, + RowCompactedSerializer::Create(data_schema, pool)); + std::shared_ptr> wrapper = + std::make_shared( + std::move(merge_function), max_level, std::move(value_serializer), value_equalizer); + return wrapper; + }; + + return std::unique_ptr( + new FullChangelogMergeTreeCompactRewriter( + max_level, partition, bucket, table_schema->Id(), trimmed_primary_keys, options, + data_schema, write_schema, std::move(dv_factory), path_factory_cache, + std::move(merge_file_split_read), std::move(merge_function_wrapper_factory), + std::move(changelog_merge_function_wrapper_factory), cancellation_controller, pool)); +} + +Result FullChangelogMergeTreeCompactRewriter::Rewrite( + int32_t output_level, bool drop_delete, const std::vector>& sections) { + if (output_level == max_level_ && !drop_delete) { + return Status::Invalid( + "Delete records should be dropped from result of full compaction. This is " + "unexpected."); + } + return ChangelogMergeTreeRewriter::Rewrite(output_level, drop_delete, sections); +} + +} // namespace paimon diff --git a/src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.h b/src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.h new file mode 100644 index 000000000..70f159485 --- /dev/null +++ b/src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.h @@ -0,0 +1,67 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#pragma once + +#include + +#include "paimon/core/mergetree/compact/changelog_merge_tree_rewriter.h" + +namespace paimon { + +/// A `MergeTreeCompactRewriter` which produces changelog files for each full compaction. +class FullChangelogMergeTreeCompactRewriter : public ChangelogMergeTreeRewriter { + public: + static Result> Create( + int32_t max_level, int32_t bucket, const BinaryRow& partition, + const std::shared_ptr& table_schema, DeletionVector::Factory dv_factory, + const std::shared_ptr& path_factory_cache, + const CoreOptions& options, + const std::shared_ptr& cancellation_controller, + const std::shared_ptr& pool); + + Result Rewrite(int32_t output_level, bool drop_delete, + const std::vector>& sections) override; + + private: + FullChangelogMergeTreeCompactRewriter( + int32_t max_level, const BinaryRow& partition, int32_t bucket, int64_t schema_id, + const std::vector& trimmed_primary_keys, const CoreOptions& options, + const std::shared_ptr& data_schema, + const std::shared_ptr& write_schema, DeletionVector::Factory dv_factory, + const std::shared_ptr& path_factory_cache, + std::unique_ptr&& merge_file_split_read, + MergeFunctionWrapperFactory merge_function_wrapper_factory, + ChangelogMergeFunctionWrapperFactory changelog_merge_function_wrapper_factory, + const std::shared_ptr& cancellation_controller, + const std::shared_ptr& pool); + + bool RewriteChangelog(int32_t output_level, bool drop_delete, + const std::vector>& sections) const override { + return output_level == max_level_; + } + + UpgradeStrategy GenerateUpgradeStrategy( + int32_t output_level, const std::shared_ptr& file) const override { + return output_level == max_level_ ? UpgradeStrategy::ChangelogNoRewrite() + : UpgradeStrategy::NoChangelogNoRewrite(); + } +}; + +} // namespace paimon diff --git a/src/paimon/core/mergetree/compact/lookup_merge_tree_compact_rewriter.cpp b/src/paimon/core/mergetree/compact/lookup_merge_tree_compact_rewriter.cpp index 1a818dd01..390f95d84 100644 --- a/src/paimon/core/mergetree/compact/lookup_merge_tree_compact_rewriter.cpp +++ b/src/paimon/core/mergetree/compact/lookup_merge_tree_compact_rewriter.cpp @@ -93,8 +93,8 @@ LookupMergeTreeCompactRewriter::Create( MergeFileSplitRead::Create(path_factory, internal_context, pool, CreateDefaultExecutor())); MergeFunctionWrapperFactory merge_function_wrapper_factory = - [data_schema, options, trimmed_primary_keys, pool]( - int32_t /*output_level*/) -> Result>> { + [data_schema, options, trimmed_primary_keys, + pool]() -> Result>> { PAIMON_ASSIGN_OR_RAISE(std::unique_ptr merge_function, PrimaryKeyTableUtils::CreateMergeFunction( data_schema, trimmed_primary_keys, options, pool)); diff --git a/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory.cpp b/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory.cpp index 258336698..727a7a613 100644 --- a/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory.cpp +++ b/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory.cpp @@ -25,6 +25,7 @@ #include "paimon/core/mergetree/compact/aggregate/aggregate_merge_function.h" #include "paimon/core/mergetree/compact/early_full_compaction.h" #include "paimon/core/mergetree/compact/force_up_level0_compaction.h" +#include "paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.h" #include "paimon/core/mergetree/compact/internal_row_equalizer.h" #include "paimon/core/mergetree/compact/lookup_merge_tree_compact_rewriter.h" #include "paimon/core/mergetree/compact/merge_tree_compact_manager.h" @@ -165,7 +166,14 @@ Result> MergeTreeCompactManagerFactory::CreateR auto path_factory_cache = std::make_shared(root_path_, table_schema_, options_, pool_); if (options_.GetChangelogProducer() == ChangelogProducer::FULL_COMPACTION) { - return Status::NotImplemented("not support full changelog merge tree compact rewriter"); + int32_t max_level = options_.GetNumLevels() - 1; + auto dv_factory = DeletionVector::CreateFactory(dv_maintainer); + PAIMON_ASSIGN_OR_RAISE( + std::unique_ptr rewriter, + FullChangelogMergeTreeCompactRewriter::Create( + max_level, bucket, partition, table_schema_, std::move(dv_factory), + path_factory_cache, options_, cancellation_controller, pool_)); + return std::shared_ptr(std::move(rewriter)); } if (options_.NeedLookup()) { // Lazily create the global lookup file cache diff --git a/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory_test.cpp b/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory_test.cpp index 06a4ec197..61ef11322 100644 --- a/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory_test.cpp +++ b/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory_test.cpp @@ -330,12 +330,13 @@ TEST_F(MergeTreeCompactManagerFactoryWriteTest, } TEST_F(MergeTreeCompactManagerFactoryWriteTest, - TestCreateFileStoreWriteShouldFailWhenFullCompactionChangelogConfigured) { - ASSERT_NOK_WITH_MSG(CreateSingleStringFileStoreWrite( - {{"bucket", "1"}, {Options::CHANGELOG_PRODUCER, "full-compaction"}}, - /*with_io_manager=*/false), - "C++ Paimon only supports 'none', 'input' and 'lookup' " - "changelog-producer now"); + TestWriteShouldSucceedWhenFullCompactionChangelogConfigured) { + ASSERT_OK_AND_ASSIGN(auto file_store_write, + CreateSingleStringFileStoreWrite( + {{"bucket", "1"}, {Options::CHANGELOG_PRODUCER, "full-compaction"}}, + /*with_io_manager=*/false)); + ASSERT_OK(WriteSingleStringRow(file_store_write.get(), /*bucket=*/0, "k1")); + ASSERT_OK(file_store_write->PrepareCommit(/*wait_compaction=*/true).status()); } TEST_F(MergeTreeCompactManagerFactoryWriteTest, diff --git a/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.cpp b/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.cpp index 7224f21d4..6a4a1fd85 100644 --- a/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.cpp +++ b/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.cpp @@ -92,7 +92,7 @@ Result> MergeTreeCompactRewriter::Crea std::unique_ptr merge_file_split_read, MergeFileSplitRead::Create(path_factory, internal_context, pool, CreateDefaultExecutor())); auto merge_function_wrapper_factory = - [](int32_t output_level) -> Result>> { + []() -> Result>> { return std::shared_ptr>(); }; @@ -201,7 +201,7 @@ MergeTreeCompactRewriter::CreateRawSortMergeReaderForSection( } Status MergeTreeCompactRewriter::MergeReadAndWrite( - int32_t output_level, bool drop_delete, const std::vector& section, + bool drop_delete, const std::vector& section, const MergeTreeCompactRewriter::KeyValueConsumerCreator& create_consumer, MergeTreeCompactRewriter::KeyValueRollingFileWriter* rolling_writer) { if (!merge_file_split_read_) { @@ -212,7 +212,7 @@ Status MergeTreeCompactRewriter::MergeReadAndWrite( PAIMON_ASSIGN_OR_RAISE(std::shared_ptr data_file_path_factory, CreateDataFilePathFactory(options_.GetFileFormat()->Identifier())); PAIMON_ASSIGN_OR_RAISE(std::shared_ptr> wrapper, - merge_function_wrapper_factory_(output_level)); + merge_function_wrapper_factory_()); merge_file_split_read_->SetMergeFunctionWrapper(wrapper); PAIMON_ASSIGN_OR_RAISE(std::unique_ptr sort_merge_reader, merge_file_split_read_->CreateSortMergeReaderForSection( diff --git a/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.h b/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.h index 513987ffb..27e4932a6 100644 --- a/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.h +++ b/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.h @@ -37,7 +37,7 @@ namespace paimon { class MergeTreeCompactRewriter : public CompactRewriter { public: using MergeFunctionWrapperFactory = - std::function>>(int32_t)>; + std::function>>()>; static Result> Create( int32_t bucket, const BinaryRow& partition, @@ -95,8 +95,7 @@ class MergeTreeCompactRewriter : public CompactRewriter { Result GenerateKeyValueConsumer() const; - Status MergeReadAndWrite(int32_t output_level, bool drop_delete, - const std::vector& section, + Status MergeReadAndWrite(bool drop_delete, const std::vector& section, const KeyValueConsumerCreator& create_consumer, KeyValueRollingFileWriter* rolling_writer); diff --git a/src/paimon/core/schema/schema_validation.cpp b/src/paimon/core/schema/schema_validation.cpp index 90f508b47..5e524a376 100644 --- a/src/paimon/core/schema/schema_validation.cpp +++ b/src/paimon/core/schema/schema_validation.cpp @@ -344,11 +344,6 @@ Status SchemaValidation::ValidateChangelogProducer(const TableSchema& schema, changelog_producer == ChangelogProducer::FULL_COMPACTION, "'{}' is only valid for 'lookup' or 'full-compaction' changelog producer.", Options::CHANGELOG_PRODUCER_ROW_DEDUPLICATE)); - PAIMON_RETURN_NOT_OK(Preconditions::CheckState( - changelog_producer == ChangelogProducer::NONE || - changelog_producer == ChangelogProducer::INPUT || - changelog_producer == ChangelogProducer::LOOKUP, - "C++ Paimon only supports 'none', 'input' and 'lookup' changelog-producer now.")); return Preconditions::CheckState( options.GetMergeEngine() != MergeEngine::FIRST_ROW || changelog_producer == ChangelogProducer::NONE || diff --git a/src/paimon/core/schema/schema_validation_test.cpp b/src/paimon/core/schema/schema_validation_test.cpp index 2137970a3..054f76da3 100644 --- a/src/paimon/core/schema/schema_validation_test.cpp +++ b/src/paimon/core/schema/schema_validation_test.cpp @@ -722,8 +722,8 @@ TEST(SchemaValidationTest, ValidateDeletionVector) { std::shared_ptr table_schema, TableSchema::Create(/*schema_id=*/0, schema, partition_keys, primary_keys, options)); ASSERT_NOK_WITH_MSG(SchemaValidation::ValidateTableSchema(*table_schema), - "C++ Paimon only supports 'none', 'input' and 'lookup' " - "changelog-producer now"); + "Deletion vectors mode is only supported for " + "NONE/INPUT/LOOKUP changelog producer now"); } { std::map options = {{Options::BUCKET, "2"}, @@ -954,16 +954,6 @@ TEST(SchemaValidationTest, ValidateInvalidConfiguration) { "Only support 'none' and 'lookup' changelog-producer on FIRST_ROW " "merge engine"); } - { - std::map options = { - {Options::CHANGELOG_PRODUCER, "full-compaction"}}; - ASSERT_OK_AND_ASSIGN(std::shared_ptr table_schema, - TableSchema::Create(/*schema_id=*/0, schema, /*partition_keys=*/{}, - /*primary_keys=*/{"f0"}, options)); - ASSERT_NOK_WITH_MSG( - SchemaValidation::ValidateTableSchema(*table_schema), - "C++ Paimon only supports 'none', 'input' and 'lookup' changelog-producer now."); - } // test for row tracking { std::map options = {{Options::ROW_TRACKING_ENABLED, "true"}, diff --git a/src/paimon/core/table/source/data_table_stream_scan.cpp b/src/paimon/core/table/source/data_table_stream_scan.cpp index 1657c9c7b..bae9800d7 100644 --- a/src/paimon/core/table/source/data_table_stream_scan.cpp +++ b/src/paimon/core/table/source/data_table_stream_scan.cpp @@ -57,7 +57,11 @@ Result> DataTableStreamScan::CreatePlan() { Result> DataTableStreamScan::TryFirstPlan() { std::shared_ptr scan_result; if (core_options_.GetChangelogProducer() == ChangelogProducer::FULL_COMPACTION) { - return Status::NotImplemented("do not support full compaction changelog producer"); + int32_t max_level = core_options_.GetNumLevels() - 1; + snapshot_reader_->WithLevelFilter( + [max_level](int32_t level) -> bool { return level == max_level; }); + PAIMON_ASSIGN_OR_RAISE(scan_result, starting_scanner_->Scan(snapshot_reader_)); + snapshot_reader_->WithLevelFilter([](int32_t) -> bool { return true; }); } else if (core_options_.GetChangelogProducer() == ChangelogProducer::LOOKUP) { // Level-0 files will be compacted later to produce changelog records. Exclude them from // the initial full scan so that the same changes are not emitted both in the full phase @@ -136,11 +140,10 @@ Status DataTableStreamScan::InitScanner() { follow_up_scanner_ = std::make_shared(); return Status::OK(); case ChangelogProducer::INPUT: + case ChangelogProducer::FULL_COMPACTION: case ChangelogProducer::LOOKUP: follow_up_scanner_ = std::make_shared(); return Status::OK(); - case ChangelogProducer::FULL_COMPACTION: - return Status::NotImplemented("do not support full compaction changelog producer"); default: return Status::NotImplemented("unknown changelog producer"); } diff --git a/test/inte/write_and_read_inte_test.cpp b/test/inte/write_and_read_inte_test.cpp index dae25488a..8d7d14bbe 100644 --- a/test/inte/write_and_read_inte_test.cpp +++ b/test/inte/write_and_read_inte_test.cpp @@ -727,6 +727,249 @@ TEST_P(WriteAndReadInteTest, TestInputChangelogStreamRead) { ASSERT_TRUE(success); } +TEST_P(WriteAndReadInteTest, TestFullCompactionChangelogStreamRead) { + auto [file_format, file_system] = GetParam(); + arrow::FieldVector fields = {arrow::field("pk", arrow::utf8()), + arrow::field("value", arrow::int32())}; + std::map options = { + {Options::MANIFEST_FORMAT, "avro"}, {Options::FILE_FORMAT, file_format}, + {Options::TARGET_FILE_SIZE, "1024"}, {Options::BUCKET, "1"}, + {Options::FILE_SYSTEM, file_system}, {Options::CHANGELOG_PRODUCER, "full-compaction"}}; + if (file_system == "jindo") { + options = AddOptionsForJindo(options); + } + ASSERT_OK_AND_ASSIGN( + auto helper, + TestHelper::Create(test_dir_, arrow::schema(fields), /*partition_keys=*/{}, + /*primary_keys=*/{"pk"}, options, /*is_streaming_mode=*/true)); + + ASSERT_OK_AND_ASSIGN(std::vector> initial_splits, + helper->NewScan(StartupMode::Latest(), /*snapshot_id=*/std::nullopt)); + ASSERT_TRUE(initial_splits.empty()); + + ASSERT_OK_AND_ASSIGN( + auto initial_batch, + TestHelper::MakeRecordBatch(arrow::struct_(fields), R"([["Alice", 10], ["Bob", 20]])", + /*partition_map=*/{}, /*bucket=*/0, {})); + ASSERT_OK(helper->WriteAndCommit(std::move(initial_batch), /*commit_identifier=*/0, + /*expected_commit_messages=*/std::nullopt)); + std::string table_path = PathUtil::JoinPath(test_dir_, "foo.db/bar"); + ASSERT_OK(CompactAndCommit(table_path, options, /*commit_identifier=*/1)); + + auto expected_type = + arrow::struct_({arrow::field("_VALUE_KIND", arrow::int8()), fields[0], fields[1]}); + ASSERT_OK_AND_ASSIGN(std::vector> changelog_splits, helper->Scan()); + ASSERT_FALSE(changelog_splits.empty()); + ASSERT_OK_AND_ASSIGN(bool initial_success, + helper->ReadAndCheckResult(expected_type, changelog_splits, + R"([[0, "Alice", 10], [0, "Bob", 20]])")); + ASSERT_TRUE(initial_success); + + ASSERT_OK_AND_ASSIGN( + auto change_batch, + TestHelper::MakeRecordBatch(arrow::struct_(fields), + R"([["Alice", 11], ["Bob", 0], ["Carol", 30]])", + /*partition_map=*/{}, /*bucket=*/0, + {RecordBatch::RowKind::INSERT, RecordBatch::RowKind::DELETE, + RecordBatch::RowKind::INSERT})); + ASSERT_OK(helper->WriteAndCommit(std::move(change_batch), /*commit_identifier=*/2, + /*expected_commit_messages=*/std::nullopt)); + ASSERT_OK_AND_ASSIGN(changelog_splits, helper->Scan()); + ASSERT_TRUE(changelog_splits.empty()); + ASSERT_OK(CompactAndCommit(table_path, options, /*commit_identifier=*/3)); + + ASSERT_OK_AND_ASSIGN(changelog_splits, helper->Scan()); + ASSERT_FALSE(changelog_splits.empty()); + ASSERT_OK_AND_ASSIGN(bool update_success, + helper->ReadAndCheckResult(expected_type, changelog_splits, + R"([[1, "Alice", 10], [2, "Alice", 11], + [3, "Bob", 20], [0, "Carol", 30]])")); + ASSERT_TRUE(update_success); +} + +TEST_P(WriteAndReadInteTest, TestFullCompactionChangelogInitialScanOnlyReadsMaxLevel) { + auto [file_format, file_system] = GetParam(); + arrow::FieldVector fields = {arrow::field("pk", arrow::utf8()), + arrow::field("value", arrow::int32())}; + std::map options = { + {Options::MANIFEST_FORMAT, "avro"}, {Options::FILE_FORMAT, file_format}, + {Options::TARGET_FILE_SIZE, "1024"}, {Options::BUCKET, "1"}, + {Options::FILE_SYSTEM, file_system}, {Options::CHANGELOG_PRODUCER, "full-compaction"}}; + if (file_system == "jindo") { + options = AddOptionsForJindo(options); + } + ASSERT_OK_AND_ASSIGN( + auto helper, + TestHelper::Create(test_dir_, arrow::schema(fields), /*partition_keys=*/{}, + /*primary_keys=*/{"pk"}, options, /*is_streaming_mode=*/true)); + + ASSERT_OK_AND_ASSIGN(auto initial_batch, + TestHelper::MakeRecordBatch(arrow::struct_(fields), R"([["Alice", 10]])", + /*partition_map=*/{}, /*bucket=*/0, {})); + ASSERT_OK(helper->WriteAndCommit(std::move(initial_batch), /*commit_identifier=*/0, + /*expected_commit_messages=*/std::nullopt)); + std::string table_path = PathUtil::JoinPath(test_dir_, "foo.db/bar"); + ASSERT_OK(CompactAndCommit(table_path, options, /*commit_identifier=*/1)); + + ASSERT_OK_AND_ASSIGN( + auto pending_batch, + TestHelper::MakeRecordBatch(arrow::struct_(fields), R"([["Alice", 20], ["Bob", 30]])", + /*partition_map=*/{}, /*bucket=*/0, {})); + ASSERT_OK(helper->WriteAndCommit(std::move(pending_batch), /*commit_identifier=*/2, + /*expected_commit_messages=*/std::nullopt)); + + ASSERT_OK_AND_ASSIGN(std::vector> initial_full_splits, + helper->NewScan(StartupMode::LatestFull(), /*snapshot_id=*/std::nullopt)); + ASSERT_FALSE(initial_full_splits.empty()); + ASSERT_OK_AND_ASSIGN(CoreOptions core_options, CoreOptions::FromMap(options)); + int32_t max_level = core_options.GetNumLevels() - 1; + for (const auto& split : initial_full_splits) { + auto data_split = std::dynamic_pointer_cast(split); + ASSERT_TRUE(data_split); + for (const auto& file : data_split->DataFiles()) { + ASSERT_EQ(max_level, file->level); + } + } + + auto expected_type = + arrow::struct_({arrow::field("_VALUE_KIND", arrow::int8()), fields[0], fields[1]}); + ASSERT_OK_AND_ASSIGN( + bool success, + helper->ReadAndCheckResult(expected_type, initial_full_splits, R"([[0, "Alice", 10]])")); + ASSERT_TRUE(success); +} + +TEST_P(WriteAndReadInteTest, TestFullCompactionChangelogRowDeduplicate) { + auto [file_format, file_system] = GetParam(); + arrow::FieldVector fields = {arrow::field("pk", arrow::utf8()), + arrow::field("value", arrow::int32())}; + std::map options = { + {Options::MANIFEST_FORMAT, "avro"}, + {Options::FILE_FORMAT, file_format}, + {Options::TARGET_FILE_SIZE, "1024"}, + {Options::BUCKET, "1"}, + {Options::FILE_SYSTEM, file_system}, + {Options::CHANGELOG_PRODUCER, "full-compaction"}, + {Options::CHANGELOG_PRODUCER_ROW_DEDUPLICATE, "true"}}; + if (file_system == "jindo") { + options = AddOptionsForJindo(options); + } + ASSERT_OK_AND_ASSIGN( + auto helper, + TestHelper::Create(test_dir_, arrow::schema(fields), /*partition_keys=*/{}, + /*primary_keys=*/{"pk"}, options, /*is_streaming_mode=*/true)); + + ASSERT_OK_AND_ASSIGN(auto initial_batch, + TestHelper::MakeRecordBatch(arrow::struct_(fields), R"([["Alice", 10]])", + /*partition_map=*/{}, /*bucket=*/0, {})); + ASSERT_OK(helper->WriteAndCommit(std::move(initial_batch), /*commit_identifier=*/0, + /*expected_commit_messages=*/std::nullopt)); + std::string table_path = PathUtil::JoinPath(test_dir_, "foo.db/bar"); + ASSERT_OK(CompactAndCommit(table_path, options, /*commit_identifier=*/1)); + + ASSERT_OK_AND_ASSIGN(auto initial_splits, + helper->NewScan(StartupMode::Latest(), /*snapshot_id=*/std::nullopt)); + ASSERT_TRUE(initial_splits.empty()); + ASSERT_OK_AND_ASSIGN(auto unchanged_batch, + TestHelper::MakeRecordBatch(arrow::struct_(fields), R"([["Alice", 10]])", + /*partition_map=*/{}, /*bucket=*/0, {})); + ASSERT_OK(helper->WriteAndCommit(std::move(unchanged_batch), /*commit_identifier=*/2, + /*expected_commit_messages=*/std::nullopt)); + ASSERT_OK(CompactAndCommit(table_path, options, /*commit_identifier=*/3)); + ASSERT_OK_AND_ASSIGN(auto empty_splits, helper->Scan()); + ASSERT_TRUE(empty_splits.empty()); + + ASSERT_OK_AND_ASSIGN(auto changed_batch, + TestHelper::MakeRecordBatch(arrow::struct_(fields), R"([["Alice", 20]])", + /*partition_map=*/{}, /*bucket=*/0, {})); + ASSERT_OK(helper->WriteAndCommit(std::move(changed_batch), /*commit_identifier=*/4, + /*expected_commit_messages=*/std::nullopt)); + ASSERT_OK(CompactAndCommit(table_path, options, /*commit_identifier=*/5)); + ASSERT_OK_AND_ASSIGN(auto changelog_splits, helper->Scan()); + ASSERT_FALSE(changelog_splits.empty()); + auto expected_type = + arrow::struct_({arrow::field("_VALUE_KIND", arrow::int8()), fields[0], fields[1]}); + ASSERT_OK_AND_ASSIGN(bool success, helper->ReadAndCheckResult(expected_type, changelog_splits, + R"([[1, "Alice", 10], + [2, "Alice", 20]])")); + ASSERT_TRUE(success); +} + +TEST_P(WriteAndReadInteTest, TestFullCompactionChangelogWithSharedShredding) { + auto [file_format, file_system] = GetParam(); + if (file_format == "avro" || file_format == "mosaic") { + return; + } + + auto map_type = arrow::map(arrow::utf8(), arrow::int64()); + arrow::FieldVector fields = {arrow::field("pk", arrow::int32()), + arrow::field("tags", map_type)}; + std::map options = { + {Options::MANIFEST_FORMAT, "avro"}, + {Options::FILE_FORMAT, file_format}, + {Options::TARGET_FILE_SIZE, "1024"}, + {Options::BUCKET, "1"}, + {Options::FILE_SYSTEM, file_system}, + {Options::CHANGELOG_PRODUCER, "full-compaction"}, + {"fields.tags.map.storage-layout", "shared-shredding"}, + {"fields.tags.map.shared-shredding.max-columns", "1"}}; + if (file_system == "jindo") { + options = AddOptionsForJindo(options); + } + ASSERT_OK_AND_ASSIGN( + auto helper, + TestHelper::Create(test_dir_, arrow::schema(fields), /*partition_keys=*/{}, + /*primary_keys=*/{"pk"}, options, /*is_streaming_mode=*/true)); + + ASSERT_OK_AND_ASSIGN( + auto initial_batch, + TestHelper::MakeRecordBatch(arrow::struct_(fields), + R"([[1, [["a", 10], ["z", 11]]], [2, [["b", 20]]]])", + /*partition_map=*/{}, /*bucket=*/0, {})); + ASSERT_OK(helper->WriteAndCommit(std::move(initial_batch), /*commit_identifier=*/0, + /*expected_commit_messages=*/std::nullopt)); + std::string table_path = PathUtil::JoinPath(test_dir_, "foo.db/bar"); + ASSERT_OK(CompactAndCommit(table_path, options, /*commit_identifier=*/1)); + + ASSERT_OK_AND_ASSIGN(auto initial_splits, + helper->NewScan(StartupMode::Latest(), /*snapshot_id=*/std::nullopt)); + ASSERT_TRUE(initial_splits.empty()); + ASSERT_OK_AND_ASSIGN( + auto change_batch, + TestHelper::MakeRecordBatch( + arrow::struct_(fields), + R"([[1, [["a", 100], ["z", 101]]], [2, [["b", 20]]], [3, [["c", 30]]]])", + /*partition_map=*/{}, /*bucket=*/0, + {RecordBatch::RowKind::INSERT, RecordBatch::RowKind::DELETE, + RecordBatch::RowKind::INSERT})); + ASSERT_OK(helper->WriteAndCommit(std::move(change_batch), /*commit_identifier=*/2, + /*expected_commit_messages=*/std::nullopt)); + ASSERT_OK(CompactAndCommit(table_path, options, /*commit_identifier=*/3)); + + ASSERT_OK_AND_ASSIGN(auto changelog_splits, helper->Scan()); + ASSERT_FALSE(changelog_splits.empty()); + for (const auto& split : changelog_splits) { + auto data_split = std::dynamic_pointer_cast(split); + ASSERT_TRUE(data_split); + for (const auto& file : data_split->DataFiles()) { + ASSERT_OK_AND_ASSIGN( + MapSharedShreddingFieldMeta meta, + ReadShreddingMeta(std::make_pair(data_split->BucketPath(), file), "tags", options)); + ASSERT_EQ(1, meta.num_columns); + } + } + + auto expected_type = + arrow::struct_({arrow::field("_VALUE_KIND", arrow::int8()), fields[0], fields[1]}); + ASSERT_OK_AND_ASSIGN(bool success, + helper->ReadAndCheckResult(expected_type, changelog_splits, + R"([[1, 1, [["a", 10], ["z", 11]]], + [2, 1, [["a", 100], ["z", 101]]], + [3, 2, [["b", 20]]], + [0, 3, [["c", 30]]]])")); + ASSERT_TRUE(success); +} + TEST_P(WriteAndReadInteTest, TestLookupChangelogStreamRead) { auto [file_format, file_system] = GetParam(); arrow::FieldVector fields = { From 4c6b660ba8281617ad17f21301bc18f62fe338df Mon Sep 17 00:00:00 2001 From: "lisizhuo.lsz" Date: Mon, 31 Aug 2026 18:03:43 +0800 Subject: [PATCH 2/5] fix pre-commit --- .../compact/full_changelog_merge_tree_compact_rewriter.cpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.cpp b/src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.cpp index a85c1a0fb..9fc0744ef 100644 --- a/src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.cpp +++ b/src/paimon/core/mergetree/compact/full_changelog_merge_tree_compact_rewriter.cpp @@ -97,7 +97,7 @@ FullChangelogMergeTreeCompactRewriter::Create( } ChangelogMergeFunctionWrapperFactory changelog_merge_function_wrapper_factory = [data_schema, trimmed_primary_keys, options, max_level, value_equalizer, - pool](int32_t /*output_level*/) + pool]([[maybe_unused]] int32_t output_level) -> Result>> { PAIMON_ASSIGN_OR_RAISE(std::unique_ptr merge_function, PrimaryKeyTableUtils::CreateMergeFunction( From 6c9f3ee6c4620bd72f11b0e9e658dd4878ed0273 Mon Sep 17 00:00:00 2001 From: "lisizhuo.lsz" Date: Tue, 1 Sep 2026 10:29:43 +0800 Subject: [PATCH 3/5] fix ci --- test/inte/write_and_read_inte_test.cpp | 2 ++ 1 file changed, 2 insertions(+) diff --git a/test/inte/write_and_read_inte_test.cpp b/test/inte/write_and_read_inte_test.cpp index 8d7d14bbe..72bc37659 100644 --- a/test/inte/write_and_read_inte_test.cpp +++ b/test/inte/write_and_read_inte_test.cpp @@ -759,6 +759,8 @@ TEST_P(WriteAndReadInteTest, TestFullCompactionChangelogStreamRead) { auto expected_type = arrow::struct_({arrow::field("_VALUE_KIND", arrow::int8()), fields[0], fields[1]}); ASSERT_OK_AND_ASSIGN(std::vector> changelog_splits, helper->Scan()); + ASSERT_TRUE(changelog_splits.empty()); + ASSERT_OK_AND_ASSIGN(changelog_splits, helper->Scan()); ASSERT_FALSE(changelog_splits.empty()); ASSERT_OK_AND_ASSIGN(bool initial_success, helper->ReadAndCheckResult(expected_type, changelog_splits, From 076c99dcee89aee2aefce280ac2cbbb18a60388b Mon Sep 17 00:00:00 2001 From: "lisizhuo.lsz" Date: Tue, 1 Sep 2026 15:08:13 +0800 Subject: [PATCH 4/5] fix xinyu comment --- ..._changelog_merge_function_wrapper_test.cpp | 115 +++++++++++------- ...erge_tree_compact_manager_factory_test.cpp | 6 +- 2 files changed, 72 insertions(+), 49 deletions(-) diff --git a/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp b/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp index dd61bee85..561457d5d 100644 --- a/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp +++ b/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp @@ -32,7 +32,7 @@ namespace paimon::test { namespace { -constexpr int32_t MAX_LEVEL = 3; +constexpr int32_t kMaxLevel = 3; KeyValue MakeKeyValue(const RowKind* kind, int64_t sequence_number, int32_t level, int32_t key, int32_t value, const std::shared_ptr& pool) { @@ -52,7 +52,7 @@ std::unique_ptr CreateWrapper( const std::shared_ptr& pool, FieldsComparator::FieldComparatorFunc value_equalizer = {}) { return std::make_unique( - std::make_unique(/*ignore_delete=*/false), MAX_LEVEL, + std::make_unique(/*ignore_delete=*/false), kMaxLevel, CreateValueSerializer(pool), std::move(value_equalizer)); } @@ -71,31 +71,37 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestSingleRecord) { auto wrapper = CreateWrapper(pool); wrapper->Reset(); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, /*level=*/0, 1, 10, pool))); + ASSERT_OK( + wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, /*level=*/0, /*key=*/1, + /*value=*/10, pool))); ASSERT_OK_AND_ASSIGN(std::optional insert_result, wrapper->GetResult()); ASSERT_TRUE(insert_result); ASSERT_TRUE(insert_result->result); ASSERT_EQ(1, insert_result->changelogs.size()); - CheckKeyValue(insert_result->changelogs[0], RowKind::Insert(), 1, KeyValue::UNKNOWN_LEVEL, 10); - CheckKeyValue(*insert_result->result, RowKind::Insert(), 1, 0, 10); + CheckKeyValue(insert_result->changelogs[0], RowKind::Insert(), /*sequence_number=*/1, + /*level=*/KeyValue::UNKNOWN_LEVEL, /*value=*/10); + CheckKeyValue(*insert_result->result, RowKind::Insert(), /*sequence_number=*/1, /*level=*/0, + /*value=*/10); wrapper->Reset(); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Delete(), /*sequence_number=*/2, /*level=*/0, 2, 20, pool))); + ASSERT_OK( + wrapper->Add(MakeKeyValue(RowKind::Delete(), /*sequence_number=*/2, /*level=*/0, /*key=*/2, + /*value=*/20, pool))); ASSERT_OK_AND_ASSIGN(std::optional delete_result, wrapper->GetResult()); ASSERT_TRUE(delete_result); ASSERT_FALSE(delete_result->result); ASSERT_TRUE(delete_result->changelogs.empty()); wrapper->Reset(); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/3, MAX_LEVEL, 3, 30, pool))); + ASSERT_OK(wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/3, + /*level=*/kMaxLevel, /*key=*/3, + /*value=*/30, pool))); ASSERT_OK_AND_ASSIGN(std::optional top_level_result, wrapper->GetResult()); ASSERT_TRUE(top_level_result); ASSERT_TRUE(top_level_result->result); ASSERT_TRUE(top_level_result->changelogs.empty()); - CheckKeyValue(*top_level_result->result, RowKind::Insert(), 3, MAX_LEVEL, 30); + CheckKeyValue(*top_level_result->result, RowKind::Insert(), /*sequence_number=*/3, + /*level=*/kMaxLevel, /*value=*/30); } TEST(FullChangelogMergeFunctionWrapperTest, TestInsertUpdateAndDelete) { @@ -103,42 +109,52 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestInsertUpdateAndDelete) { auto wrapper = CreateWrapper(pool); wrapper->Reset(); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool))); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, 1, 20, pool))); + ASSERT_OK(wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, + /*level=*/kMaxLevel, /*key=*/1, + /*value=*/10, pool))); + ASSERT_OK( + wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, /*key=*/1, + /*value=*/20, pool))); ASSERT_OK_AND_ASSIGN(std::optional update_result, wrapper->GetResult()); ASSERT_TRUE(update_result); ASSERT_TRUE(update_result->result); ASSERT_EQ(2, update_result->changelogs.size()); - CheckKeyValue(update_result->changelogs[0], RowKind::UpdateBefore(), 1, KeyValue::UNKNOWN_LEVEL, - 10); - CheckKeyValue(update_result->changelogs[1], RowKind::UpdateAfter(), 2, KeyValue::UNKNOWN_LEVEL, - 20); - CheckKeyValue(*update_result->result, RowKind::Insert(), 2, 0, 20); + CheckKeyValue(update_result->changelogs[0], RowKind::UpdateBefore(), /*sequence_number=*/1, + /*level=*/KeyValue::UNKNOWN_LEVEL, /*value=*/10); + CheckKeyValue(update_result->changelogs[1], RowKind::UpdateAfter(), /*sequence_number=*/2, + /*level=*/KeyValue::UNKNOWN_LEVEL, /*value=*/20); + CheckKeyValue(*update_result->result, RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, + /*value=*/20); wrapper->Reset(); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/3, MAX_LEVEL, 2, 30, pool))); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Delete(), /*sequence_number=*/4, /*level=*/0, 2, 30, pool))); + ASSERT_OK(wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/3, + /*level=*/kMaxLevel, /*key=*/2, + /*value=*/30, pool))); + ASSERT_OK( + wrapper->Add(MakeKeyValue(RowKind::Delete(), /*sequence_number=*/4, /*level=*/0, /*key=*/2, + /*value=*/30, pool))); ASSERT_OK_AND_ASSIGN(std::optional delete_result, wrapper->GetResult()); ASSERT_TRUE(delete_result); ASSERT_FALSE(delete_result->result); ASSERT_EQ(1, delete_result->changelogs.size()); - CheckKeyValue(delete_result->changelogs[0], RowKind::Delete(), 3, KeyValue::UNKNOWN_LEVEL, 30); + CheckKeyValue(delete_result->changelogs[0], RowKind::Delete(), /*sequence_number=*/3, + /*level=*/KeyValue::UNKNOWN_LEVEL, /*value=*/30); wrapper->Reset(); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/5, /*level=*/0, 3, 40, pool))); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/6, /*level=*/0, 3, 50, pool))); + ASSERT_OK( + wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/5, /*level=*/0, /*key=*/3, + /*value=*/40, pool))); + ASSERT_OK( + wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/6, /*level=*/0, /*key=*/3, + /*value=*/50, pool))); ASSERT_OK_AND_ASSIGN(std::optional insert_result, wrapper->GetResult()); ASSERT_TRUE(insert_result); ASSERT_TRUE(insert_result->result); ASSERT_EQ(1, insert_result->changelogs.size()); - CheckKeyValue(insert_result->changelogs[0], RowKind::Insert(), 6, KeyValue::UNKNOWN_LEVEL, 50); - CheckKeyValue(*insert_result->result, RowKind::Insert(), 6, 0, 50); + CheckKeyValue(insert_result->changelogs[0], RowKind::Insert(), /*sequence_number=*/6, + /*level=*/KeyValue::UNKNOWN_LEVEL, /*value=*/50); + CheckKeyValue(*insert_result->result, RowKind::Insert(), /*sequence_number=*/6, /*level=*/0, + /*value=*/50); } TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicate) { @@ -147,9 +163,11 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicate) { auto wrapper_without_deduplicate = CreateWrapper(pool); wrapper_without_deduplicate->Reset(); ASSERT_OK(wrapper_without_deduplicate->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool))); + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, /*level=*/kMaxLevel, /*key=*/1, + /*value=*/10, pool))); ASSERT_OK(wrapper_without_deduplicate->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, 1, 10, pool))); + MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, /*key=*/1, + /*value=*/10, pool))); ASSERT_OK_AND_ASSIGN(std::optional result_without_deduplicate, wrapper_without_deduplicate->GetResult()); ASSERT_TRUE(result_without_deduplicate); @@ -161,15 +179,18 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicate) { auto wrapper = CreateWrapper(pool, std::move(value_equalizer)); wrapper->Reset(); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool))); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, 1, 10, pool))); + ASSERT_OK(wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, + /*level=*/kMaxLevel, /*key=*/1, + /*value=*/10, pool))); + ASSERT_OK( + wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, /*key=*/1, + /*value=*/10, pool))); ASSERT_OK_AND_ASSIGN(std::optional result, wrapper->GetResult()); ASSERT_TRUE(result); ASSERT_TRUE(result->result); ASSERT_TRUE(result->changelogs.empty()); - CheckKeyValue(*result->result, RowKind::Insert(), 2, 0, 10); + CheckKeyValue(*result->result, RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, + /*value=*/10); } TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicateWithIgnoreFields) { @@ -181,11 +202,11 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicateWithIgnoreFields) ASSERT_OK_AND_ASSIGN(std::unique_ptr value_serializer, RowCompactedSerializer::Create(value_schema, pool)); FullChangelogMergeFunctionWrapper wrapper( - std::make_unique(/*ignore_delete=*/false), MAX_LEVEL, + std::make_unique(/*ignore_delete=*/false), kMaxLevel, std::move(value_serializer), std::move(value_equalizer)); wrapper.Reset(); - ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, + ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/1, kMaxLevel, BinaryRowGenerator::GenerateRowPtr({1}, pool.get()), BinaryRowGenerator::GenerateRowPtr({10, 1}, pool.get())))); ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, @@ -199,7 +220,7 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestRowDeduplicateWithIgnoreFields) ASSERT_EQ(2, ignored_field_result->result->value->GetInt(1)); wrapper.Reset(); - ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/3, MAX_LEVEL, + ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/3, kMaxLevel, BinaryRowGenerator::GenerateRowPtr({1}, pool.get()), BinaryRowGenerator::GenerateRowPtr({10, 1}, pool.get())))); ASSERT_OK(wrapper.Add(KeyValue(RowKind::Insert(), /*sequence_number=*/4, /*level=*/0, @@ -218,11 +239,13 @@ TEST(FullChangelogMergeFunctionWrapperTest, TestRejectMultipleTopLevelRecords) { auto wrapper = CreateWrapper(pool); wrapper->Reset(); - ASSERT_OK(wrapper->Add( - MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, MAX_LEVEL, 1, 10, pool))); - ASSERT_NOK_WITH_MSG(wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, - MAX_LEVEL, 1, 20, pool)), - "Top level key-value already exists"); + ASSERT_OK(wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/1, + /*level=*/kMaxLevel, /*key=*/1, + /*value=*/10, pool))); + ASSERT_NOK_WITH_MSG( + wrapper->Add(MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, + /*level=*/kMaxLevel, /*key=*/1, /*value=*/20, pool)), + "Top level key-value already exists"); } } // namespace paimon::test diff --git a/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory_test.cpp b/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory_test.cpp index 61ef11322..223abafed 100644 --- a/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory_test.cpp +++ b/src/paimon/core/mergetree/compact/merge_tree_compact_manager_factory_test.cpp @@ -336,7 +336,7 @@ TEST_F(MergeTreeCompactManagerFactoryWriteTest, {{"bucket", "1"}, {Options::CHANGELOG_PRODUCER, "full-compaction"}}, /*with_io_manager=*/false)); ASSERT_OK(WriteSingleStringRow(file_store_write.get(), /*bucket=*/0, "k1")); - ASSERT_OK(file_store_write->PrepareCommit(/*wait_compaction=*/true).status()); + ASSERT_OK(file_store_write->PrepareCommit(/*wait_compaction=*/true)); } TEST_F(MergeTreeCompactManagerFactoryWriteTest, @@ -348,7 +348,7 @@ TEST_F(MergeTreeCompactManagerFactoryWriteTest, /*with_io_manager=*/true)); ASSERT_OK(WriteSingleStringRow(file_store_write.get(), /*bucket=*/0, "k1")); - ASSERT_OK(file_store_write->PrepareCommit(/*wait_compaction=*/true).status()); + ASSERT_OK(file_store_write->PrepareCommit(/*wait_compaction=*/true)); } TEST_F(MergeTreeCompactManagerFactoryWriteTest, @@ -377,7 +377,7 @@ TEST_F(MergeTreeCompactManagerFactoryWriteTest, /*with_io_manager=*/true)); ASSERT_OK(WriteStringAndInt64Row(file_store_write.get(), /*bucket=*/0, "k1", 1)); - ASSERT_OK(file_store_write->PrepareCommit(/*wait_compaction=*/true).status()); + ASSERT_OK(file_store_write->PrepareCommit(/*wait_compaction=*/true)); } TEST_F(MergeTreeCompactManagerFactoryWriteTest, From d1df3371f99f33982bead01aa5c53b1f61eec837 Mon Sep 17 00:00:00 2001 From: "lisizhuo.lsz" Date: Tue, 1 Sep 2026 16:15:25 +0800 Subject: [PATCH 5/5] fix rebase --- .../core/mergetree/compact/merge_tree_compact_rewriter.cpp | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.cpp b/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.cpp index 6a4a1fd85..afe50ddfb 100644 --- a/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.cpp +++ b/src/paimon/core/mergetree/compact/merge_tree_compact_rewriter.cpp @@ -276,8 +276,8 @@ Result MergeTreeCompactRewriter::RewriteCompaction( }); for (const auto& section : sections) { - PAIMON_RETURN_NOT_OK(MergeReadAndWrite(output_level, drop_delete, section, create_consumer, - rolling_writer.get())); + PAIMON_RETURN_NOT_OK( + MergeReadAndWrite(drop_delete, section, create_consumer, rolling_writer.get())); } PAIMON_RETURN_NOT_OK(rolling_writer->Close());