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

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 9265910e55 [core] Stream row ID reassignment manifest rewrites (#9443)
9265910e55 is described below

commit 9265910e557467a45c09b88241ddaf1f6ec821ab
Author: Jingsong Lee <[email protected]>
AuthorDate: Fri Aug 28 14:16:12 2026 +0800

    [core] Stream row ID reassignment manifest rewrites (#9443)
---
 .../DataEvolutionRowIdReassigner.java              | 60 ++++++++++++++--------
 .../DataEvolutionRowIdReassignerTest.java          | 36 +++++++++++++
 2 files changed, 76 insertions(+), 20 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionRowIdReassigner.java
 
b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionRowIdReassigner.java
index c7e5e77a69..ec12b3fd21 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionRowIdReassigner.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionRowIdReassigner.java
@@ -29,7 +29,9 @@ import org.apache.paimon.index.IndexFileMeta;
 import org.apache.paimon.manifest.FileKind;
 import org.apache.paimon.manifest.IndexManifestEntry;
 import org.apache.paimon.manifest.IndexManifestFile;
+import org.apache.paimon.manifest.ManifestAvroWriter;
 import org.apache.paimon.manifest.ManifestEntry;
+import org.apache.paimon.manifest.ManifestEntrySerializer;
 import org.apache.paimon.manifest.ManifestFile;
 import org.apache.paimon.manifest.ManifestFileMeta;
 import org.apache.paimon.manifest.ManifestList;
@@ -665,32 +667,50 @@ public class DataEvolutionRowIdReassigner {
     private RewrittenDataManifest rewriteDataManifest(
             Assignment assignment, ManifestFile manifestFile, ManifestFileMeta 
manifestMeta) {
         beforeManifestRewrite.accept(manifestMeta);
-        List<ManifestEntry> entries =
-                manifestFile.read(manifestMeta.fileName(), 
manifestMeta.fileSize());
+        ManifestEntrySerializer serializer = new ManifestEntrySerializer();
+        ManifestAvroWriter writer = manifestFile.createAvroWriter();
         long reassignedAddFileCount = 0L;
         boolean hasRewrittenEntry = false;
-        for (int i = 0; i < entries.size(); i++) {
-            ManifestEntry entry = entries.get(i);
-            RowRangeMappingIndex mapping = 
assignment.rowIdMappings.get(entry.partition());
-            if (mapping == null) {
-                continue;
-            }
-            Optional<Range> reassignedRange = 
mapping.map(entry.file().nonNullRowIdRange());
-            if (reassignedRange.isPresent()) {
-                validatePlanningEntry(entry);
-                entries.set(i, 
entry.assignFirstRowId(reassignedRange.get().from));
-                hasRewrittenEntry = true;
-                if (entry.kind() == FileKind.ADD) {
-                    reassignedAddFileCount++;
+        List<ManifestFileMeta> replacements;
+        try (CloseableIterator<ProjectedManifestEntry> entries =
+                manifestFile.scan(
+                        manifestMeta.fileName(), 
ProjectedManifestEntry.fullProjection())) {
+            while (entries.hasNext()) {
+                ProjectedManifestEntry entry = entries.next();
+                ManifestEntry output = entry;
+                RowRangeMappingIndex mapping = 
assignment.rowIdMappings.get(entry.partition());
+                if (mapping != null) {
+                    Optional<Range> reassignedRange = 
mapping.map(entry.file().nonNullRowIdRange());
+                    if (reassignedRange.isPresent()) {
+                        validatePlanningEntry(entry);
+                        output =
+                                serializer
+                                        .fromRow(entry.fullRow())
+                                        
.assignFirstRowId(reassignedRange.get().from);
+                        hasRewrittenEntry = true;
+                        if (entry.kind() == FileKind.ADD) {
+                            reassignedAddFileCount++;
+                        }
+                    }
                 }
+                writer.write(output);
             }
+            checkState(
+                    hasRewrittenEntry,
+                    "Cannot find entries to reassign in planned manifest %s.",
+                    manifestMeta.fileName());
+            writer.close();
+            replacements = writer.result();
+        } catch (RuntimeException | Error failure) {
+            writer.abort(failure);
+            throw failure;
+        } catch (Exception failure) {
+            writer.abort(failure);
+            throw new RuntimeException(
+                    "Failed to stream manifest file " + 
manifestMeta.fileName(), failure);
         }
-        checkState(
-                hasRewrittenEntry,
-                "Cannot find entries to reassign in planned manifest %s.",
-                manifestMeta.fileName());
         return new RewrittenDataManifest(
-                manifestMeta.fileName(), manifestFile.write(entries), 
reassignedAddFileCount);
+                manifestMeta.fileName(), replacements, reassignedAddFileCount);
     }
 
     private void validatePlanningEntry(ManifestEntry entry) {
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionRowIdReassignerTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionRowIdReassignerTest.java
index 0b6fbb6141..54e89507c9 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionRowIdReassignerTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionRowIdReassignerTest.java
@@ -49,6 +49,7 @@ import org.apache.paimon.manifest.ManifestFile;
 import org.apache.paimon.manifest.ManifestFileMeta;
 import org.apache.paimon.manifest.ManifestList;
 import org.apache.paimon.operation.FileStoreCommitImpl;
+import org.apache.paimon.options.MemorySize;
 import org.apache.paimon.partition.PartitionPredicate;
 import org.apache.paimon.predicate.Predicate;
 import org.apache.paimon.predicate.PredicateBuilder;
@@ -68,6 +69,7 @@ import org.apache.paimon.types.DataTypes;
 import org.apache.paimon.types.RowType;
 import org.apache.paimon.utils.Pair;
 import org.apache.paimon.utils.Range;
+import org.apache.paimon.utils.SegmentsCache;
 import org.apache.paimon.utils.SnapshotManager;
 
 import org.junit.jupiter.api.Test;
@@ -1416,6 +1418,40 @@ public class DataEvolutionRowIdReassignerTest extends 
TableTestBase {
                 .containsEntry("pt=b/", Arrays.asList(1L, 3L));
     }
 
+    @Test
+    public void testReassignStreamsManifestRewritesWithoutPopulatingCache() 
throws Exception {
+        FileStoreTable originalTable = createTableWithInterleavedPartitions();
+        FileStoreTable table =
+                originalTable.copy(
+                        
Collections.singletonMap(CoreOptions.SCAN_MANIFEST_PARALLELISM.key(), "2"));
+        List<String> originalManifestFiles = dataManifestFileNames(table);
+        assertThat(originalManifestFiles).hasSizeGreaterThan(1);
+        SegmentsCache<Path> manifestCache =
+                new SegmentsCache<>(1024, MemorySize.ofMebiBytes(64), 
Long.MAX_VALUE, null, false);
+        table.setManifestCache(manifestCache);
+
+        DataEvolutionRowIdReassigner.Result result =
+                new DataEvolutionRowIdReassigner(table)
+                        
.reassign("test-streaming-manifest-rewrite-with-cache");
+
+        Set<String> currentManifestFiles = new 
HashSet<>(dataManifestFileNames(table));
+        List<String> replacedManifestFiles = new ArrayList<>();
+        for (String fileName : originalManifestFiles) {
+            if (!currentManifestFiles.contains(fileName)) {
+                replacedManifestFiles.add(fileName);
+                assertThat(
+                                manifestCache.getIfPresents(
+                                        
table.store().pathFactory().toManifestFilePath(fileName)))
+                        .isNull();
+            }
+        }
+        assertThat(replacedManifestFiles).isNotEmpty();
+        assertThat(result.fileCount).isEqualTo(5L);
+        assertThat(rowIdsByPartition(table))
+                .containsEntry("pt=a/", Arrays.asList(5L, 6L, 7L))
+                .containsEntry("pt=b/", Arrays.asList(8L, 9L));
+    }
+
     @Test
     public void testReassignDoesNotCompactManifests() throws Exception {
         testReassignSkipsManifestOptimization(false);

Reply via email to