This is an automated email from the ASF dual-hosted git repository.

lxy-9602 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-cpp.git


The following commit(s) were added to refs/heads/main by this push:
     new a127a7ff feat(changelog): support full-compaction mode changelog 
producer (#266)
a127a7ff is described below

commit a127a7fff762a08fe57f8d05978a0acec099cacf
Author: lszskye <[email protected]>
AuthorDate: Tue Sep 1 02:28:45 2026 -0700

    feat(changelog): support full-compaction mode changelog producer (#266)
---
 include/paimon/defs.h                              |   1 -
 src/paimon/CMakeLists.txt                          |   2 +
 .../full_changelog_merge_function_wrapper.h        | 142 ++++++++++++
 .../full_changelog_merge_function_wrapper_test.cpp | 251 +++++++++++++++++++++
 .../full_changelog_merge_tree_compact_rewriter.cpp | 131 +++++++++++
 .../full_changelog_merge_tree_compact_rewriter.h   |  67 ++++++
 .../compact/lookup_merge_tree_compact_rewriter.cpp |   4 +-
 .../compact/merge_tree_compact_manager_factory.cpp |  10 +-
 .../merge_tree_compact_manager_factory_test.cpp    |  17 +-
 .../compact/merge_tree_compact_rewriter.cpp        |  10 +-
 .../compact/merge_tree_compact_rewriter.h          |   5 +-
 src/paimon/core/schema/schema_validation.cpp       |   5 -
 src/paimon/core/schema/schema_validation_test.cpp  |  14 +-
 .../core/table/source/data_table_stream_scan.cpp   |   9 +-
 test/inte/write_and_read_inte_test.cpp             | 245 ++++++++++++++++++++
 15 files changed, 873 insertions(+), 40 deletions(-)

diff --git a/include/paimon/defs.h b/include/paimon/defs.h
index e05faefd..0069ded7 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 051eba32..54b5aafb 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 00000000..f5698b58
--- /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 <memory>
+#include <optional>
+#include <utility>
+
+#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<ChangelogResult> {
+ public:
+    FullChangelogMergeFunctionWrapper(std::unique_ptr<MergeFunction>&& 
merge_function,
+                                      int32_t max_level,
+                                      
std::unique_ptr<RowCompactedSerializer>&& 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<std::optional<ChangelogResult>> GetResult() override {
+        std::optional<KeyValue> 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<ChangelogResult>(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<KeyValue> CloneKeyValue(const KeyValue& from, const RowKind* 
value_kind) const {
+        PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<Bytes> bytes,
+                               
value_serializer_->SerializeToBytes(*from.value));
+        PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<InternalRow> value,
+                               value_serializer_->Deserialize(bytes));
+        return KeyValue(value_kind, from.sequence_number, 
KeyValue::UNKNOWN_LEVEL, from.key,
+                        std::move(value));
+    }
+
+    std::unique_ptr<MergeFunction> merge_function_;
+    int32_t max_level_;
+    std::unique_ptr<RowCompactedSerializer> value_serializer_;
+    FieldsComparator::FieldComparatorFunc value_equalizer_;
+    std::optional<KeyValue> top_level_kv_;
+    std::optional<KeyValue> 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 00000000..561457d5
--- /dev/null
+++ 
b/src/paimon/core/mergetree/compact/full_changelog_merge_function_wrapper_test.cpp
@@ -0,0 +1,251 @@
+/*
+ * 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 <memory>
+#include <utility>
+
+#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 kMaxLevel = 3;
+
+KeyValue MakeKeyValue(const RowKind* kind, int64_t sequence_number, int32_t 
level, int32_t key,
+                      int32_t value, const std::shared_ptr<MemoryPool>& pool) {
+    return KeyValue(kind, sequence_number, level,
+                    BinaryRowGenerator::GenerateRowPtr({key}, pool.get()),
+                    BinaryRowGenerator::GenerateRowPtr({value}, pool.get()));
+}
+
+std::unique_ptr<RowCompactedSerializer> CreateValueSerializer(
+    const std::shared_ptr<MemoryPool>& pool) {
+    return RowCompactedSerializer::Create(arrow::schema({arrow::field("value", 
arrow::int32())}),
+                                          pool)
+        .value();
+}
+
+std::unique_ptr<FullChangelogMergeFunctionWrapper> CreateWrapper(
+    const std::shared_ptr<MemoryPool>& pool,
+    FieldsComparator::FieldComparatorFunc value_equalizer = {}) {
+    return std::make_unique<FullChangelogMergeFunctionWrapper>(
+        std::make_unique<DeduplicateMergeFunction>(/*ignore_delete=*/false), 
kMaxLevel,
+        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, /*key=*/1,
+                                  /*value=*/10, pool)));
+    ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> 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(), 
/*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, /*key=*/2,
+                                  /*value=*/20, pool)));
+    ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> 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,
+                                        /*level=*/kMaxLevel, /*key=*/3,
+                                        /*value=*/30, pool)));
+    ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> 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(), 
/*sequence_number=*/3,
+                  /*level=*/kMaxLevel, /*value=*/30);
+}
+
+TEST(FullChangelogMergeFunctionWrapperTest, TestInsertUpdateAndDelete) {
+    auto pool = GetDefaultPool();
+    auto wrapper = CreateWrapper(pool);
+
+    wrapper->Reset();
+    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<ChangelogResult> 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(), 
/*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,
+                                        /*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<ChangelogResult> 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(), 
/*sequence_number=*/3,
+                  /*level=*/KeyValue::UNKNOWN_LEVEL, /*value=*/30);
+
+    wrapper->Reset();
+    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<ChangelogResult> 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(), 
/*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) {
+    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, 
/*level=*/kMaxLevel, /*key=*/1,
+                     /*value=*/10, pool)));
+    ASSERT_OK(wrapper_without_deduplicate->Add(
+        MakeKeyValue(RowKind::Insert(), /*sequence_number=*/2, /*level=*/0, 
/*key=*/1,
+                     /*value=*/10, pool)));
+    ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> 
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,
+                                        /*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<ChangelogResult> result, 
wrapper->GetResult());
+    ASSERT_TRUE(result);
+    ASSERT_TRUE(result->result);
+    ASSERT_TRUE(result->changelogs.empty());
+    CheckKeyValue(*result->result, RowKind::Insert(), /*sequence_number=*/2, 
/*level=*/0,
+                  /*value=*/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<RowCompactedSerializer> 
value_serializer,
+                         RowCompactedSerializer::Create(value_schema, pool));
+    FullChangelogMergeFunctionWrapper wrapper(
+        std::make_unique<DeduplicateMergeFunction>(/*ignore_delete=*/false), 
kMaxLevel,
+        std::move(value_serializer), std::move(value_equalizer));
+
+    wrapper.Reset();
+    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,
+                                   BinaryRowGenerator::GenerateRowPtr({1}, 
pool.get()),
+                                   BinaryRowGenerator::GenerateRowPtr({10, 2}, 
pool.get()))));
+    ASSERT_OK_AND_ASSIGN(std::optional<ChangelogResult> 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, 
kMaxLevel,
+                                   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<ChangelogResult> 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,
+                                        /*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/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 00000000..9fc0744e
--- /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 <utility>
+
+#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<std::string>& trimmed_primary_keys, const CoreOptions& 
options,
+    const std::shared_ptr<arrow::Schema>& data_schema,
+    const std::shared_ptr<arrow::Schema>& write_schema, 
DeletionVector::Factory dv_factory,
+    const std::shared_ptr<FileStorePathFactoryCache>& path_factory_cache,
+    std::unique_ptr<MergeFileSplitRead>&& merge_file_split_read,
+    MergeFunctionWrapperFactory merge_function_wrapper_factory,
+    ChangelogMergeFunctionWrapperFactory 
changelog_merge_function_wrapper_factory,
+    const std::shared_ptr<CancellationController>& cancellation_controller,
+    const std::shared_ptr<MemoryPool>& 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<std::unique_ptr<FullChangelogMergeTreeCompactRewriter>>
+FullChangelogMergeTreeCompactRewriter::Create(
+    int32_t max_level, int32_t bucket, const BinaryRow& partition,
+    const std::shared_ptr<TableSchema>& table_schema, DeletionVector::Factory 
dv_factory,
+    const std::shared_ptr<FileStorePathFactoryCache>& path_factory_cache,
+    const CoreOptions& options,
+    const std::shared_ptr<CancellationController>& cancellation_controller,
+    const std::shared_ptr<MemoryPool>& pool) {
+    PAIMON_ASSIGN_OR_RAISE(std::vector<std::string> trimmed_primary_keys,
+                           table_schema->TrimmedPrimaryKeys());
+    std::shared_ptr<arrow::Schema> data_schema =
+        DataField::ConvertDataFieldsToArrowSchema(table_schema->Fields());
+    std::shared_ptr<arrow::Schema> 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<ReadContext> read_context,
+                           read_context_builder.Finish());
+    PAIMON_ASSIGN_OR_RAISE(
+        std::shared_ptr<InternalReadContext> internal_context,
+        InternalReadContext::Create(read_context, table_schema, 
options.ToMap()));
+    PAIMON_ASSIGN_OR_RAISE(
+        std::shared_ptr<FileStorePathFactory> path_factory,
+        
path_factory_cache->GetOrCreatePathFactory(options.GetFileFormat()->Identifier()));
+    PAIMON_ASSIGN_OR_RAISE(
+        std::unique_ptr<MergeFileSplitRead> merge_file_split_read,
+        MergeFileSplitRead::Create(path_factory, internal_context, pool, 
CreateDefaultExecutor()));
+
+    MergeFunctionWrapperFactory merge_function_wrapper_factory =
+        []() -> Result<std::shared_ptr<MergeFunctionWrapper<KeyValue>>> {
+        return std::shared_ptr<MergeFunctionWrapper<KeyValue>>();
+    };
+
+    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]([[maybe_unused]] int32_t output_level)
+        -> Result<std::shared_ptr<MergeFunctionWrapper<ChangelogResult>>> {
+        PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<MergeFunction> merge_function,
+                               PrimaryKeyTableUtils::CreateMergeFunction(
+                                   data_schema, trimmed_primary_keys, options, 
pool));
+        PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<RowCompactedSerializer> 
value_serializer,
+                               RowCompactedSerializer::Create(data_schema, 
pool));
+        std::shared_ptr<MergeFunctionWrapper<ChangelogResult>> wrapper =
+            std::make_shared<FullChangelogMergeFunctionWrapper>(
+                std::move(merge_function), max_level, 
std::move(value_serializer), value_equalizer);
+        return wrapper;
+    };
+
+    return std::unique_ptr<FullChangelogMergeTreeCompactRewriter>(
+        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<CompactResult> FullChangelogMergeTreeCompactRewriter::Rewrite(
+    int32_t output_level, bool drop_delete, const 
std::vector<std::vector<SortedRun>>& 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 00000000..70f15948
--- /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 <memory>
+
+#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<std::unique_ptr<FullChangelogMergeTreeCompactRewriter>> 
Create(
+        int32_t max_level, int32_t bucket, const BinaryRow& partition,
+        const std::shared_ptr<TableSchema>& table_schema, 
DeletionVector::Factory dv_factory,
+        const std::shared_ptr<FileStorePathFactoryCache>& path_factory_cache,
+        const CoreOptions& options,
+        const std::shared_ptr<CancellationController>& cancellation_controller,
+        const std::shared_ptr<MemoryPool>& pool);
+
+    Result<CompactResult> Rewrite(int32_t output_level, bool drop_delete,
+                                  const std::vector<std::vector<SortedRun>>& 
sections) override;
+
+ private:
+    FullChangelogMergeTreeCompactRewriter(
+        int32_t max_level, const BinaryRow& partition, int32_t bucket, int64_t 
schema_id,
+        const std::vector<std::string>& trimmed_primary_keys, const 
CoreOptions& options,
+        const std::shared_ptr<arrow::Schema>& data_schema,
+        const std::shared_ptr<arrow::Schema>& write_schema, 
DeletionVector::Factory dv_factory,
+        const std::shared_ptr<FileStorePathFactoryCache>& path_factory_cache,
+        std::unique_ptr<MergeFileSplitRead>&& merge_file_split_read,
+        MergeFunctionWrapperFactory merge_function_wrapper_factory,
+        ChangelogMergeFunctionWrapperFactory 
changelog_merge_function_wrapper_factory,
+        const std::shared_ptr<CancellationController>& cancellation_controller,
+        const std::shared_ptr<MemoryPool>& pool);
+
+    bool RewriteChangelog(int32_t output_level, bool drop_delete,
+                          const std::vector<std::vector<SortedRun>>& sections) 
const override {
+        return output_level == max_level_;
+    }
+
+    UpgradeStrategy GenerateUpgradeStrategy(
+        int32_t output_level, const std::shared_ptr<DataFileMeta>& 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 1a818dd0..390f95d8 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<T>::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<std::shared_ptr<MergeFunctionWrapper<KeyValue>>> {
+        [data_schema, options, trimmed_primary_keys,
+         pool]() -> Result<std::shared_ptr<MergeFunctionWrapper<KeyValue>>> {
         PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<MergeFunction> 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 25833669..727a7a61 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<std::shared_ptr<CompactRewriter>> 
MergeTreeCompactManagerFactory::CreateR
     auto path_factory_cache =
         std::make_shared<FileStorePathFactoryCache>(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<FullChangelogMergeTreeCompactRewriter> rewriter,
+            FullChangelogMergeTreeCompactRewriter::Create(
+                max_level, bucket, partition, table_schema_, 
std::move(dv_factory),
+                path_factory_cache, options_, cancellation_controller, pool_));
+        return std::shared_ptr<CompactRewriter>(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 06a4ec19..223abafe 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));
 }
 
 TEST_F(MergeTreeCompactManagerFactoryWriteTest,
@@ -347,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,
@@ -376,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,
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 7224f21d..afe50ddf 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<std::unique_ptr<MergeTreeCompactRewriter>> 
MergeTreeCompactRewriter::Crea
         std::unique_ptr<MergeFileSplitRead> merge_file_split_read,
         MergeFileSplitRead::Create(path_factory, internal_context, pool, 
CreateDefaultExecutor()));
     auto merge_function_wrapper_factory =
-        [](int32_t output_level) -> 
Result<std::shared_ptr<MergeFunctionWrapper<KeyValue>>> {
+        []() -> Result<std::shared_ptr<MergeFunctionWrapper<KeyValue>>> {
         return std::shared_ptr<MergeFunctionWrapper<KeyValue>>();
     };
 
@@ -201,7 +201,7 @@ 
MergeTreeCompactRewriter::CreateRawSortMergeReaderForSection(
 }
 
 Status MergeTreeCompactRewriter::MergeReadAndWrite(
-    int32_t output_level, bool drop_delete, const std::vector<SortedRun>& 
section,
+    bool drop_delete, const std::vector<SortedRun>& 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<DataFilePathFactory> 
data_file_path_factory,
                            
CreateDataFilePathFactory(options_.GetFileFormat()->Identifier()));
     PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<MergeFunctionWrapper<KeyValue>> 
wrapper,
-                           merge_function_wrapper_factory_(output_level));
+                           merge_function_wrapper_factory_());
     merge_file_split_read_->SetMergeFunctionWrapper(wrapper);
     PAIMON_ASSIGN_OR_RAISE(std::unique_ptr<SortMergeReader> sort_merge_reader,
                            
merge_file_split_read_->CreateSortMergeReaderForSection(
@@ -276,8 +276,8 @@ Result<CompactResult> 
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());
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 513987ff..27e4932a 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<Result<std::shared_ptr<MergeFunctionWrapper<KeyValue>>>(int32_t)>;
+        
std::function<Result<std::shared_ptr<MergeFunctionWrapper<KeyValue>>>()>;
 
     static Result<std::unique_ptr<MergeTreeCompactRewriter>> Create(
         int32_t bucket, const BinaryRow& partition,
@@ -95,8 +95,7 @@ class MergeTreeCompactRewriter : public CompactRewriter {
 
     Result<KeyValueConsumerCreator> GenerateKeyValueConsumer() const;
 
-    Status MergeReadAndWrite(int32_t output_level, bool drop_delete,
-                             const std::vector<SortedRun>& section,
+    Status MergeReadAndWrite(bool drop_delete, const std::vector<SortedRun>& 
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 90f508b4..5e524a37 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 2137970a..054f76da 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<TableSchema> 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<std::string, std::string> 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<std::string, std::string> options = {
-            {Options::CHANGELOG_PRODUCER, "full-compaction"}};
-        ASSERT_OK_AND_ASSIGN(std::shared_ptr<TableSchema> 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<std::string, std::string> 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 1657c9c7..bae9800d 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<std::shared_ptr<Plan>> 
DataTableStreamScan::CreatePlan() {
 Result<std::shared_ptr<Plan>> DataTableStreamScan::TryFirstPlan() {
     std::shared_ptr<StartingScanner::ScanResult> 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<DeltaFollowUpScanner>();
             return Status::OK();
         case ChangelogProducer::INPUT:
+        case ChangelogProducer::FULL_COMPACTION:
         case ChangelogProducer::LOOKUP:
             follow_up_scanner_ = std::make_shared<ChangelogFollowUpScanner>();
             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 dae25488..72bc3765 100644
--- a/test/inte/write_and_read_inte_test.cpp
+++ b/test/inte/write_and_read_inte_test.cpp
@@ -727,6 +727,251 @@ 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<std::string, std::string> 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<std::shared_ptr<Split>> 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<std::shared_ptr<Split>> 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,
+                                                    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<std::string, std::string> 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<std::shared_ptr<Split>> 
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<DataSplitImpl>(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<std::string, std::string> 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<std::string, std::string> 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<DataSplitImpl>(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 = {

Reply via email to