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 = {