Skip to content

Commit 8481ef3

Browse files
committed
fix
1 parent ccda76a commit 8481ef3

3 files changed

Lines changed: 13 additions & 23 deletions

File tree

src/paimon/core/operation/commit/overwrite_changes_provider.cpp

Lines changed: 1 addition & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -16,10 +16,8 @@
1616

1717
#include "paimon/core/operation/commit/overwrite_changes_provider.h"
1818

19-
#include <unordered_set>
2019
#include <utility>
2120

22-
#include "paimon/core/manifest/file_entry.h"
2321
#include "paimon/core/manifest/file_kind.h"
2422

2523
namespace paimon {
@@ -46,21 +44,12 @@ Result<std::shared_ptr<CommitChanges>> OverwriteChangesProvider::Provide(
4644

4745
PAIMON_ASSIGN_OR_RAISE(std::vector<ManifestEntry> entries,
4846
manifest_scan_(latest_snapshot.value()));
49-
std::unordered_set<FileEntry::Identifier> existing_identifiers;
50-
existing_identifiers.reserve(entries.size());
5147
for (const auto& entry : entries) {
52-
existing_identifiers.insert(entry.CreateIdentifier());
5348
delta_files.emplace_back(FileKind::Delete(), entry.Partition(), entry.Bucket(),
5449
entry.TotalBuckets(), entry.File());
5550
}
5651

57-
for (const auto& change : changes_) {
58-
if (change.Kind() == FileKind::Add() &&
59-
existing_identifiers.find(change.CreateIdentifier()) != existing_identifiers.end()) {
60-
continue;
61-
}
62-
delta_files.push_back(change);
63-
}
52+
delta_files.insert(delta_files.end(), changes_.begin(), changes_.end());
6453

6554
PAIMON_ASSIGN_OR_RAISE(std::vector<IndexManifestEntry> previous_index_entries,
6655
index_scan_(latest_snapshot.value()));

src/paimon/core/operation/commit/overwrite_changes_provider_test.cpp

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -138,7 +138,7 @@ TEST(OverwriteChangesProviderTest, TestProvideWithoutLatestSnapshotUsesChangesOn
138138
EXPECT_EQ("index-new", provided_index_entries[0].index_file->FileName());
139139
}
140140

141-
TEST(OverwriteChangesProviderTest, TestProvideWithLatestSnapshotAddsDeletesAndDeduplicatesAdds) {
141+
TEST(OverwriteChangesProviderTest, TestProvideWithLatestSnapshotAddsDeletesAndAppendsAllChanges) {
142142
std::vector<ManifestEntry> changes = {
143143
CreateManifestEntry("existing-a", FileKind::Add(), /*partition_value=*/1),
144144
CreateManifestEntry("new-b", FileKind::Add(), /*partition_value=*/1),
@@ -167,16 +167,19 @@ TEST(OverwriteChangesProviderTest, TestProvideWithLatestSnapshotAddsDeletesAndDe
167167

168168
ASSERT_TRUE(changelog_files.empty());
169169

170-
// delta = delete(existing-a), delete(existing-x), add(new-b), delete(force-delete-c)
171-
ASSERT_EQ(4u, delta_files.size());
170+
// delta = delete(existing-a), delete(existing-x),
171+
// add(existing-a), add(new-b), delete(force-delete-c)
172+
ASSERT_EQ(5u, delta_files.size());
172173
EXPECT_TRUE(delta_files[0].Kind() == FileKind::Delete());
173174
EXPECT_EQ("existing-a", delta_files[0].FileName());
174175
EXPECT_TRUE(delta_files[1].Kind() == FileKind::Delete());
175176
EXPECT_EQ("existing-x", delta_files[1].FileName());
176177
EXPECT_TRUE(delta_files[2].Kind() == FileKind::Add());
177-
EXPECT_EQ("new-b", delta_files[2].FileName());
178-
EXPECT_TRUE(delta_files[3].Kind() == FileKind::Delete());
179-
EXPECT_EQ("force-delete-c", delta_files[3].FileName());
178+
EXPECT_EQ("existing-a", delta_files[2].FileName());
179+
EXPECT_TRUE(delta_files[3].Kind() == FileKind::Add());
180+
EXPECT_EQ("new-b", delta_files[3].FileName());
181+
EXPECT_TRUE(delta_files[4].Kind() == FileKind::Delete());
182+
EXPECT_EQ("force-delete-c", delta_files[4].FileName());
180183

181184
// index = provided new + delete old
182185
ASSERT_EQ(2u, provided_index_entries.size());

src/paimon/core/operation/file_store_commit_impl_test.cpp

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -2108,11 +2108,9 @@ TEST_F(FileStoreCommitImplTest, TestOverwriteWithSameFile) {
21082108
ASSERT_TRUE(IsStringInSet(file_names, "data-6828284c-e707-49b5-af6b-69be79af120c-0.orc"));
21092109
ASSERT_TRUE(IsStringInSet(file_names, "data-8dc7f04c-3c98-48b2-9d56-834d746c4a40-0.orc"));
21102110

2111-
// same file delete, then add, file will also be removed in result
2112-
ASSERT_OK(commit_impl->Overwrite({}, msgs1, 2));
2113-
ASSERT_OK_AND_ASSIGN(auto snapshot2, commit_impl->snapshot_manager_->LatestSnapshot());
2114-
ASSERT_OK_AND_ASSIGN(auto entries2, commit_impl->GetAllFiles(snapshot2.value(), {}));
2115-
ASSERT_EQ(0u, entries2.size());
2111+
// Java parity: overwrite provider adds DELETE(old) + ADD(newChanges) without dedup,
2112+
// so same-file overwrite fails on duplicate add.
2113+
ASSERT_NOK_WITH_MSG(commit_impl->Overwrite({}, msgs1, 2), "Trying to add file");
21162114
}
21172115

21182116
TEST_F(FileStoreCommitImplTest, TestAppendDiscardDuplicateFiles) {

0 commit comments

Comments
 (0)