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 d69a7e654e [spark] Retry data evolution compaction after merge 
conflicts (#9197)
d69a7e654e is described below

commit d69a7e654e9800a31ea42aae1472a937a39239a4
Author: Jingsong Lee <[email protected]>
AuthorDate: Thu Aug 13 20:05:36 2026 +0800

    [spark] Retry data evolution compaction after merge conflicts (#9197)
---
 .../org/apache/paimon/append/AppendOnlyWriter.java |  70 ++-
 .../DataEvolutionNormalCompactTask.java            |   2 +
 .../paimon/operation/BaseAppendFileStoreWrite.java |   9 +-
 .../commit/DataEvolutionConflictDetection.java     |   2 +-
 .../DataEvolutionRowRangeConflictException.java    |  27 +
 .../operation/commit/ConflictDetectionTest.java    |  21 +
 .../paimon/table/DataEvolutionTableTest.java       |   2 +
 .../paimon/spark/procedure/CompactProcedure.java   | 111 +++-
 .../procedure/DataEvolutionRewriteExecutor.java    | 229 +++++++-
 ...DataEvolutionCompactMergeConflictRewriter.scala | 458 +++++++++++++++
 .../spark/procedure/CompactProcedureTestBase.scala | 628 ++++++++++++++++++++-
 11 files changed, 1518 insertions(+), 41 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/append/AppendOnlyWriter.java 
b/paimon-core/src/main/java/org/apache/paimon/append/AppendOnlyWriter.java
index 5f47789a26..98bbc618d6 100644
--- a/paimon-core/src/main/java/org/apache/paimon/append/AppendOnlyWriter.java
+++ b/paimon-core/src/main/java/org/apache/paimon/append/AppendOnlyWriter.java
@@ -86,6 +86,7 @@ public class AppendOnlyWriter implements BatchRecordWriter, 
MemoryOwner {
     private final boolean forceCompact;
     private final boolean asyncFileWrite;
     private final boolean statsDenseStore;
+    private final FileSource fileSource;
     @Nullable private final FileFormat rowSidecarFileFormat;
     @Nullable private final BlobFileContext blobContext;
     private final List<DataFileMeta> newFiles;
@@ -134,6 +135,70 @@ public class AppendOnlyWriter implements 
BatchRecordWriter, MemoryOwner {
             boolean dataEvolutionEnabled,
             @Nullable FileFormat rowSidecarFileFormat,
             @Nullable BlobFileContext blobContext) {
+        this(
+                fileIO,
+                ioManager,
+                schemaId,
+                fileFormat,
+                vectorFileFormat,
+                targetFileSize,
+                blobTargetFileSize,
+                vectorTargetFileSize,
+                targetFileRowNum,
+                writeSchema,
+                writeCols,
+                maxSequenceNumber,
+                compactManager,
+                dataFileRead,
+                forceCompact,
+                pathFactory,
+                increment,
+                useWriteBuffer,
+                spillable,
+                fileCompression,
+                spillCompression,
+                statsCollectorFactories,
+                maxDiskSize,
+                fileIndexOptions,
+                asyncFileWrite,
+                statsDenseStore,
+                dataEvolutionEnabled,
+                rowSidecarFileFormat,
+                blobContext,
+                FileSource.APPEND);
+    }
+
+    public AppendOnlyWriter(
+            FileIO fileIO,
+            @Nullable IOManager ioManager,
+            long schemaId,
+            FileFormat fileFormat,
+            @Nullable FileFormat vectorFileFormat,
+            long targetFileSize,
+            long blobTargetFileSize,
+            long vectorTargetFileSize,
+            long targetFileRowNum,
+            RowType writeSchema,
+            @Nullable List<String> writeCols,
+            long maxSequenceNumber,
+            CompactManager compactManager,
+            IOFunction<List<DataFileMeta>, RecordReaderIterator<InternalRow>> 
dataFileRead,
+            boolean forceCompact,
+            DataFilePathFactory pathFactory,
+            @Nullable CommitIncrement increment,
+            boolean useWriteBuffer,
+            boolean spillable,
+            String fileCompression,
+            CompressOptions spillCompression,
+            StatsCollectorFactories statsCollectorFactories,
+            MemorySize maxDiskSize,
+            FileIndexOptions fileIndexOptions,
+            boolean asyncFileWrite,
+            boolean statsDenseStore,
+            boolean dataEvolutionEnabled,
+            @Nullable FileFormat rowSidecarFileFormat,
+            @Nullable BlobFileContext blobContext,
+            FileSource fileSource) {
         this.fileIO = fileIO;
         this.schemaId = schemaId;
         this.fileFormat = fileFormat;
@@ -150,6 +215,7 @@ public class AppendOnlyWriter implements BatchRecordWriter, 
MemoryOwner {
         this.forceCompact = forceCompact;
         this.asyncFileWrite = asyncFileWrite;
         this.statsDenseStore = statsDenseStore;
+        this.fileSource = fileSource;
         this.rowSidecarFileFormat = dataEvolutionEnabled ? 
rowSidecarFileFormat : null;
         this.blobContext = blobContext;
         this.newFiles = new ArrayList<>();
@@ -335,7 +401,7 @@ public class AppendOnlyWriter implements BatchRecordWriter, 
MemoryOwner {
                     fileCompression,
                     statsCollectorFactories,
                     fileIndexOptions,
-                    FileSource.APPEND,
+                    fileSource,
                     statsDenseStore,
                     blobContext);
         }
@@ -350,7 +416,7 @@ public class AppendOnlyWriter implements BatchRecordWriter, 
MemoryOwner {
                 fileCompression,
                 
statsCollectorFactories.statsCollectors(writeSchema.getFieldNames()),
                 fileIndexOptions,
-                FileSource.APPEND,
+                fileSource,
                 asyncFileWrite,
                 statsDenseStore,
                 writeCols,
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java
 
b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java
index 9dc611217f..9f4e41c34d 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java
@@ -23,6 +23,7 @@ import org.apache.paimon.CoreOptions;
 import org.apache.paimon.data.BinaryRow;
 import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.manifest.FileSource;
 import org.apache.paimon.operation.AppendFileStoreWrite;
 import org.apache.paimon.reader.RecordReader;
 import org.apache.paimon.table.FileStoreTable;
@@ -98,6 +99,7 @@ public class DataEvolutionNormalCompactTask extends 
DataEvolutionCompactTask {
                 
store.newDataEvolutionRead().withReadType(readWriteType).createReader(dataSplit);
         AppendFileStoreWrite storeWrite = (AppendFileStoreWrite) 
store.newWrite(commitUser);
         storeWrite.withWriteType(readWriteType);
+        storeWrite.withFileSource(FileSource.COMPACT);
         RecordWriter<InternalRow> writer = storeWrite.createWriter(partition, 
0);
 
         reader.forEachRemaining(
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
 
b/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
index c3ad22b65c..cb2dd42cfa 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
@@ -86,6 +86,7 @@ public abstract class BaseAppendFileStoreWrite extends 
MemoryFileStoreWrite<Inte
     private @Nullable BlobFetchMetrics blobFetchMetrics;
     private RowType writeType;
     private @Nullable List<String> writeCols;
+    private FileSource fileSource = FileSource.APPEND;
     private boolean forceBufferSpill = false;
 
     public BaseAppendFileStoreWrite(
@@ -181,7 +182,13 @@ public abstract class BaseAppendFileStoreWrite extends 
MemoryFileStoreWrite<Inte
                 options.statsDenseStore(),
                 options.dataEvolutionEnabled(),
                 rowSidecarFileFormat(),
-                blobContext);
+                blobContext,
+                fileSource);
+    }
+
+    public BaseAppendFileStoreWrite withFileSource(FileSource fileSource) {
+        this.fileSource = fileSource;
+        return this;
     }
 
     @Override
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionConflictDetection.java
 
b/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionConflictDetection.java
index c04e3d4651..bfbe8a4f57 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionConflictDetection.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionConflictDetection.java
@@ -361,7 +361,7 @@ public class DataEvolutionConflictDetection extends 
ConflictDetection {
         for (List<SimpleFileEntry> dataFileGroup : 
rangeHelper.mergeOverlappingRanges(dataFiles)) {
             if (!rangeHelper.areAllRangesSame(dataFileGroup)) {
                 return Optional.of(
-                        new RuntimeException(
+                        new DataEvolutionRowRangeConflictException(
                                 "For Data Evolution table, multiple 'MERGE 
INTO' and 'COMPACT' "
                                         + "operations "
                                         + "have encountered conflicts, data 
files: "
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionRowRangeConflictException.java
 
b/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionRowRangeConflictException.java
new file mode 100644
index 0000000000..f958ff0fef
--- /dev/null
+++ 
b/paimon-core/src/main/java/org/apache/paimon/operation/commit/DataEvolutionRowRangeConflictException.java
@@ -0,0 +1,27 @@
+/*
+ * 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.
+ */
+
+package org.apache.paimon.operation.commit;
+
+/** Conflict caused by incompatible row-range boundaries in a data-evolution 
table. */
+public final class DataEvolutionRowRangeConflictException extends 
RuntimeException {
+
+    public DataEvolutionRowRangeConflictException(String message) {
+        super(message);
+    }
+}
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/operation/commit/ConflictDetectionTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/operation/commit/ConflictDetectionTest.java
index e200faaa0e..21e9f71ab9 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/operation/commit/ConflictDetectionTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/operation/commit/ConflictDetectionTest.java
@@ -1186,6 +1186,26 @@ class ConflictDetectionTest {
         assertThat(result.get()).hasMessageContaining("Row ID existence 
conflict");
     }
 
+    @Test
+    void testCheckRowIdRangeConflictsUsesRetryableExceptionForDataFiles() {
+        DataEvolutionConflictDetection detection = createConflictDetection();
+
+        Optional<RuntimeException> exception =
+                detection.checkConflicts(
+                        snapshot(1),
+                        Arrays.asList(
+                                createFileEntryWithRowId("f1", ADD, 0L, 2L),
+                                createFileEntryWithRowId("f2", ADD, 2L, 2L)),
+                        Collections.singletonList(
+                                createFileEntryWithRowId("compacted", ADD, 0L, 
4L)),
+                        Collections.emptyList(),
+                        null,
+                        Snapshot.CommitKind.COMPACT);
+
+        assertThat(exception).isPresent();
+        
assertThat(exception.get()).isInstanceOf(DataEvolutionRowRangeConflictException.class);
+    }
+
     @Test
     void testCheckRowIdRangeConflictsReportsDedicatedFileSpanningDataFiles() {
         DataEvolutionConflictDetection detection = createConflictDetection();
@@ -1203,6 +1223,7 @@ class ConflictDetectionTest {
 
         assertThat(exception).isPresent();
         assertThat(exception.get())
+                .isNotInstanceOf(DataEvolutionRowRangeConflictException.class)
                 .hasMessageContaining("dedicated file")
                 .hasMessageContaining("p1.blob")
                 .hasMessageContaining("spans multiple data file ranges")
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionTableTest.java 
b/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionTableTest.java
index c90c69c8cb..2cd688c9e9 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionTableTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/DataEvolutionTableTest.java
@@ -32,6 +32,7 @@ import org.apache.paimon.index.IndexFileMeta;
 import org.apache.paimon.index.IndexPathFactory;
 import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.io.DataFilePathFactory;
+import org.apache.paimon.manifest.FileSource;
 import org.apache.paimon.manifest.ManifestEntry;
 import org.apache.paimon.manifest.ManifestFileMeta;
 import org.apache.paimon.partition.PartitionPredicate;
@@ -1646,6 +1647,7 @@ public class DataEvolutionTableTest extends 
DataEvolutionTestBase {
         }
 
         assertThat(entries.size()).isEqualTo(1);
+        
assertThat(entries.get(0).file().fileSource()).contains(FileSource.COMPACT);
         assertThat(entries.get(0).file().nonNullFirstRowId()).isEqualTo(0);
         assertThat(entries.get(0).file().rowCount()).isEqualTo(500000L);
     }
diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
index b2817bce5c..f9368cb50a 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
@@ -38,6 +38,7 @@ import org.apache.paimon.operation.BaseAppendFileStoreWrite;
 import org.apache.paimon.partition.PartitionPredicate;
 import org.apache.paimon.partition.PartitionValuesTimeExpireStrategy;
 import org.apache.paimon.spark.SparkUtils;
+import 
org.apache.paimon.spark.commands.DataEvolutionCompactMergeConflictRewriter;
 import org.apache.paimon.spark.commands.PaimonSparkWriter;
 import org.apache.paimon.spark.sort.TableSorter;
 import org.apache.paimon.spark.util.ScanPlanHelper$;
@@ -94,6 +95,7 @@ import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.concurrent.atomic.AtomicReference;
 import java.util.function.Function;
 import java.util.function.Predicate;
 import java.util.stream.Collectors;
@@ -293,7 +295,11 @@ public class CompactProcedure extends BaseProcedure {
                 case BUCKET_UNAWARE:
                     if (table.coreOptions().dataEvolutionEnabled()) {
                         compactDataEvolutionTable(
-                                table, partitionPredicate, partitionIdleTime, 
javaSparkContext);
+                                table,
+                                relation,
+                                partitionPredicate,
+                                partitionIdleTime,
+                                javaSparkContext);
                     } else if (clusterIncrementalEnabled) {
                         clusterIncrementalUnAwareBucketTable(
                                 table, partitionPredicate, fullCompact, 
relation);
@@ -548,11 +554,30 @@ public class CompactProcedure extends BaseProcedure {
 
     private void compactDataEvolutionTable(
             FileStoreTable table,
+            DataSourceV2Relation relation,
             @Nullable PartitionPredicate partitionPredicate,
             @Nullable Duration partitionIdleTime,
             JavaSparkContext javaSparkContext) {
         executeDataEvolutionCompaction(
-                table, partitionPredicate, partitionIdleTime, 
javaSparkContext, spark());
+                table, relation, partitionPredicate, partitionIdleTime, 
javaSparkContext, spark());
+    }
+
+    static void executeDataEvolutionCompaction(
+            FileStoreTable table,
+            DataSourceV2Relation relation,
+            @Nullable PartitionPredicate partitionPredicate,
+            @Nullable Duration partitionIdleTime,
+            JavaSparkContext javaSparkContext,
+            SparkSession sparkSession) {
+        executeDataEvolutionCompaction(
+                table,
+                relation,
+                partitionPredicate,
+                partitionIdleTime,
+                javaSparkContext,
+                sparkSession,
+                null,
+                commit -> {});
     }
 
     static void executeDataEvolutionCompaction(
@@ -562,7 +587,14 @@ public class CompactProcedure extends BaseProcedure {
             JavaSparkContext javaSparkContext,
             SparkSession sparkSession) {
         executeDataEvolutionCompaction(
-                table, partitionPredicate, partitionIdleTime, 
javaSparkContext, sparkSession, null);
+                table,
+                null,
+                partitionPredicate,
+                partitionIdleTime,
+                javaSparkContext,
+                sparkSession,
+                null,
+                commit -> {});
     }
 
     static void executeDataEvolutionCompaction(
@@ -572,33 +604,68 @@ public class CompactProcedure extends BaseProcedure {
             JavaSparkContext javaSparkContext,
             SparkSession sparkSession,
             @Nullable Integer candidateFilesPerBatch) {
+        executeDataEvolutionCompaction(
+                table,
+                null,
+                partitionPredicate,
+                partitionIdleTime,
+                javaSparkContext,
+                sparkSession,
+                candidateFilesPerBatch,
+                commit -> {});
+    }
+
+    static void executeDataEvolutionCompaction(
+            FileStoreTable table,
+            @Nullable DataSourceV2Relation relation,
+            @Nullable PartitionPredicate partitionPredicate,
+            @Nullable Duration partitionIdleTime,
+            JavaSparkContext javaSparkContext,
+            SparkSession sparkSession,
+            @Nullable Integer candidateFilesPerBatch,
+            DataEvolutionRewriteExecutor.CommitConfigurer commitConfigurer) {
         DataEvolutionCompactCoordinator.validateOptions(table.coreOptions());
         Snapshot snapshot = table.snapshotManager().latestSnapshot();
         if (snapshot == null) {
             LOG.info("Table {} has no snapshot yet, skip this compact job.", 
table.fullName());
             return;
         }
-        DataEvolutionCompactCoordinator coordinator =
-                candidateFilesPerBatch == null
-                        ? new DataEvolutionCompactCoordinator(
-                                table,
-                                partitionPredicate,
-                                table.coreOptions().blobCompactionEnabled(),
-                                false,
-                                snapshot)
-                        : new DataEvolutionCompactCoordinator(
-                                table,
-                                partitionPredicate,
-                                table.coreOptions().blobCompactionEnabled(),
-                                false,
-                                snapshot,
-                                candidateFilesPerBatch);
+        AtomicReference<DataEvolutionCompactCoordinator> coordinatorRef = new 
AtomicReference<>();
         Function<Snapshot, List<DataEvolutionCompactTask>> taskPlanner =
-                ignored ->
-                        filterIdlePartitions(
-                                coordinator.plan(), table, partitionPredicate, 
partitionIdleTime);
+                planningSnapshot -> {
+                    DataEvolutionCompactCoordinator coordinator = 
coordinatorRef.get();
+                    if (coordinator == null
+                            || coordinator.snapshot().id() != 
planningSnapshot.id()) {
+                        coordinator =
+                                candidateFilesPerBatch == null
+                                        ? new DataEvolutionCompactCoordinator(
+                                                table,
+                                                partitionPredicate,
+                                                
table.coreOptions().blobCompactionEnabled(),
+                                                false,
+                                                planningSnapshot)
+                                        : new DataEvolutionCompactCoordinator(
+                                                table,
+                                                partitionPredicate,
+                                                
table.coreOptions().blobCompactionEnabled(),
+                                                false,
+                                                planningSnapshot,
+                                                candidateFilesPerBatch);
+                        coordinatorRef.set(coordinator);
+                    }
+                    return filterIdlePartitions(
+                            coordinator.plan(), table, partitionPredicate, 
partitionIdleTime);
+                };
         DataEvolutionRewriteExecutor.execute(
-                table, snapshot, taskPlanner, javaSparkContext, sparkSession, 
commit -> {});
+                table,
+                snapshot,
+                taskPlanner,
+                javaSparkContext,
+                sparkSession,
+                commitConfigurer,
+                relation == null
+                        ? null
+                        : new DataEvolutionCompactMergeConflictRewriter(table, 
relation)::rewrite);
     }
 
     private static List<DataEvolutionCompactTask> filterIdlePartitions(
diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/DataEvolutionRewriteExecutor.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/DataEvolutionRewriteExecutor.java
index 0fa3dd52fe..b664b82f7e 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/DataEvolutionRewriteExecutor.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/DataEvolutionRewriteExecutor.java
@@ -22,11 +22,14 @@ import org.apache.paimon.Snapshot;
 import org.apache.paimon.append.dataevolution.DataEvolutionCompactTask;
 import 
org.apache.paimon.append.dataevolution.DataEvolutionCompactTaskSerializer;
 import 
org.apache.paimon.append.dataevolution.DataEvolutionCompactionCommitPreparation;
+import 
org.apache.paimon.operation.commit.DataEvolutionRowRangeConflictException;
 import org.apache.paimon.table.FileStoreTable;
 import org.apache.paimon.table.sink.CommitMessage;
 import org.apache.paimon.table.sink.CommitMessageSerializer;
 import org.apache.paimon.table.sink.TableCommitImpl;
 import org.apache.paimon.table.source.EndOfScanException;
+import org.apache.paimon.utils.ExceptionUtils;
+import org.apache.paimon.utils.RetryWaiter;
 
 import org.apache.spark.api.java.JavaRDD;
 import org.apache.spark.api.java.JavaSparkContext;
@@ -35,10 +38,14 @@ import org.apache.spark.sql.SparkSession;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import javax.annotation.Nullable;
+
 import java.io.IOException;
 import java.util.ArrayList;
+import java.util.Collections;
 import java.util.Iterator;
 import java.util.List;
+import java.util.Optional;
 import java.util.function.Function;
 
 import static org.apache.paimon.CoreOptions.createCommitUser;
@@ -59,6 +66,24 @@ final class DataEvolutionRewriteExecutor {
             JavaSparkContext javaSparkContext,
             SparkSession sparkSession,
             CommitConfigurer commitConfigurer) {
+        execute(
+                table,
+                initialSnapshot,
+                taskPlanner,
+                javaSparkContext,
+                sparkSession,
+                commitConfigurer,
+                null);
+    }
+
+    static void execute(
+            FileStoreTable table,
+            Snapshot initialSnapshot,
+            Function<Snapshot, List<DataEvolutionCompactTask>> taskPlanner,
+            JavaSparkContext javaSparkContext,
+            SparkSession sparkSession,
+            CommitConfigurer commitConfigurer,
+            @Nullable CommitMessageRewriter commitMessageRewriter) {
         CommitMessageSerializer messageSerializer = new 
CommitMessageSerializer();
         String commitUser = 
createCommitUser(table.coreOptions().toConfiguration());
         Snapshot preparationSnapshot = initialSnapshot;
@@ -130,22 +155,19 @@ final class DataEvolutionRewriteExecutor {
                                                 });
 
                 List<byte[]> serializedMessages = new 
ArrayList<>(commitMessageJavaRDD.collect());
-                try (TableCommitImpl commit = table.newCommit(commitUser)) {
-                    commitConfigurer.configure(commit);
-                    List<CommitMessage> messages =
+                try {
+                    List<CommitMessage> compactMessages =
                             deserializeCommitMessagesAndReleaseSerializedBytes(
                                     messageSerializer, serializedMessages);
-                    messages.addAll(
-                            new 
DataEvolutionCompactionCommitPreparation(table, preparationSnapshot)
-                                    .prepare(messages));
-                    commit.commit(messages);
                     Snapshot committedSnapshot =
-                            table.snapshotManager()
-                                    .latestSnapshotOfUser(commitUser)
-                                    .orElseThrow(
-                                            () ->
-                                                    new IllegalStateException(
-                                                            "Cannot find the 
committed data evolution rewrite snapshot."));
+                            commitWithMergeConflictRetry(
+                                    table,
+                                    preparationSnapshot,
+                                    compactMessages,
+                                    commitUser,
+                                    sparkSession,
+                                    commitConfigurer,
+                                    commitMessageRewriter);
                     checkArgument(
                             committedSnapshot.id() > preparationSnapshot.id(),
                             "Committed data evolution rewrite snapshot %s must 
be newer than preparation snapshot %s.",
@@ -165,6 +187,173 @@ final class DataEvolutionRewriteExecutor {
         }
     }
 
+    private static Snapshot commitWithMergeConflictRetry(
+            FileStoreTable table,
+            Snapshot taskSnapshot,
+            List<CommitMessage> compactMessages,
+            String commitUser,
+            SparkSession sparkSession,
+            CommitConfigurer commitConfigurer,
+            @Nullable CommitMessageRewriter commitMessageRewriter) {
+        int retryCount = 0;
+        long startMillis = System.currentTimeMillis();
+        RetryWaiter retryWaiter =
+                new RetryWaiter(
+                        table.coreOptions().commitMinRetryWait(),
+                        table.coreOptions().commitMaxRetryWait());
+        RuntimeException lastConflict = null;
+
+        while (true) {
+            if (lastConflict != null
+                    && System.currentTimeMillis() - startMillis
+                            > table.coreOptions().commitTimeout()) {
+                throw lastConflict;
+            }
+
+            Snapshot attemptSnapshot = taskSnapshot;
+            List<CommitMessage> attemptMessages = compactMessages;
+            List<CommitMessage> retryArtifacts = Collections.emptyList();
+            Snapshot latestSnapshot = table.snapshotManager().latestSnapshot();
+            if (commitMessageRewriter != null
+                    && latestSnapshot != null
+                    && latestSnapshot.id() > taskSnapshot.id()) {
+                Optional<List<CommitMessage>> rewritten;
+                try {
+                    rewritten =
+                            commitMessageRewriter.rewrite(
+                                    sparkSession, taskSnapshot, 
latestSnapshot, compactMessages);
+                } catch (RuntimeException rewriteError) {
+                    if (lastConflict == null) {
+                        throw rewriteError;
+                    }
+                    RuntimeException failure =
+                            new RuntimeException(
+                                    lastConflict.getMessage() + " " + 
rewriteError.getMessage(),
+                                    rewriteError);
+                    failure.addSuppressed(lastConflict);
+                    throw failure;
+                }
+                if (rewritten.isPresent()) {
+                    attemptSnapshot = latestSnapshot;
+                    attemptMessages = rewritten.get();
+                    retryArtifacts = retryArtifacts(compactMessages, 
attemptMessages);
+                    LOG.info(
+                            "Rebased staged data evolution compact files 
against compatible "
+                                    + "concurrent partial-column files "
+                                    + "through snapshot {} for table {}.",
+                            latestSnapshot.id(),
+                            table.fullName());
+                } else if (lastConflict != null) {
+                    throw lastConflict;
+                }
+                if (lastConflict != null
+                        && System.currentTimeMillis() - startMillis
+                                > table.coreOptions().commitTimeout()) {
+                    abortRetryArtifacts(table, commitUser, retryArtifacts, 
lastConflict);
+                    throw lastConflict;
+                }
+            }
+
+            List<CommitMessage> preparedMessages = new 
ArrayList<>(attemptMessages);
+            List<CommitMessage> preparationArtifacts =
+                    new DataEvolutionCompactionCommitPreparation(table, 
attemptSnapshot)
+                            .prepare(preparedMessages);
+            preparedMessages.addAll(preparationArtifacts);
+            List<CommitMessage> abortMessages = new 
ArrayList<>(retryArtifacts);
+            abortMessages.addAll(preparationArtifacts);
+            try (TableCommitImpl commit = table.newCommit(commitUser)) {
+                commitConfigurer.configure(commit);
+                try {
+                    commit.commit(preparedMessages);
+                } catch (RuntimeException conflict) {
+                    if (isMergeConflict(conflict)) {
+                        abortRetryArtifacts(commit, abortMessages, conflict, 
table);
+                    }
+                    throw conflict;
+                }
+                return table.snapshotManager()
+                        .latestSnapshotOfUser(commitUser)
+                        .orElseThrow(
+                                () ->
+                                        new IllegalStateException(
+                                                "Cannot find the committed 
data evolution rewrite snapshot."));
+            } catch (RuntimeException conflict) {
+                if (commitMessageRewriter == null
+                        || !isMergeConflict(conflict)
+                        || System.currentTimeMillis() - startMillis
+                                > table.coreOptions().commitTimeout()
+                        || retryCount >= 
table.coreOptions().commitMaxRetries()) {
+                    throw conflict;
+                }
+                lastConflict = conflict;
+                retryWaiter.retryWait(retryCount);
+                retryCount++;
+            } catch (Exception e) {
+                throw new RuntimeException(e);
+            }
+        }
+    }
+
+    private static List<CommitMessage> retryArtifacts(
+            List<CommitMessage> compactMessages, List<CommitMessage> 
rewrittenMessages) {
+        checkArgument(
+                rewrittenMessages.size() >= compactMessages.size(),
+                "Rewritten commit messages must retain all staged compact 
messages.");
+        for (int i = 0; i < compactMessages.size(); i++) {
+            checkArgument(
+                    rewrittenMessages.get(i) == compactMessages.get(i),
+                    "Rewritten commit messages must retain staged compact 
message %s.",
+                    i);
+        }
+        return new ArrayList<>(
+                rewrittenMessages.subList(compactMessages.size(), 
rewrittenMessages.size()));
+    }
+
+    private static boolean isMergeConflict(RuntimeException conflict) {
+        return ExceptionUtils.findThrowable(conflict, 
DataEvolutionRowRangeConflictException.class)
+                .isPresent();
+    }
+
+    private static void abortRetryArtifacts(
+            TableCommitImpl commit,
+            List<CommitMessage> abortMessages,
+            RuntimeException conflict,
+            FileStoreTable table) {
+        if (abortMessages.isEmpty()) {
+            return;
+        }
+        try {
+            commit.abort(abortMessages);
+        } catch (RuntimeException abortFailure) {
+            conflict.addSuppressed(abortFailure);
+            LOG.warn(
+                    "Failed to abort {} staged compact retry artifacts for 
table {}.",
+                    abortMessages.size(),
+                    table.fullName(),
+                    abortFailure);
+        }
+    }
+
+    private static void abortRetryArtifacts(
+            FileStoreTable table,
+            String commitUser,
+            List<CommitMessage> abortMessages,
+            RuntimeException conflict) {
+        if (abortMessages.isEmpty()) {
+            return;
+        }
+        try (TableCommitImpl commit = table.newCommit(commitUser)) {
+            abortRetryArtifacts(commit, abortMessages, conflict, table);
+        } catch (Exception abortFailure) {
+            conflict.addSuppressed(abortFailure);
+            LOG.warn(
+                    "Failed to close the commit after aborting staged compact 
retry artifacts "
+                            + "for table {}.",
+                    table.fullName(),
+                    abortFailure);
+        }
+    }
+
     private static List<CommitMessage> 
deserializeCommitMessagesAndReleaseSerializedBytes(
             CommitMessageSerializer serializer, List<byte[]> 
serializedMessages)
             throws IOException {
@@ -181,4 +370,18 @@ final class DataEvolutionRewriteExecutor {
 
         void configure(TableCommitImpl commit);
     }
+
+    @FunctionalInterface
+    interface CommitMessageRewriter {
+
+        /**
+         * Returns the original compact messages followed by any newly staged 
retry artifacts. Retry
+         * artifacts are aborted if the rebased commit still conflicts.
+         */
+        Optional<List<CommitMessage>> rewrite(
+                SparkSession sparkSession,
+                Snapshot taskSnapshot,
+                Snapshot latestSnapshot,
+                List<CommitMessage> compactMessages);
+    }
 }
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionCompactMergeConflictRewriter.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionCompactMergeConflictRewriter.scala
new file mode 100644
index 0000000000..f41cbfc3ae
--- /dev/null
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionCompactMergeConflictRewriter.scala
@@ -0,0 +1,458 @@
+/*
+ * 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.
+ */
+
+package org.apache.paimon.spark.commands
+
+import org.apache.paimon.Snapshot
+import org.apache.paimon.Snapshot.CommitKind
+import org.apache.paimon.data.BinaryRow
+import org.apache.paimon.format.blob.BlobFileFormat.isBlobFile
+import org.apache.paimon.io.{CompactIncrement, DataFileMeta, DataIncrement}
+import org.apache.paimon.manifest.FileSource
+import org.apache.paimon.spark.util.ScanPlanHelper
+import org.apache.paimon.table.{FileStoreTable, SpecialFields}
+import org.apache.paimon.table.sink.{CommitMessage, CommitMessageImpl}
+import org.apache.paimon.table.source.{DataSplit, IncrementalSplit}
+import org.apache.paimon.table.source.snapshot.SnapshotReader
+import org.apache.paimon.types.VectorType.isVectorStoreFile
+import org.apache.paimon.utils.{Range, RowRangeIndex}
+
+import org.apache.spark.sql.{functions, SparkSession}
+import org.apache.spark.sql.PaimonUtils.createDataset
+import org.apache.spark.sql.catalyst.analysis.SimpleAnalyzer.resolver
+import org.apache.spark.sql.catalyst.expressions.AttributeReference
+import org.apache.spark.sql.execution.datasources.v2.DataSourceV2Relation
+import org.apache.spark.sql.functions.{col, udf}
+import org.apache.spark.sql.paimon.shims.SparkShimLoader
+
+import java.util.{Collections, List => JList, Optional => JOptional}
+
+import scala.collection.JavaConverters._
+import scala.collection.mutable
+
+/** Rebases MERGE-compatible partial-column files onto staged compact output 
boundaries. */
+class DataEvolutionCompactMergeConflictRewriter(
+    table: FileStoreTable,
+    targetRelation: DataSourceV2Relation)
+  extends ScanPlanHelper {
+
+  import DataEvolutionCompactMergeConflictRewriter._
+
+  def rewrite(
+      sparkSession: SparkSession,
+      baseSnapshot: Snapshot,
+      latestSnapshot: Snapshot,
+      compactMessages: JList[CommitMessage]): JOptional[JList[CommitMessage]] 
= {
+    if (
+      table.coreOptions().deletionVectorsEnabled() ||
+      latestSnapshot.schemaId() != baseSnapshot.schemaId() ||
+      latestSnapshot.id() <= baseSnapshot.id()
+    ) {
+      return JOptional.empty()
+    }
+
+    val messageImpls = compactMessages.asScala.collect {
+      case message: CommitMessageImpl => message
+    }
+    if (messageImpls.size != compactMessages.size()) {
+      return JOptional.empty()
+    }
+
+    val targets = messageImpls
+      .flatMap(
+        message =>
+          normalRowIdFiles(message.compactIncrement().compactAfter().asScala)
+            .map(file => CompactTarget(message, file)))
+      .toSeq
+    if (targets.isEmpty) {
+      return JOptional.empty()
+    }
+    val targetIndex = new CompactTargetIndex(targets)
+    if (!targetIndex.valid) {
+      return JOptional.empty()
+    }
+    val targetScan = new TargetScan(targets)
+
+    // Snapshot.operation is optional and Python MERGE currently does not 
persist it. CommitKind
+    // and the portable partial-column file contract are available across 
engines.
+    val additions = mergeAdditions(baseSnapshot, latestSnapshot, targetScan, 
targetIndex) match {
+      case Some(files) => files
+      case None => return JOptional.empty()
+    }
+
+    val additionsByTarget =
+      mutable.HashMap.empty[CompactTarget, mutable.ArrayBuffer[AddedFile]]
+    additions.foreach {
+      addition =>
+        val intersectingTargets = targetIndex.intersecting(addition)
+        if (intersectingTargets.nonEmpty) {
+          if (
+            !isRegularPartialFile(addition.file) ||
+            intersectingTargets.length != 1 ||
+            !intersectingTargets.head.contains(addition)
+          ) {
+            return JOptional.empty()
+          }
+          additionsByTarget
+            .getOrElseUpdate(intersectingTargets.head, 
mutable.ArrayBuffer.empty)
+            .append(addition)
+        }
+    }
+    if (additionsByTarget.isEmpty) {
+      return JOptional.empty()
+    }
+
+    val targetRewrites = targets.flatMap {
+      target =>
+        val files = 
additionsByTarget.get(target).map(_.toSeq).getOrElse(Seq.empty)
+        if (files.nonEmpty) {
+          val updatedFields = table
+            .rowType()
+            .getFieldNames
+            .asScala
+            .filter(name => files.exists(_.file.writeCols().contains(name)))
+            .toSeq
+          if (updatedFields.isEmpty) {
+            return JOptional.empty()
+          }
+          Some(TargetRewrite(target, files.toSeq, updatedFields))
+        } else {
+          None
+        }
+    }
+    if (targetRewrites.isEmpty) {
+      return JOptional.empty()
+    }
+
+    val currentSplits = targetReader(latestSnapshot, targetScan)
+      .read()
+      .splits()
+      .asScala
+      .collect { case split: DataSplit => split }
+      .toSeq
+
+    val rewrittenMessages = targetRewrites
+      .groupBy(_.updatedFields)
+      .toSeq
+      .flatMap {
+        case (updatedFields, rewrites) =>
+          rewriteFiles(sparkSession, updatedFields, rewrites.toSeq, 
currentSplits)
+      }
+
+    JOptional.of((messageImpls ++ 
rewrittenMessages).map(_.asInstanceOf[CommitMessage]).asJava)
+  }
+
+  private def targetReader(snapshot: Snapshot, targetScan: TargetScan): 
SnapshotReader = {
+    table
+      .newSnapshotReader()
+      .withSnapshot(snapshot)
+      .withPartitionFilter(targetScan.partitions)
+      .withBucketFilter(bucket => targetScan.buckets.contains(bucket))
+      .withRowRangeIndex(targetScan.rowRangeIndex)
+  }
+
+  private def mergeAdditions(
+      baseSnapshot: Snapshot,
+      latestSnapshot: Snapshot,
+      targetScan: TargetScan,
+      targetIndex: CompactTargetIndex): Option[Seq[AddedFile]] = {
+    val additions = mutable.ArrayBuffer.empty[AddedFile]
+    val snapshotManager = table.snapshotManager()
+    var snapshotId = baseSnapshot.id() + 1
+    while (snapshotId <= latestSnapshot.id()) {
+      if (!snapshotManager.snapshotExists(snapshotId)) {
+        return None
+      }
+
+      val snapshot = snapshotManager.snapshot(snapshotId)
+      val changes = targetReader(snapshot, targetScan)
+        .readChanges()
+        .splits()
+        .asScala
+        .collect { case split: IncrementalSplit => split }
+      changes.foreach {
+        split =>
+          val before = split
+            .beforeFiles()
+            .asScala
+            .map(file => AddedFile(split.partition(), split.bucket(), file))
+          val after = split
+            .afterFiles()
+            .asScala
+            .map(file => AddedFile(split.partition(), split.bucket(), file))
+          if (snapshot.commitKind() != CommitKind.APPEND) {
+            if ((before ++ after).exists(file => 
targetIndex.intersecting(file).nonEmpty)) {
+              return None
+            }
+          } else {
+            if (before.exists(file => 
targetIndex.intersecting(file).nonEmpty)) {
+              return None
+            }
+            additions ++= after
+          }
+      }
+      snapshotId += 1
+    }
+    Some(additions.toSeq)
+  }
+
+  private def rewriteFiles(
+      sparkSession: SparkSession,
+      updatedFields: Seq[String],
+      rewrites: Seq[TargetRewrite],
+      currentSplits: Seq[DataSplit]): Seq[CommitMessageImpl] = {
+    val targetIndex = new CompactTargetIndex(rewrites.map(_.target))
+    val relevantSplits = currentSplits.flatMap {
+      split =>
+        val filtered = split.filterDataFile(
+          file =>
+            isNormalRowIdFile(file) &&
+              targetIndex.intersects(split.partition(), split.bucket(), 
file.nonNullRowIdRange()))
+        if (filtered.isPresent) Some(filtered.get()) else None
+    }
+
+    val relationAttributes = (targetRelation.output ++ 
targetRelation.metadataOutput).collect {
+      case attribute: AttributeReference => attribute
+    }
+    def attribute(name: String): AttributeReference = {
+      relationAttributes
+        .find(attr => resolver(attr.name, name))
+        .getOrElse(throw new RuntimeException(s"Cannot find column $name for 
compact rebase."))
+    }
+
+    val rowIdAttribute = attribute(ROW_ID_NAME)
+    val readOutput = updatedFields.map(attribute) :+ rowIdAttribute
+    val relation = createNewScanPlan(relevantSplits, targetRelation)
+    val readPlan =
+      SparkShimLoader.shim.copyDataSourceV2Relation(relation, relation.table, 
readOutput)
+    val targetRanges = rewrites.map(_.target.range).toArray
+    val rangeIndex = new CompactRowIdRangeIndex(targetRanges)
+    val firstRowId = udf((rowId: Long) => rangeIndex.firstRowId(rowId))
+    val rewrittenRows = createDataset(sparkSession, readPlan)
+      .select((updatedFields.map(quotedColumn) :+ quotedColumn(ROW_ID_NAME)): 
_*)
+      .withColumn(FIRST_ROW_ID_NAME, firstRowId(quotedColumn(ROW_ID_NAME)))
+      .filter(quotedColumn(FIRST_ROW_ID_NAME).isNotNull)
+      .repartition(col(FIRST_ROW_ID_NAME))
+      .sortWithinPartitions(FIRST_ROW_ID_NAME, ROW_ID_NAME)
+
+    val targetSplits = rewrites.map {
+      rewrite =>
+        val target = rewrite.target
+        DataSplit
+          .builder()
+          .withPartition(target.message.partition())
+          .withBucket(target.message.bucket())
+          .withTotalBuckets(target.message.totalBuckets())
+          .withBucketPath(
+            table
+              .store()
+              .pathFactory()
+              .bucketPath(target.message.partition(), target.message.bucket())
+              .toString)
+          .withDataFiles(Collections.singletonList(target.file))
+          .rawConvertible(true)
+          .build()
+    }
+
+    val written = DataEvolutionPaimonWriter(table, 
targetSplits).writePartialFields(
+      rewrittenRows,
+      updatedFields)
+    written.map {
+      case message: CommitMessageImpl =>
+        val newFiles = 
normalRowIdFiles(message.newFilesIncrement().newFiles().asScala)
+        if (newFiles.size != message.newFilesIncrement().newFiles().size()) {
+          throw new UnsupportedOperationException(
+            "Compact MERGE conflict rebase does not support dedicated files.")
+        }
+        val rewrite = rewrites
+          .find(
+            rewrite =>
+              rewrite.target.sameBucket(message.partition(), message.bucket()) 
&&
+                newFiles.forall(_.nonNullRowIdRange() == rewrite.target.range))
+          .getOrElse(throw new IllegalStateException(
+            s"Cannot match rebased files $newFiles to a staged compact 
range."))
+        // The rebased file only materializes the source state. Preserve its 
sequence range so a
+        // MERGE committed after this rewrite remains newer even if this 
COMPACT commits last.
+        val sourceFiles = rewrite.target.file +: rewrite.mergeFiles.map(_.file)
+        val minSequenceNumber = sourceFiles.map(_.minSequenceNumber()).min
+        val maxSequenceNumber = sourceFiles.map(_.maxSequenceNumber()).max
+        val rebasedFiles =
+          newFiles.map(_.assignSequenceNumber(minSequenceNumber, 
maxSequenceNumber))
+        new CommitMessageImpl(
+          message.partition(),
+          message.bucket(),
+          message.totalBuckets(),
+          DataIncrement.emptyIncrement(),
+          new CompactIncrement(
+            rewrite.mergeFiles.map(_.file).asJava,
+            rebasedFiles.asJava,
+            Collections.emptyList())
+        )
+      case other =>
+        throw new UnsupportedOperationException(
+          s"Unsupported compact MERGE conflict commit message: $other")
+    }
+  }
+
+}
+
+private object DataEvolutionCompactMergeConflictRewriter {
+
+  private val ROW_ID_NAME = "_ROW_ID"
+  private val FIRST_ROW_ID_NAME = "_FIRST_ROW_ID"
+
+  private case class AddedFile(partition: BinaryRow, bucket: Int, file: 
DataFileMeta)
+
+  private case class CompactTarget(message: CommitMessageImpl, file: 
DataFileMeta) {
+
+    val range: Range = file.nonNullRowIdRange()
+
+    def sameBucket(partition: BinaryRow, bucket: Int): Boolean = {
+      message.partition() == partition && message.bucket() == bucket
+    }
+
+    def contains(added: AddedFile): Boolean = {
+      sameBucket(added.partition, added.bucket) && 
containsRange(added.file.nonNullRowIdRange())
+    }
+
+    def containsRange(other: Range): Boolean = {
+      range.from <= other.from && other.to <= range.to
+    }
+  }
+
+  private case class TargetRewrite(
+      target: CompactTarget,
+      mergeFiles: Seq[AddedFile],
+      updatedFields: Seq[String])
+
+  private case class Bucket(partition: BinaryRow, bucket: Int)
+
+  private class TargetScan(targets: Seq[CompactTarget]) {
+
+    val partitions: JList[BinaryRow] =
+      targets.map(_.message.partition()).distinct.asJava
+    val buckets: Set[Int] = targets.map(_.message.bucket()).toSet
+    val rowRangeIndex: RowRangeIndex =
+      RowRangeIndex.create(targets.map(_.range).distinct.asJava)
+  }
+
+  private class CompactTargetIndex(targets: Seq[CompactTarget]) {
+
+    private val targetsByBucket = targets
+      .groupBy(target => Bucket(target.message.partition(), 
target.message.bucket()))
+      .map {
+        case (bucket, bucketTargets) =>
+          bucket -> bucketTargets.sortBy(_.range.from).toArray
+      }
+
+    val valid: Boolean = targetsByBucket.values.forall {
+      bucketTargets =>
+        bucketTargets.indices.drop(1).forall {
+          index => !bucketTargets(index - 
1).range.hasIntersection(bucketTargets(index).range)
+        }
+    }
+
+    def intersecting(added: AddedFile): Array[CompactTarget] = {
+      if (added.file.firstRowId() == null) {
+        Array.empty
+      } else {
+        intersecting(Bucket(added.partition, added.bucket), 
added.file.nonNullRowIdRange())
+      }
+    }
+
+    def intersects(partition: BinaryRow, bucket: Int, range: Range): Boolean = 
{
+      intersecting(Bucket(partition, bucket), range).nonEmpty
+    }
+
+    private def intersecting(bucket: Bucket, range: Range): 
Array[CompactTarget] = {
+      targetsByBucket.get(bucket) match {
+        case None => Array.empty
+        case Some(bucketTargets) =>
+          val first = firstPossible(bucketTargets, range)
+          val matches = mutable.ArrayBuffer.empty[CompactTarget]
+          var index = first
+          while (index < bucketTargets.length && 
bucketTargets(index).range.from <= range.to) {
+            matches.append(bucketTargets(index))
+            index += 1
+          }
+          matches.toArray
+      }
+    }
+
+    private def firstPossible(targets: Array[CompactTarget], range: Range): 
Int = {
+      var low = 0
+      var high = targets.length
+      while (low < high) {
+        val mid = (low + high) >>> 1
+        if (targets(mid).range.to < range.from) {
+          low = mid + 1
+        } else {
+          high = mid
+        }
+      }
+      low
+    }
+  }
+
+  private def normalRowIdFiles(files: Iterable[DataFileMeta]): 
Seq[DataFileMeta] = {
+    files.filter(isNormalRowIdFile).toSeq
+  }
+
+  private def isNormalRowIdFile(file: DataFileMeta): Boolean = {
+    file.firstRowId() != null && !isBlobFile(file.fileName()) && 
!isVectorStoreFile(file.fileName())
+  }
+
+  private def isRegularPartialFile(file: DataFileMeta): Boolean = {
+    isNormalRowIdFile(file) &&
+    file.fileSource().orElse(null) == FileSource.APPEND &&
+    file.writeCols() != null &&
+    !file.writeCols().isEmpty &&
+    file.writeCols().asScala.forall(column => 
!SpecialFields.isSystemField(column))
+  }
+
+  private def quotedColumn(name: String) = {
+    functions.col("`" + name.replace("`", "``") + "`")
+  }
+}
+
+private[spark] class CompactRowIdRangeIndex(inputRanges: Seq[Range]) extends 
Serializable {
+
+  private val ranges = inputRanges.sortBy(_.from).toArray
+  require(
+    ranges.indices.drop(1).forall(index => !ranges(index - 
1).hasIntersection(ranges(index))),
+    "Staged compact row ID ranges must not overlap.")
+
+  def firstRowId(rowId: Long): java.lang.Long = {
+    var low = 0
+    var high = ranges.length
+    while (low < high) {
+      val mid = (low + high) >>> 1
+      if (ranges(mid).from <= rowId) {
+        low = mid + 1
+      } else {
+        high = mid
+      }
+    }
+    val index = low - 1
+    if (index >= 0 && rowId <= ranges(index).to) {
+      java.lang.Long.valueOf(ranges(index).from)
+    } else {
+      null
+    }
+  }
+}
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
index 77dbe8cb66..9c5f464bbc 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
+++ 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
@@ -18,19 +18,31 @@
 
 package org.apache.paimon.spark.procedure
 
+import org.apache.paimon.Snapshot
 import org.apache.paimon.Snapshot.CommitKind
+import org.apache.paimon.append.dataevolution.{DataEvolutionCompactTask, 
DataEvolutionNormalCompactTask}
+import org.apache.paimon.data.BinaryRow
 import 
org.apache.paimon.deletionvectors.DeletionVectorsIndexFile.DELETION_VECTORS_INDEX
+import org.apache.paimon.format.blob.BlobFileFormat
 import org.apache.paimon.fs.Path
+import org.apache.paimon.io.DataFileMeta
+import org.apache.paimon.manifest.FileSource
+import 
org.apache.paimon.operation.commit.DataEvolutionRowRangeConflictException
 import org.apache.paimon.partition.PartitionPredicate
 import org.apache.paimon.spark.PaimonSparkTestBase
+import org.apache.paimon.spark.catalyst.analysis.PaimonRelation
+import 
org.apache.paimon.spark.commands.{DataEvolutionCompactMergeConflictRewriter, 
DataEvolutionPaimonWriter, PaimonSparkWriter}
+import org.apache.paimon.spark.commands.CompactRowIdRangeIndex
 import org.apache.paimon.spark.utils.SparkProcedureUtils
 import org.apache.paimon.table.FileStoreTable
-import org.apache.paimon.table.source.DataSplit
+import org.apache.paimon.table.source.{DataSplit, EndOfScanException, 
IncrementalSplit}
 import org.apache.paimon.table.source.snapshot.SnapshotReader
+import org.apache.paimon.utils.Range
 
 import org.apache.spark.api.java.JavaSparkContext
 import org.apache.spark.scheduler.{SparkListener, SparkListenerStageSubmitted}
 import org.apache.spark.sql.{Dataset, Row}
+import org.apache.spark.sql.functions.{col, udf}
 import org.apache.spark.sql.paimon.shims.memstream.MemoryStream
 import org.apache.spark.sql.streaming.StreamTest
 import org.assertj.core.api.Assertions
@@ -41,7 +53,7 @@ import java.lang.reflect.{InvocationHandler, Method, Proxy}
 import java.time.LocalDate
 import java.time.format.DateTimeFormatter
 import java.util
-import java.util.concurrent.atomic.AtomicBoolean
+import java.util.concurrent.atomic.{AtomicBoolean, AtomicInteger, AtomicLong, 
AtomicReference}
 
 import scala.collection.JavaConverters._
 import scala.util.Random
@@ -1771,6 +1783,560 @@ abstract class CompactProcedureTestBase extends 
PaimonSparkTestBase with StreamT
     }
   }
 
+  test("Paimon Procedure: rebase data evolution compact after operation-less 
partial update") {
+    withTable("T") {
+      sql("""
+            |CREATE TABLE T (id INT, value INT)
+            |TBLPROPERTIES (
+            |  'bucket' = '-1',
+            |  'file.format' = 'avro',
+            |  'row-tracking.enabled' = 'true',
+            |  'data-evolution.enabled' = 'true',
+            |  'compaction.min.file-num' = '2')
+            |""".stripMargin)
+      sql("INSERT INTO T VALUES (1, 10)")
+      sql("INSERT INTO T VALUES (2, 20)")
+
+      val table = loadTable("T")
+      val relation =
+        
PaimonRelation.getPaimonRelation(spark.table("T").queryExecution.analyzed)
+      val updated = new AtomicBoolean(false)
+
+      CompactProcedure.executeDataEvolutionCompaction(
+        table,
+        relation,
+        null,
+        null,
+        new JavaSparkContext(spark.sparkContext),
+        spark,
+        null,
+        _ => {
+          if (updated.compareAndSet(false, true)) {
+            // Python MERGE commits the same regular partial-column files 
without an operation.
+            val updateSnapshot = table.latestSnapshot().get()
+            val dataSplits = table
+              .newSnapshotReader()
+              .withSnapshot(updateSnapshot)
+              .read()
+              .splits()
+              .asScala
+              .collect { case split: DataSplit => split }
+              .toSeq
+            val firstRowIds = dataSplits
+              .flatMap(_.dataFiles().asScala)
+              .map(_.nonNullFirstRowId())
+              .sorted
+            val firstRowId = udf((rowId: Long) => firstRowIds.takeWhile(_ <= 
rowId).last)
+            val updateRows = sql("SELECT value + 1 AS value, _ROW_ID FROM T")
+              .withColumn("_FIRST_ROW_ID", firstRowId(col("_ROW_ID")))
+              .select("value", "_FIRST_ROW_ID", "_ROW_ID")
+            val updateMessages =
+              DataEvolutionPaimonWriter(table, dataSplits)
+                .writePartialFields(updateRows, Seq("value"))
+
+            val writer = PaimonSparkWriter(table)
+            writer.rowIdCheckConflict(updateSnapshot.id())
+            writer.commit(updateMessages)
+            assert(table.latestSnapshot().get().operation() == null)
+          }
+        }
+      )
+
+      checkAnswer(sql("SELECT id, value FROM T ORDER BY id"), Seq(Row(1, 11), 
Row(2, 21)))
+      val ranges = table
+        .newSnapshotReader()
+        .read()
+        .dataSplits()
+        .asScala
+        .flatMap(_.dataFiles().asScala)
+        .map(file => (file.nonNullFirstRowId(), file.rowCount()))
+        .distinct
+      assert(ranges == Seq((0L, 2L)), ranges)
+    }
+  }
+
+  test("Paimon Procedure: reject data evolution compact rebase after schema 
evolution") {
+    withTable("T") {
+      sql("""
+            |CREATE TABLE T (id INT, value INT)
+            |TBLPROPERTIES (
+            |  'bucket' = '-1',
+            |  'file.format' = 'avro',
+            |  'row-tracking.enabled' = 'true',
+            |  'data-evolution.enabled' = 'true',
+            |  'compaction.min.file-num' = '2')
+            |""".stripMargin)
+      sql("INSERT INTO T VALUES (1, 10)")
+      sql("INSERT INTO T VALUES (2, 20)")
+
+      val table = loadTable("T")
+      val relation =
+        
PaimonRelation.getPaimonRelation(spark.table("T").queryExecution.analyzed)
+      val updated = new AtomicBoolean(false)
+
+      val exception = intercept[RuntimeException] {
+        CompactProcedure.executeDataEvolutionCompaction(
+          table,
+          relation,
+          null,
+          null,
+          new JavaSparkContext(spark.sparkContext),
+          spark,
+          null,
+          _ => {
+            if (updated.compareAndSet(false, true)) {
+              sql("ALTER TABLE T ADD COLUMN extra INT")
+              val evolvedTable = loadTable("T")
+              val updateSnapshot = evolvedTable.latestSnapshot().get()
+              val dataSplits = evolvedTable
+                .newSnapshotReader()
+                .withSnapshot(updateSnapshot)
+                .read()
+                .splits()
+                .asScala
+                .collect { case split: DataSplit => split }
+                .toSeq
+              val firstRowIds = dataSplits
+                .flatMap(_.dataFiles().asScala)
+                .map(_.nonNullFirstRowId())
+                .sorted
+              val firstRowId = udf((rowId: Long) => firstRowIds.takeWhile(_ <= 
rowId).last)
+              val updateRows =
+                sql("SELECT value + 1 AS value, 99 AS extra, _ROW_ID FROM T 
WHERE id = 1")
+                  .withColumn("_FIRST_ROW_ID", firstRowId(col("_ROW_ID")))
+                  .select("value", "extra", "_FIRST_ROW_ID", "_ROW_ID")
+              val updateMessages =
+                DataEvolutionPaimonWriter(evolvedTable, dataSplits)
+                  .writePartialFields(updateRows, Seq("value", "extra"))
+
+              val writer = PaimonSparkWriter(evolvedTable)
+              writer.rowIdCheckConflict(updateSnapshot.id())
+              writer.commit(updateMessages)
+            }
+          }
+        )
+      }
+
+      Assertions
+        .assertThat(exception)
+        
.hasRootCauseInstanceOf(classOf[DataEvolutionRowRangeConflictException])
+      checkAnswer(
+        sql("SELECT id, value, extra FROM T ORDER BY id"),
+        Seq(Row(1, 11, 99), Row(2, 20, null)))
+    }
+  }
+
+  test("Paimon Procedure: retry rebased compact after same-boundary partial 
update") {
+    withTable("T") {
+      val table = createCompactMergeRaceTable()
+      val relation =
+        
PaimonRelation.getPaimonRelation(spark.table("T").queryExecution.analyzed)
+      val javaSparkContext = new JavaSparkContext(spark.sparkContext)
+      val attempts = new AtomicInteger()
+      val rewriteSnapshotId = new AtomicLong(-1L)
+      val mergeFileAfterRewrite = new AtomicReference[DataFileMeta]()
+      partialUpdate(table, "SELECT * FROM VALUES (1, 10), (2, 20) AS S(id, 
value)")
+      val normalFiles = normalDataFiles(table)
+      assert(normalFiles.size == 2)
+      val stagedTask = new DataEvolutionNormalCompactTask(
+        BinaryRow.EMPTY_ROW,
+        normalFiles.asJava
+      )
+      val taskSnapshot = table.latestSnapshot().get()
+      val planned = new AtomicBoolean(false)
+      val rewriter = new DataEvolutionCompactMergeConflictRewriter(table, 
relation)
+
+      val planner: java.util.function.Function[Snapshot, 
util.List[DataEvolutionCompactTask]] =
+        _ => {
+          if (planned.compareAndSet(false, true)) {
+            util.Collections.singletonList(stagedTask)
+          } else {
+            throw new EndOfScanException()
+          }
+        }
+      val configurer: DataEvolutionRewriteExecutor.CommitConfigurer =
+        _ => {
+          attempts.getAndIncrement() match {
+            case 0 =>
+              partialUpdate(
+                table,
+                "SELECT * FROM VALUES (1, 10), (2, 20) AS S(id, value)"
+              )
+              throw new DataEvolutionRowRangeConflictException("Injected MERGE 
range conflict.")
+            case 1 =>
+              // This MERGE lands after the retry file has been rewritten. The 
retry file must
+              // preserve its source sequence so this newer update remains the 
winner.
+              rewriteSnapshotId.set(table.latestSnapshot().get().id())
+              val mergeFile = partialUpdate(
+                table,
+                "SELECT * FROM VALUES (1, 11), (2, 21) AS S(id, value)"
+              )
+              mergeFileAfterRewrite.set(mergeFile)
+              assert(mergeFile.nonNullFirstRowId() == 0L)
+              assert(mergeFile.rowCount() == 2L)
+            case _ =>
+          }
+        }
+      val messageRewriter: DataEvolutionRewriteExecutor.CommitMessageRewriter =
+        (session, base, latest, messages) => rewriter.rewrite(session, base, 
latest, messages)
+
+      DataEvolutionRewriteExecutor.execute(
+        table,
+        taskSnapshot,
+        planner,
+        javaSparkContext,
+        spark,
+        configurer,
+        messageRewriter
+      )
+
+      Assertions.assertThat(attempts.get()).isEqualTo(2)
+      checkAnswer(sql("SELECT id, value FROM T ORDER BY id"), Seq(Row(1, 11), 
Row(2, 21)))
+
+      val mergeFile = mergeFileAfterRewrite.get()
+      val bridgeFiles = normalDataFiles(table).filter(
+        file =>
+          file.fileSource().orElse(null) == FileSource.APPEND &&
+            file.fileName() != mergeFile.fileName())
+      assert(bridgeFiles.size == 1, bridgeFiles)
+      val bridgeFile = bridgeFiles.head
+      assert(bridgeFile.maxSequenceNumber() == rewriteSnapshotId.get())
+      assert(bridgeFile.maxSequenceNumber() < mergeFile.maxSequenceNumber())
+      assert(mergeFile.maxSequenceNumber() < table.latestSnapshot().get().id())
+    }
+  }
+
+  test("Paimon Procedure: reject compact rebase over concurrent compact") {
+    withTable("T") {
+      val table = createCompactMergeRaceTable()
+      val relation =
+        
PaimonRelation.getPaimonRelation(spark.table("T").queryExecution.analyzed)
+      partialUpdate(table, "SELECT * FROM VALUES (1, 10), (2, 20) AS S(id, 
value)")
+      val normalFiles = normalDataFiles(table)
+      assert(normalFiles.size == 2)
+
+      val stagedTask = new DataEvolutionNormalCompactTask(
+        BinaryRow.EMPTY_ROW,
+        normalFiles.asJava
+      )
+      val stagedUser = "staged-compact"
+      val stagedMessage = stagedTask.doCompact(table, stagedUser)
+      val baseSnapshot = table.latestSnapshot().get()
+
+      val mergeFile = partialUpdate(
+        table,
+        "SELECT * FROM VALUES (1, 11), (2, 21) AS S(id, value)"
+      )
+      assert(mergeFile.fileSource().get() == FileSource.APPEND)
+
+      val concurrentTask = new DataEvolutionNormalCompactTask(
+        BinaryRow.EMPTY_ROW,
+        normalDataFiles(table).asJava
+      )
+      val concurrentUser = "concurrent-compact"
+      val concurrentMessage = concurrentTask.doCompact(table, concurrentUser)
+      val concurrentCommit = table.newCommit(concurrentUser)
+      try {
+        
concurrentCommit.commit(util.Collections.singletonList(concurrentMessage))
+      } finally {
+        concurrentCommit.close()
+      }
+
+      val latestSnapshot = table.latestSnapshot().get()
+      assert(latestSnapshot.commitKind() == CommitKind.COMPACT)
+      assert(concurrentTask.compactAfter().get(0).fileSource().get() == 
FileSource.COMPACT)
+
+      val rewritten = new DataEvolutionCompactMergeConflictRewriter(table, 
relation)
+        .rewrite(
+          spark,
+          baseSnapshot,
+          latestSnapshot,
+          util.Collections.singletonList(stagedMessage)
+        )
+      assert(!rewritten.isPresent)
+      checkAnswer(sql("SELECT id, value FROM T ORDER BY id"), Seq(Row(1, 11), 
Row(2, 21)))
+    }
+  }
+
+  test("Paimon Procedure: reject compact rebase when compact lands after 
rewrite") {
+    withTable("T") {
+      val table = createCompactMergeRaceTable()
+      val relation =
+        
PaimonRelation.getPaimonRelation(spark.table("T").queryExecution.analyzed)
+      val javaSparkContext = new JavaSparkContext(spark.sparkContext)
+      val attempts = new AtomicInteger()
+      partialUpdate(table, "SELECT * FROM VALUES (1, 10), (2, 20) AS S(id, 
value)")
+      val normalFiles = normalDataFiles(table)
+      assert(normalFiles.size == 2)
+      val stagedTask = new DataEvolutionNormalCompactTask(
+        BinaryRow.EMPTY_ROW,
+        normalFiles.asJava
+      )
+      val taskSnapshot = table.latestSnapshot().get()
+      val planned = new AtomicBoolean(false)
+      val rewriter = new DataEvolutionCompactMergeConflictRewriter(table, 
relation)
+
+      val planner: java.util.function.Function[Snapshot, 
util.List[DataEvolutionCompactTask]] =
+        _ => {
+          if (planned.compareAndSet(false, true)) {
+            util.Collections.singletonList(stagedTask)
+          } else {
+            throw new EndOfScanException()
+          }
+        }
+      val configurer: DataEvolutionRewriteExecutor.CommitConfigurer =
+        _ => {
+          attempts.getAndIncrement() match {
+            case 0 =>
+              partialUpdate(
+                table,
+                "SELECT * FROM VALUES (1, 10), (2, 20) AS S(id, value)"
+              )
+              throw new DataEvolutionRowRangeConflictException("Injected MERGE 
range conflict.")
+            case 1 =>
+              val concurrentTask = new DataEvolutionNormalCompactTask(
+                BinaryRow.EMPTY_ROW,
+                normalDataFiles(table).asJava
+              )
+              val compactUser = "concurrent-compact"
+              val compactMessage = concurrentTask.doCompact(table, compactUser)
+              val compactCommit = table.newCommit(compactUser)
+              try {
+                
compactCommit.commit(util.Collections.singletonList(compactMessage))
+              } finally {
+                compactCommit.close()
+              }
+              val compactedFile = concurrentTask.compactAfter().get(0)
+              assert(table.latestSnapshot().get().commitKind() == 
CommitKind.COMPACT)
+              assert(compactedFile.fileSource().get() == FileSource.COMPACT)
+              assert(compactedFile.writeCols().asScala == Seq("id", "value"))
+
+              val mergeFile = partialUpdate(
+                table,
+                "SELECT * FROM VALUES (1, 11), (2, 21) AS S(id, value)"
+              )
+              assert(mergeFile.fileSource().get() == FileSource.APPEND)
+            case _ =>
+          }
+        }
+      val messageRewriter: DataEvolutionRewriteExecutor.CommitMessageRewriter =
+        (session, base, latest, messages) => rewriter.rewrite(session, base, 
latest, messages)
+
+      assertThatThrownBy(
+        () =>
+          DataEvolutionRewriteExecutor.execute(
+            table,
+            taskSnapshot,
+            planner,
+            javaSparkContext,
+            spark,
+            configurer,
+            messageRewriter
+          )).hasMessageContaining("Execute data evolution rewrite failed")
+
+      Assertions.assertThat(attempts.get()).isEqualTo(2)
+      checkAnswer(sql("SELECT id, value FROM T ORDER BY id"), Seq(Row(1, 11), 
Row(2, 21)))
+    }
+  }
+
+  test("Paimon Procedure: compact rebase range lookup handles high 
cardinality") {
+    val ranges = (0 until 100000).map {
+      index =>
+        val firstRowId = index.toLong * 3
+        new Range(firstRowId, firstRowId + 1)
+    }
+    val rangeIndex = new CompactRowIdRangeIndex(ranges)
+
+    assert(rangeIndex.firstRowId(0) == Long.box(0))
+    assert(rangeIndex.firstRowId(150001) == Long.box(150000))
+    assert(rangeIndex.firstRowId(299998) == Long.box(299997))
+    assert(rangeIndex.firstRowId(2) == null)
+    assert(rangeIndex.firstRowId(300000) == null)
+  }
+
+  test("Paimon Procedure: rebase later compact planner batch after partial 
update") {
+    withTable("T") {
+      sql("""
+            |CREATE TABLE T (id INT, value INT, pt STRING)
+            |TBLPROPERTIES (
+            |  'bucket' = '-1',
+            |  'file.format' = 'avro',
+            |  'row-tracking.enabled' = 'true',
+            |  'data-evolution.enabled' = 'true',
+            |  'compaction.min.file-num' = '2',
+            |  'snapshot.num-retained.min' = '1',
+            |  'snapshot.num-retained.max' = '1',
+            |  'snapshot.time-retained' = '0 ms',
+            |  'commit.max-retries' = '2',
+            |  'commit.min-retry-wait' = '1 ms',
+            |  'commit.max-retry-wait' = '1 ms')
+            |PARTITIONED BY (pt)
+            |""".stripMargin)
+      sql("INSERT INTO T VALUES (1, 10, 'p0')")
+      sql("INSERT INTO T VALUES (2, 20, 'p0')")
+      sql("INSERT INTO T VALUES (3, 30, 'p1')")
+      sql("INSERT INTO T VALUES (4, 40, 'p1')")
+
+      val table = loadTable("T")
+      val relation =
+        
PaimonRelation.getPaimonRelation(spark.table("T").queryExecution.analyzed)
+      val updated = new AtomicBoolean(false)
+
+      CompactProcedure.executeDataEvolutionCompaction(
+        table,
+        relation,
+        null,
+        null,
+        new JavaSparkContext(spark.sparkContext),
+        spark,
+        Int.box(2),
+        _ => {
+          if (updated.compareAndSet(false, true)) {
+            val updateSnapshot = table.latestSnapshot().get()
+            val dataSplits = table
+              .newSnapshotReader()
+              .withSnapshot(updateSnapshot)
+              .read()
+              .splits()
+              .asScala
+              .collect { case split: DataSplit => split }
+              .toSeq
+            val stagedPartitions = dataSplits.filter {
+              split =>
+                val liveFiles = 
split.dataFiles().asScala.map(_.fileName()).toSet
+                val bucketPath = table
+                  .store()
+                  .pathFactory()
+                  .bucketPath(split.partition(), split.bucket())
+                table
+                  .fileIO()
+                  .listStatus(bucketPath)
+                  .filterNot(_.isDir)
+                  .map(_.getPath.getName)
+                  .exists(name => !liveFiles.contains(name))
+            }
+            Assertions.assertThat(stagedPartitions.size).isEqualTo(1)
+            val stagedPartition = 
stagedPartitions.head.partition().getString(0).toString
+            val updatePartition = if (stagedPartition == "p0") "p1" else "p0"
+            val updateId = if (updatePartition == "p0") 1 else 3
+            val firstRowIds = dataSplits
+              .flatMap(_.dataFiles().asScala)
+              .map(_.nonNullFirstRowId())
+              .sorted
+            val firstRowId = udf((rowId: Long) => firstRowIds.takeWhile(_ <= 
rowId).last)
+            val updateRows =
+              sql(s"SELECT value + 1 AS value, _ROW_ID FROM T WHERE id = 
$updateId")
+                .withColumn("_FIRST_ROW_ID", firstRowId(col("_ROW_ID")))
+                .select("value", "_FIRST_ROW_ID", "_ROW_ID")
+            val updateMessages =
+              DataEvolutionPaimonWriter(table, dataSplits)
+                .writePartialFields(updateRows, Seq("value"))
+
+            val writer = PaimonSparkWriter(table)
+            writer.rowIdCheckConflict(updateSnapshot.id())
+            writer.commit(updateMessages)
+          }
+        }
+      )
+
+      checkAnswer(sql("SELECT sum(value), count(*) FROM T"), Seq(Row(101L, 
4L)))
+    }
+  }
+
+  test("Paimon Procedure: abort failed compact rebase files before retry") {
+    withTable("T") {
+      sql("""
+            |CREATE TABLE T (id INT, value INT)
+            |TBLPROPERTIES (
+            |  'bucket' = '-1',
+            |  'file.format' = 'avro',
+            |  'row-tracking.enabled' = 'true',
+            |  'data-evolution.enabled' = 'true',
+            |  'compaction.min.file-num' = '2',
+            |  'commit.max-retries' = '3',
+            |  'commit.min-retry-wait' = '1 ms',
+            |  'commit.max-retry-wait' = '1 ms')
+            |""".stripMargin)
+      sql("INSERT INTO T VALUES (1, 10)")
+      sql("INSERT INTO T VALUES (2, 20)")
+
+      val table = loadTable("T")
+      val relation =
+        
PaimonRelation.getPaimonRelation(spark.table("T").queryExecution.analyzed)
+      val attempts = new AtomicInteger()
+      var originalStagedFiles = Set.empty[String]
+      var firstRetryFiles = Set.empty[String]
+
+      CompactProcedure.executeDataEvolutionCompaction(
+        table,
+        relation,
+        null,
+        null,
+        new JavaSparkContext(spark.sparkContext),
+        spark,
+        null,
+        _ => {
+          val attempt = attempts.getAndIncrement()
+          val updateSnapshot = table.latestSnapshot().get()
+          val dataSplits = table
+            .newSnapshotReader()
+            .withSnapshot(updateSnapshot)
+            .read()
+            .splits()
+            .asScala
+            .collect { case split: DataSplit => split }
+            .toSeq
+          val liveFiles =
+            dataSplits.flatMap(_.dataFiles().asScala).map(_.fileName()).toSet
+          val physicalFiles = dataSplits
+            .map(
+              split =>
+                table
+                  .store()
+                  .pathFactory()
+                  .bucketPath(split.partition(), split.bucket()))
+            .distinct
+            .flatMap(table.fileIO().listStatus)
+            .filterNot(_.isDir)
+            .map(_.getPath.getName)
+            .toSet
+          val stagedFiles = physicalFiles -- liveFiles
+          if (attempt == 0) {
+            originalStagedFiles = stagedFiles
+            assert(originalStagedFiles.nonEmpty)
+          } else if (attempt == 1) {
+            firstRetryFiles = stagedFiles -- originalStagedFiles
+            assert(firstRetryFiles.nonEmpty)
+          } else {
+            assert(firstRetryFiles.forall(file => 
!physicalFiles.contains(file)))
+          }
+
+          if (attempt < 2) {
+            val firstRowIds = dataSplits
+              .flatMap(_.dataFiles().asScala)
+              .map(_.nonNullFirstRowId())
+              .sorted
+            val firstRowId = udf((rowId: Long) => firstRowIds.takeWhile(_ <= 
rowId).last)
+            val updateRows = sql("SELECT value + 1 AS value, _ROW_ID FROM T 
WHERE id = 1")
+              .withColumn("_FIRST_ROW_ID", firstRowId(col("_ROW_ID")))
+              .select("value", "_FIRST_ROW_ID", "_ROW_ID")
+            val updateMessages =
+              DataEvolutionPaimonWriter(table, dataSplits)
+                .writePartialFields(updateRows, Seq("value"))
+
+            val writer = PaimonSparkWriter(table)
+            writer.rowIdCheckConflict(updateSnapshot.id())
+            writer.commit(updateMessages)
+          }
+        }
+      )
+
+      Assertions.assertThat(attempts.get()).isEqualTo(4)
+      Assertions.assertThat(firstRetryFiles.nonEmpty).isTrue
+      checkAnswer(sql("SELECT id, value FROM T ORDER BY id"), Seq(Row(1, 12), 
Row(2, 20)))
+    }
+  }
+
   test("Paimon Procedure: reject legacy row id rewrite for empty data 
evolution table") {
     withTable("T") {
       sql("""
@@ -1863,6 +2429,64 @@ abstract class CompactProcedureTestBase extends 
PaimonSparkTestBase with StreamT
     }
   }
 
+  private def createCompactMergeRaceTable(): FileStoreTable = {
+    sql("""
+          |CREATE TABLE T (id INT, value INT, picture BINARY)
+          |TBLPROPERTIES (
+          |  'bucket' = '-1',
+          |  'row-tracking.enabled' = 'true',
+          |  'data-evolution.enabled' = 'true',
+          |  'blob-field' = 'picture',
+          |  'compaction.min.file-num' = '2',
+          |  'commit.max-retries' = '3',
+          |  'commit.min-retry-wait' = '1 ms',
+          |  'commit.max-retry-wait' = '1 ms')
+          |""".stripMargin)
+    sql("""
+          |INSERT INTO T
+          |SELECT /*+ REPARTITION(1) */ id, value, CAST(NULL AS BINARY)
+          |FROM VALUES (1, 10), (2, 20) AS S(id, value)
+          |""".stripMargin)
+    loadTable("T")
+  }
+
+  private def partialUpdate(table: FileStoreTable, sourceQuery: String): 
DataFileMeta = {
+    val beforeMerge = table.latestSnapshot().get()
+    sql(sourceQuery).createOrReplaceTempView("merge_source")
+    try {
+      sql("""
+            |MERGE INTO T
+            |USING merge_source AS S
+            |ON T.id = S.id
+            |WHEN MATCHED THEN UPDATE SET T.id = S.id, T.value = S.value
+            |""".stripMargin)
+    } finally {
+      spark.catalog.dropTempView("merge_source")
+    }
+    table
+      .newSnapshotReader()
+      .withSnapshot(table.latestSnapshot().get())
+      .readIncrementalDiff(beforeMerge)
+      .splits()
+      .asScala
+      .collect { case split: IncrementalSplit => split }
+      .flatMap(_.afterFiles().asScala)
+      .find(file => !BlobFileFormat.isBlobFile(file.fileName()))
+      .get
+  }
+
+  private def normalDataFiles(table: FileStoreTable): Seq[DataFileMeta] = {
+    table
+      .newSnapshotReader()
+      .read()
+      .dataSplits()
+      .asScala
+      .flatMap(_.dataFiles().asScala)
+      .filterNot(file => BlobFileFormat.isBlobFile(file.fileName()))
+      .sortBy(_.maxSequenceNumber())
+      .toSeq
+  }
+
   private def executeDataEvolutionCompaction(
       table: FileStoreTable,
       candidateFilesPerBatch: Int): Unit = {

Reply via email to