3333#include " paimon/core/manifest/manifest_entry.h"
3434#include " paimon/core/operation/commit/overwrite_changes_provider.h"
3535#include " paimon/core/operation/file_store_scan.h"
36+ #include " paimon/core/table/bucket_mode.h"
3637#include " paimon/scan_context.h"
3738
3839namespace paimon {
@@ -90,13 +91,8 @@ Result<std::vector<ManifestEntry>> CommitScanner::ReadAllEntriesFromChangedParti
9091 std::vector<std::map<std::string, std::string>> partition_filters;
9192 PAIMON_ASSIGN_OR_RAISE (partition_filters, ToPartitionFilters (changed_partitions));
9293
93- auto scan_filter = std::make_shared<ScanFilter>(/* predicate=*/ nullptr , partition_filters,
94- /* bucket_filter=*/ std::nullopt );
95- if (!scan_supplier_) {
96- return Status::Invalid (" CommitScanner requires non-empty scan supplier." );
97- }
98-
99- PAIMON_ASSIGN_OR_RAISE (std::unique_ptr<FileStoreScan> scan, scan_supplier_ (scan_filter));
94+ PAIMON_ASSIGN_OR_RAISE (std::unique_ptr<FileStoreScan> scan,
95+ NewScan (partition_filters, /* for_overwrite=*/ false ));
10096 PAIMON_ASSIGN_OR_RAISE (std::shared_ptr<FileStoreScan::RawPlan> plan,
10197 scan->WithSnapshot (snapshot)->WithKind (ScanMode::ALL )->CreatePlan ());
10298 return plan->Files ();
@@ -105,16 +101,27 @@ Result<std::vector<ManifestEntry>> CommitScanner::ReadAllEntriesFromChangedParti
105101Result<std::vector<ManifestEntry>> CommitScanner::ReadAllEntriesFromPartitions (
106102 const Snapshot& snapshot,
107103 const std::vector<std::map<std::string, std::string>>& partitions) const {
104+ PAIMON_ASSIGN_OR_RAISE (std::unique_ptr<FileStoreScan> scan,
105+ NewScan (partitions, /* for_overwrite=*/ false ));
106+ PAIMON_ASSIGN_OR_RAISE (std::shared_ptr<FileStoreScan::RawPlan> plan,
107+ scan->WithSnapshot (snapshot)->WithKind (ScanMode::ALL )->CreatePlan ());
108+ return plan->Files ();
109+ }
110+
111+ Result<std::unique_ptr<FileStoreScan>> CommitScanner::NewScan (
112+ const std::vector<std::map<std::string, std::string>>& partitions,
113+ bool for_overwrite) const {
108114 auto scan_filter = std::make_shared<ScanFilter>(/* predicate=*/ nullptr , partitions,
109115 /* bucket_filter=*/ std::nullopt );
110116 if (!scan_supplier_) {
111117 return Status::Invalid (" CommitScanner requires non-empty scan supplier." );
112118 }
113119
114120 PAIMON_ASSIGN_OR_RAISE (std::unique_ptr<FileStoreScan> scan, scan_supplier_ (scan_filter));
115- PAIMON_ASSIGN_OR_RAISE (std::shared_ptr<FileStoreScan::RawPlan> plan,
116- scan->WithSnapshot (snapshot)->WithKind (ScanMode::ALL )->CreatePlan ());
117- return plan->Files ();
121+ if (for_overwrite && core_options_.GetBucket () != BucketModeDefine::POSTPONE_BUCKET ) {
122+ scan->OnlyReadRealBuckets ();
123+ }
124+ return scan;
118125}
119126
120127Result<std::vector<IndexManifestEntry>> CommitScanner::ReadAllIndexEntriesFromPartitions (
@@ -165,8 +172,14 @@ std::shared_ptr<CommitChangesProvider> CommitScanner::OverwriteChangesProvider(
165172 const std::vector<IndexManifestEntry>& index_entries) const {
166173 return std::make_shared<paimon::OverwriteChangesProvider>(
167174 changes, index_entries,
168- [this , partitions](const Snapshot& snapshot) {
169- return ReadAllEntriesFromPartitions (snapshot, partitions);
175+ [this , partitions](const Snapshot& snapshot) -> Result<std::vector<ManifestEntry>> {
176+ PAIMON_ASSIGN_OR_RAISE (std::unique_ptr<FileStoreScan> scan,
177+ NewScan (partitions, /* for_overwrite=*/ true ));
178+ PAIMON_ASSIGN_OR_RAISE (std::shared_ptr<FileStoreScan::RawPlan> plan,
179+ scan->WithSnapshot (snapshot)
180+ ->WithKind (ScanMode::ALL )
181+ ->CreatePlan ());
182+ return plan->Files ();
170183 },
171184 [this , partitions](const Snapshot& snapshot) {
172185 return ReadAllIndexEntriesFromPartitions (snapshot, partitions);
0 commit comments