From 4699958a62713202d44f7522c29b1d497f0e92c4 Mon Sep 17 00:00:00 2001 From: zhongyujiang Date: Wed, 12 Aug 2026 20:39:11 +0800 Subject: [PATCH 1/5] [core][spark]: support rewrite_manifest procedure Add a rewrite_manifest Spark procedure that reads all manifest entries, sorts them globally by partition -> bucket -> level -> fileName (canceling ADD/DELETE pairs of the same file along the way) and writes them back as new manifest files so that ManifestFileMeta statistics become more compact for scan pruning. The sort runs distributed via a Spark sortByKey shuffle. - ManifestEntrySortKey: serializable sort key with byte[] partition for Kryo compatibility and lazy RecordComparator - FileStoreCommit#replaceManifest: optimistic-concurrency commit with conflict detection (containsAll) and concurrent delta preservation - RewriteManifestProcedure: distributed sortByKey + streaming ADD/DELETE cancellation + where clause + auto parallelism - FormatTableCommit: add unsupported replaceManifests - register rewrite_manifest in SparkProcedures - tests: ManifestEntrySortKeyTest (9), FileStoreCommitTest replaceManifest conflict scenarios (3), RewriteManifestProcedureTest (8) Co-Authored-By: Claude --- .../paimon/manifest/ManifestEntrySortKey.java | 143 ++++++ .../paimon/operation/FileStoreCommit.java | 11 + .../paimon/operation/FileStoreCommitImpl.java | 92 ++++ .../table/format/FormatTableCommit.java | 7 + .../paimon/table/sink/BatchTableCommit.java | 9 + .../paimon/table/sink/TableCommitImpl.java | 7 + .../manifest/ManifestEntrySortKeyTest.java | 173 +++++++ .../paimon/operation/FileStoreCommitTest.java | 92 ++++ .../apache/paimon/spark/SparkProcedures.java | 2 + .../procedure/RewriteManifestProcedure.java | 479 ++++++++++++++++++ .../RewriteManifestProcedureTest.scala | 363 +++++++++++++ 11 files changed, 1378 insertions(+) create mode 100644 paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntrySortKey.java create mode 100644 paimon-core/src/test/java/org/apache/paimon/manifest/ManifestEntrySortKeyTest.java create mode 100644 paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java create mode 100644 paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RewriteManifestProcedureTest.scala diff --git a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntrySortKey.java b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntrySortKey.java new file mode 100644 index 000000000000..6859e1113018 --- /dev/null +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntrySortKey.java @@ -0,0 +1,143 @@ +/* + * 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.manifest; + +import org.apache.paimon.codegen.CodeGenUtils; +import org.apache.paimon.codegen.RecordComparator; +import org.apache.paimon.data.BinaryRow; +import org.apache.paimon.types.DataType; +import org.apache.paimon.utils.SerializationUtils; + +import java.io.Serializable; +import java.util.List; +import java.util.Objects; + +/** + * A lightweight, serializable sort key for {@link ManifestEntry}. Entries are ordered by {@code + * partition -> bucket -> level -> fileName}, which co-locates the ADD and DELETE of the same file + * (they share the same key) so they always land in the same Spark partition after {@code sortByKey} + * and can be cancelled during the manifest rewrite. The {@code partition -> bucket -> level} prefix + * also keeps {@link ManifestFileMeta} statistics (partitionStats / minBucket / maxBucket / minLevel + * / maxLevel) compact for scan pruning. + * + *

The partition is kept as serialized bytes and deserialized lazily on the first comparison; the + * generated {@link RecordComparator} is also created lazily per executor to avoid being serialized + * across the shuffle. This class is intentionally small so that it can be used as the key of a + * Spark {@code sortByKey} shuffle without moving the whole {@link ManifestEntry} payload. + * + *

Kryo compatibility: the partition is stored as a plain {@code byte[]} rather than a + * {@link BinaryRow} because this key travels through Spark's shuffle, where the serializer is + * {@code KryoSerializer} (Paimon's Spark test base randomly picks Kryo). {@code BinaryRow} only + * implements Java serialization (its {@code BinarySection.writeObject/readObject} callbacks), and + * Kryo does not invoke those callbacks while also skipping the {@code transient segments} field — + * so a {@code BinaryRow} key ends up with {@code null} segments after a Kryo shuffle and NPEs on + * the first comparison. {@code byte[]} and {@code String} are transparent to both Kryo and Java + * serialization, and the {@link BinaryRow} / {@link RecordComparator} are rebuilt lazily on the + * executor after the shuffle. + */ +public class ManifestEntrySortKey implements Serializable, Comparable { + + private static final long serialVersionUID = 1L; + + private final byte[] partitionBytes; + private final int bucket; + private final int level; + private final String fileName; + + private final List partitionFieldTypes; + + private transient BinaryRow partition; + private transient RecordComparator partitionComparator; + + public ManifestEntrySortKey( + BinaryRow partition, + int bucket, + int level, + String fileName, + List partitionFieldTypes) { + this.partitionBytes = SerializationUtils.serializeBinaryRow(partition); + this.bucket = bucket; + this.level = level; + this.fileName = fileName; + this.partitionFieldTypes = partitionFieldTypes; + } + + @Override + public int compareTo(ManifestEntrySortKey other) { + // 1. partition + int cmp = partitionComparator().compare(partition(), other.partition()); + if (cmp != 0) { + return cmp; + } + // 2. bucket + cmp = Integer.compare(bucket, other.bucket); + if (cmp != 0) { + return cmp; + } + // 3. level + cmp = Integer.compare(level, other.level); + if (cmp != 0) { + return cmp; + } + // 4. fileName — co-locates ADD and DELETE of the same file (same key) + return fileName.compareTo(other.fileName); + } + + private BinaryRow partition() { + if (partition == null) { + partition = SerializationUtils.deserializeBinaryRow(partitionBytes); + } + return partition; + } + + private RecordComparator partitionComparator() { + if (partitionComparator == null) { + partitionComparator = CodeGenUtils.newRecordComparator(partitionFieldTypes); + } + return partitionComparator; + } + + public int bucket() { + return bucket; + } + + public int level() { + return level; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (o == null || getClass() != o.getClass()) { + return false; + } + ManifestEntrySortKey that = (ManifestEntrySortKey) o; + return bucket == that.bucket + && level == that.level + && Objects.deepEquals(partitionBytes, that.partitionBytes) + && Objects.equals(fileName, that.fileName); + } + + @Override + public int hashCode() { + return Objects.hash(bucket, level, fileName, java.util.Arrays.hashCode(partitionBytes)); + } +} diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommit.java b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommit.java index b039ffb9e9fc..7caa89e2da28 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommit.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommit.java @@ -23,6 +23,7 @@ import org.apache.paimon.disk.IOManager; import org.apache.paimon.fs.FileIO; import org.apache.paimon.manifest.ManifestCommittable; +import org.apache.paimon.manifest.ManifestFileMeta; import org.apache.paimon.operation.metrics.CommitMetrics; import org.apache.paimon.stats.Statistics; import org.apache.paimon.table.sink.CommitMessage; @@ -83,6 +84,16 @@ FileStoreCommit rowIdCheckConflictForMaterializeDvCompaction( /** Compact the manifest entries only. */ void compactManifest(); + /** + * Replace the manifest entries with the given rewritten manifests. The {@code removedManifests} + * are the manifests the caller read and sorted (used for conflict detection); the {@code + * addedManifests} are the sorted rewrite result produced by the caller. The commit reuses the + * optimistic concurrency mode of {@link #compactManifest()}: on conflict, new delta manifests + * added by other commits are appended to the tail of {@code addedManifests}. + */ + void replaceManifest( + List removedManifests, List addedManifests); + /** Roll back to the target snapshot and materialize it as the latest snapshot. */ boolean rollbackToAsLatest(Snapshot targetSnapshot); diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java index 505faba772e9..9307b690a4ae 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/FileStoreCommitImpl.java @@ -1578,6 +1578,98 @@ public void compactManifest() { } } + @Override + public void replaceManifest( + List removedManifests, List addedManifests) { + int retryCount = 0; + long startMillis = System.currentTimeMillis(); + Set removedPathSet = + removedManifests.stream() + .map(ManifestFileMeta::fileName) + .collect(Collectors.toSet()); + while (true) { + Snapshot latestSnapshot = snapshotManager.latestSnapshot(); + if (latestSnapshot == null) { + throw new RuntimeException("Cannot replace manifests: the table has no snapshot."); + } + + List currentBase = manifestList.readDataManifests(latestSnapshot); + Set currentPathSet = + currentBase.stream() + .map(ManifestFileMeta::fileName) + .collect(Collectors.toSet()); + + if (!currentPathSet.containsAll(removedPathSet)) { + cleanUpRewrittenManifests(addedManifests); + throw new RuntimeException( + "Manifest conflict: the current snapshot does not contain all the " + + "manifests to replace. Another manifest rewrite may have " + + "happened concurrently; please retry."); + } + + List manifestsKept = + currentBase.stream() + .filter(m -> !removedPathSet.contains(m.fileName())) + .collect(Collectors.toList()); + + List manifestsToCommit = new ArrayList<>(manifestsKept); + manifestsToCommit.addAll(addedManifests); + + Pair baseManifestList = manifestList.write(manifestsToCommit); + Pair deltaManifestList = manifestList.write(emptyList()); + + Snapshot newSnapshot = + new Snapshot( + latestSnapshot.id() + 1, + latestSnapshot.schemaId(), + baseManifestList.getLeft(), + baseManifestList.getRight(), + deltaManifestList.getLeft(), + deltaManifestList.getRight(), + null, + null, + latestSnapshot.indexManifest(), + commitUser, + Long.MAX_VALUE, + CommitKind.COMPACT, + System.currentTimeMillis(), + latestSnapshot.totalRecordCount(), + 0L, + null, + latestSnapshot.watermark(), + latestSnapshot.statistics(), + latestSnapshot.properties(), + latestSnapshot.nextRowId(), + null); + + if (commitSnapshotImpl(latestSnapshot, newSnapshot, emptyList())) { + return; + } + + manifestList.delete(deltaManifestList.getLeft()); + manifestList.delete(baseManifestList.getLeft()); + + if (System.currentTimeMillis() - startMillis > options.commitTimeout() + || retryCount >= options.commitMaxRetries()) { + cleanUpRewrittenManifests(addedManifests); + throw new RuntimeException( + String.format( + "Commit failed after %s millis with %s retries, there maybe exist commit conflicts between multiple jobs.", + options.commitTimeout(), retryCount)); + } + + retryWaiter.retryWait(retryCount); + retryCount++; + } + } + + /** Delete the rewritten manifest files produced by this rewrite. */ + private void cleanUpRewrittenManifests(List addedManifests) { + for (ManifestFileMeta manifest : addedManifests) { + manifestFile.delete(manifest.fileName()); + } + } + private boolean compactManifestOnce() { Snapshot latestSnapshot = snapshotManager.latestSnapshot(); diff --git a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableCommit.java b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableCommit.java index 278e4c4f369f..2b1c483cc75b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableCommit.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableCommit.java @@ -27,6 +27,7 @@ import org.apache.paimon.fs.FileStatus; import org.apache.paimon.fs.Path; import org.apache.paimon.fs.TwoPhaseOutputStream; +import org.apache.paimon.manifest.ManifestFileMeta; import org.apache.paimon.metrics.MetricRegistry; import org.apache.paimon.options.CatalogOptions; import org.apache.paimon.options.Options; @@ -318,6 +319,12 @@ public void compactManifests() { throw new UnsupportedOperationException(); } + @Override + public void replaceManifests( + List removedManifests, List addedManifests) { + throw new UnsupportedOperationException(); + } + @Override public TableCommit withMetricRegistry(MetricRegistry registry) { throw new UnsupportedOperationException(); diff --git a/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchTableCommit.java b/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchTableCommit.java index 2cfc181eada2..09a4c88d306b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchTableCommit.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/BatchTableCommit.java @@ -21,6 +21,7 @@ import org.apache.paimon.Snapshot; import org.apache.paimon.Snapshot.CommitKind; import org.apache.paimon.annotation.Public; +import org.apache.paimon.manifest.ManifestFileMeta; import org.apache.paimon.stats.Statistics; import java.util.List; @@ -73,6 +74,14 @@ public interface BatchTableCommit extends TableCommit { /** Compact the manifest entries. Generates a snapshot with {@link CommitKind#COMPACT}. */ void compactManifests(); + /** + * Replace the manifest entries with the given rewritten manifests. The {@code removedManifests} + * are the manifests the caller read and sorted; the {@code addedManifests} are the sorted + * rewrite result. Generates a snapshot with {@link CommitKind#COMPACT}. + */ + void replaceManifests( + List removedManifests, List addedManifests); + /** Set the logical operation type (e.g. WRITE, DELETE, MERGE) recorded in the snapshot. */ default BatchTableCommit withOperation(Snapshot.Operation operation) { return this; diff --git a/paimon-core/src/main/java/org/apache/paimon/table/sink/TableCommitImpl.java b/paimon-core/src/main/java/org/apache/paimon/table/sink/TableCommitImpl.java index 014b5e64daa1..a3a010c483b6 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/sink/TableCommitImpl.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/sink/TableCommitImpl.java @@ -28,6 +28,7 @@ import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.io.DataFilePathFactory; import org.apache.paimon.manifest.ManifestCommittable; +import org.apache.paimon.manifest.ManifestFileMeta; import org.apache.paimon.metrics.MetricRegistry; import org.apache.paimon.operation.FileStoreCommit; import org.apache.paimon.operation.PartitionExpire; @@ -235,6 +236,12 @@ public void compactManifests() { commit.compactManifest(); } + @Override + public void replaceManifests( + List removedManifests, List addedManifests) { + commit.replaceManifest(removedManifests, addedManifests); + } + public boolean rollbackToAsLatest(Tag targetTag) { checkCommitted(); boolean success = commit.rollbackToAsLatest(targetTag.trimToSnapshot()); diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestEntrySortKeyTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestEntrySortKeyTest.java new file mode 100644 index 000000000000..02030364c5e1 --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestEntrySortKeyTest.java @@ -0,0 +1,173 @@ +/* + * 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.manifest; + +import org.apache.paimon.data.BinaryRow; +import org.apache.paimon.data.BinaryRowWriter; +import org.apache.paimon.data.BinaryString; +import org.apache.paimon.types.DataType; +import org.apache.paimon.types.DataTypes; +import org.apache.paimon.types.RowType; + +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link ManifestEntrySortKey}. */ +public class ManifestEntrySortKeyTest { + + // partition type: (dt INT, region STRING) + private final List partitionFieldTypes = + RowType.of(DataTypes.INT(), DataTypes.STRING()).getFieldTypes(); + + private ManifestEntrySortKey key(int dt, String region, int bucket, int level) { + return key(dt, region, bucket, level, "f1"); + } + + private ManifestEntrySortKey key(int dt, String region, int bucket, int level, String file) { + return new ManifestEntrySortKey( + partition(dt, region), bucket, level, file, partitionFieldTypes); + } + + private BinaryRow partition(int dt, String region) { + BinaryRow row = new BinaryRow(2); + BinaryRowWriter writer = new BinaryRowWriter(row); + writer.writeInt(0, dt); + writer.writeString(1, BinaryString.fromString(region)); + writer.complete(); + return row; + } + + @Test + public void testOrderByPartitionFirst() { + // same bucket & level, different partition -> ordered by partition + assertThat(key(1, "a", 0, 0)).isLessThan(key(2, "a", 0, 0)); + assertThat(key(2, "a", 0, 0)).isGreaterThan(key(1, "a", 0, 0)); + } + + @Test + public void testPartitionComparesAllFields() { + // dt equal, region differs -> ordered by region (second partition field) + assertThat(key(1, "a", 0, 0)).isLessThan(key(1, "b", 0, 0)); + assertThat(key(1, "b", 0, 0)).isGreaterThan(key(1, "a", 0, 0)); + // dt differs -> dt wins regardless of region + assertThat(key(2, "a", 0, 0)).isGreaterThan(key(1, "z", 0, 0)); + } + + @Test + public void testOrderByBucketWhenPartitionEquals() { + assertThat(key(1, "a", 0, 0)).isLessThan(key(1, "a", 1, 0)); + assertThat(key(1, "a", 1, 0)).isGreaterThan(key(1, "a", 0, 0)); + } + + @Test + public void testOrderByLevelWhenPartitionAndBucketEqual() { + assertThat(key(1, "a", 0, 0)).isLessThan(key(1, "a", 0, 1)); + assertThat(key(1, "a", 0, 1)).isGreaterThan(key(1, "a", 0, 0)); + } + + @Test + public void testOrderByFileNameWhenPartitionBucketLevelEqual() { + // same (p,b,l), different fileName -> ordered by fileName + assertThat(key(1, "a", 0, 0, "f1")).isLessThan(key(1, "a", 0, 0, "f2")); + assertThat(key(1, "a", 0, 0, "f2")).isGreaterThan(key(1, "a", 0, 0, "f1")); + } + + @Test + public void testAddAndDeleteShareSameKey() { + // ADD and DELETE of the same file share the same sort key, so they always land in the + // same Spark partition after sortByKey and can be cancelled. + ManifestEntrySortKey add = key(1, "a", 0, 0, "f1"); + ManifestEntrySortKey delete = key(1, "a", 0, 0, "f1"); + assertThat(add.compareTo(delete)).isEqualTo(0); + assertThat(add).isEqualTo(delete); + } + + @Test + public void testGlobalSort() { + // build keys in a deliberately shuffled order + List keys = + new ArrayList<>( + Arrays.asList( + key(2, "a", 1, 0, "f1"), + key(1, "b", 0, 1, "f1"), + key(1, "a", 0, 0, "f2"), + key(1, "a", 0, 0, "f1"), + key(2, "a", 0, 0, "f2"), + key(1, "a", 0, 1, "f1"), + key(2, "a", 0, 0, "f1"), + key(1, "a", 1, 0, "f1"))); + Collections.sort(keys); + + // expected: partition asc -> bucket asc -> level asc -> fileName asc + List expected = + Arrays.asList( + key(1, "a", 0, 0, "f1"), + key(1, "a", 0, 0, "f2"), + key(1, "a", 0, 1, "f1"), + key(1, "a", 1, 0, "f1"), + key(1, "b", 0, 1, "f1"), + key(2, "a", 0, 0, "f1"), + key(2, "a", 0, 0, "f2"), + key(2, "a", 1, 0, "f1")); + assertThat(keys).isEqualTo(expected); + } + + @Test + public void testEqualsAndHashCode() { + // keys with the same partition/bucket/level/fileName are equal even if built from different + // BinaryRow instances, because the partition is compared by serialized bytes + ManifestEntrySortKey a = key(1, "a", 3, 5, "f1"); + ManifestEntrySortKey b = key(1, "a", 3, 5, "f1"); + assertThat(a).isEqualTo(b); + assertThat(a.hashCode()).isEqualTo(b.hashCode()); + + // any field difference breaks equality + assertThat(a).isNotEqualTo(key(2, "a", 3, 5, "f1")); + assertThat(a).isNotEqualTo(key(1, "b", 3, 5, "f1")); + assertThat(a).isNotEqualTo(key(1, "a", 4, 5, "f1")); + assertThat(a).isNotEqualTo(key(1, "a", 3, 6, "f1")); + assertThat(a).isNotEqualTo(key(1, "a", 3, 5, "f2")); + } + + @Test + public void testSurvivesSerialization() throws Exception { + // the comparator and partition are transient; after Java serialization they must be + // rebuilt lazily and comparison must still work + ManifestEntrySortKey original = key(1, "a", 0, 0, "f1"); + + java.io.ByteArrayOutputStream baos = new java.io.ByteArrayOutputStream(); + try (java.io.ObjectOutputStream oos = new java.io.ObjectOutputStream(baos)) { + oos.writeObject(original); + } + try (java.io.ObjectInputStream ois = + new java.io.ObjectInputStream( + new java.io.ByteArrayInputStream(baos.toByteArray()))) { + ManifestEntrySortKey roundTripped = (ManifestEntrySortKey) ois.readObject(); + assertThat(roundTripped.compareTo(key(1, "a", 0, 0, "f1"))).isEqualTo(0); + assertThat(roundTripped).isLessThan(key(2, "a", 0, 0, "f1")); + assertThat(roundTripped).isLessThan(key(1, "a", 0, 0, "f2")); + } + } +} diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java index 3297383dcbd7..46877ae17c24 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java @@ -99,6 +99,7 @@ import java.util.Map; import java.util.Objects; import java.util.Optional; +import java.util.Set; import java.util.UUID; import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.atomic.AtomicReference; @@ -2388,4 +2389,95 @@ private void logData(Supplier> supplier, String name) { } LOG.debug("========== End of " + name + " =========="); } + + // --------------------------------------------------------------------------------------------- + // Tests for replaceManifest + // --------------------------------------------------------------------------------------------- + + @Test + public void testReplaceManifestsIdentity() throws Exception { + TestFileStore store = createStore(false, 2); + store.commitData(generateDataList(20), gen::getPartition, kv -> 0); + + Snapshot snapshot = store.snapshotManager().latestSnapshot(); + ManifestList manifestList = store.manifestListFactory().create(); + List manifests = manifestList.readDataManifests(snapshot); + assertThat(manifests).isNotEmpty(); + + Set fileNamesBefore = + manifests.stream().map(ManifestFileMeta::fileName).collect(Collectors.toSet()); + + try (FileStoreCommit commit = store.newCommit()) { + commit.replaceManifest(manifests, manifests); + } + + Snapshot newSnapshot = store.snapshotManager().latestSnapshot(); + assertThat(newSnapshot.id()).isEqualTo(snapshot.id() + 1); + assertThat(newSnapshot.commitKind()).isEqualTo(Snapshot.CommitKind.COMPACT); + List newManifests = manifestList.readDataManifests(newSnapshot); + Set fileNamesAfter = + newManifests.stream().map(ManifestFileMeta::fileName).collect(Collectors.toSet()); + assertThat(fileNamesAfter).isEqualTo(fileNamesBefore); + } + + @Test + public void testReplaceManifestsPreservesConcurrentDelta() throws Exception { + TestFileStore store = createStore(false, 2); + store.commitData(generateDataList(20), gen::getPartition, kv -> 0); + + Snapshot snapshotBefore = store.snapshotManager().latestSnapshot(); + ManifestList manifestList = store.manifestListFactory().create(); + List manifestsBefore = manifestList.readDataManifests(snapshotBefore); + + store.commitData(generateDataList(20), gen::getPartition, kv -> 0); + Snapshot snapshotAfter = store.snapshotManager().latestSnapshot(); + List manifestsAfter = manifestList.readDataManifests(snapshotAfter); + + Set beforeNames = + manifestsBefore.stream() + .map(ManifestFileMeta::fileName) + .collect(Collectors.toSet()); + List deltaManifests = + manifestsAfter.stream() + .filter(m -> !beforeNames.contains(m.fileName())) + .collect(Collectors.toList()); + assertThat(deltaManifests).isNotEmpty(); + + try (FileStoreCommit commit = store.newCommit()) { + commit.replaceManifest(manifestsBefore, manifestsBefore); + } + + Snapshot finalSnapshot = store.snapshotManager().latestSnapshot(); + List finalManifests = manifestList.readDataManifests(finalSnapshot); + Set finalNames = + finalManifests.stream().map(ManifestFileMeta::fileName).collect(Collectors.toSet()); + + for (String name : beforeNames) { + assertThat(finalNames).contains(name); + } + for (ManifestFileMeta delta : deltaManifests) { + assertThat(finalNames).contains(delta.fileName()); + } + } + + @Test + public void testReplaceManifestsThrowsOnConcurrentCompact() throws Exception { + TestFileStore store = createStore(false, 2); + for (int i = 0; i < 3; i++) { + store.commitData(generateDataList(10), gen::getPartition, kv -> 0); + } + + Snapshot snapshotBefore = store.snapshotManager().latestSnapshot(); + ManifestList manifestList = store.manifestListFactory().create(); + List manifestsBefore = manifestList.readDataManifests(snapshotBefore); + + try (FileStoreCommit commit = store.newCommit()) { + commit.compactManifest(); + } + + try (FileStoreCommit commit = store.newCommit()) { + assertThatThrownBy(() -> commit.replaceManifest(manifestsBefore, manifestsBefore)) + .hasMessageContaining("Manifest conflict"); + } + } } diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkProcedures.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkProcedures.java index 60b5747e3db7..facdbec37753 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkProcedures.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkProcedures.java @@ -58,6 +58,7 @@ import org.apache.paimon.spark.procedure.RescaleProcedure; import org.apache.paimon.spark.procedure.ResetConsumerProcedure; import org.apache.paimon.spark.procedure.RewriteFileIndexProcedure; +import org.apache.paimon.spark.procedure.RewriteManifestProcedure; import org.apache.paimon.spark.procedure.RollbackProcedure; import org.apache.paimon.spark.procedure.RollbackToTimestampProcedure; import org.apache.paimon.spark.procedure.RollbackToWatermarkProcedure; @@ -124,6 +125,7 @@ private static Map> initProcedureBuilders() { procedureBuilders.put("reset_consumer", ResetConsumerProcedure::builder); procedureBuilders.put("mark_partition_done", MarkPartitionDoneProcedure::builder); procedureBuilders.put("compact_manifest", CompactManifestProcedure::builder); + procedureBuilders.put("rewrite_manifest", RewriteManifestProcedure::builder); procedureBuilders.put("clear_consumers", ClearConsumersProcedure::builder); procedureBuilders.put("alter_view_dialect", AlterViewDialectProcedure::builder); procedureBuilders.put("create_function", CreateFunctionProcedure::builder); diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java new file mode 100644 index 000000000000..8f1b32e8add4 --- /dev/null +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java @@ -0,0 +1,479 @@ +/* + * 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.procedure; + +import org.apache.paimon.Snapshot; +import org.apache.paimon.fs.Path; +import org.apache.paimon.manifest.FileKind; +import org.apache.paimon.manifest.ManifestEntry; +import org.apache.paimon.manifest.ManifestEntrySerializer; +import org.apache.paimon.manifest.ManifestEntrySortKey; +import org.apache.paimon.manifest.ManifestFile; +import org.apache.paimon.manifest.ManifestFile.ManifestEntryWriter; +import org.apache.paimon.manifest.ManifestFileMeta; +import org.apache.paimon.manifest.ManifestFileMetaSerializer; +import org.apache.paimon.manifest.ManifestList; +import org.apache.paimon.partition.PartitionPredicate; +import org.apache.paimon.predicate.Predicate; +import org.apache.paimon.predicate.PredicateBuilder; +import org.apache.paimon.spark.catalyst.analysis.expressions.ExpressionUtils; +import org.apache.paimon.stats.SimpleStats; +import org.apache.paimon.table.FileStoreTable; +import org.apache.paimon.table.Table; +import org.apache.paimon.table.sink.BatchTableCommit; +import org.apache.paimon.types.DataType; +import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.FileStorePathFactory; +import org.apache.paimon.utils.ProcedureUtils; +import org.apache.paimon.utils.StringUtils; + +import org.apache.spark.api.java.JavaPairRDD; +import org.apache.spark.api.java.JavaSparkContext; +import org.apache.spark.api.java.function.FlatMapFunction; +import org.apache.spark.api.java.function.PairFlatMapFunction; +import org.apache.spark.sql.catalyst.InternalRow; +import org.apache.spark.sql.catalyst.expressions.Expression; +import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan; +import org.apache.spark.sql.connector.catalog.Identifier; +import org.apache.spark.sql.connector.catalog.TableCatalog; +import org.apache.spark.sql.execution.datasources.v2.DataSourceV2Relation; +import org.apache.spark.sql.types.DataTypes; +import org.apache.spark.sql.types.Metadata; +import org.apache.spark.sql.types.StructField; +import org.apache.spark.sql.types.StructType; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.Iterator; +import java.util.List; +import java.util.stream.Collectors; + +import scala.Tuple2; + +import static org.apache.paimon.utils.Preconditions.checkArgument; +import static org.apache.spark.sql.types.DataTypes.StringType; + +/** + * Rewrite manifest procedure. It reads all manifest entries, sorts them globally by {@code + * partition -> bucket -> level -> fileName} (canceling ADD/DELETE pairs along the way) and writes + * them back as new manifest files so that {@link ManifestFileMeta} statistics become more compact + * for scan pruning. The sort runs distributed via a Spark {@code sortByKey} shuffle. + * + *

An optional {@code where} clause restricts the rewrite to manifests whose partition stats may + * match the predicate; the remaining manifests are left untouched. + * + *


+ *  CALL sys.rewrite_manifest(table => 'tableId')
+ *  CALL sys.rewrite_manifest(table => 'tableId', where => 'dt = "2024-01-01"')
+ * 
+ */ +public class RewriteManifestProcedure extends BaseProcedure { + + private static final ProcedureParameter[] PARAMETERS = + new ProcedureParameter[] { + ProcedureParameter.required("table", StringType), + ProcedureParameter.optional("where", StringType), + ProcedureParameter.optional("options", StringType) + }; + + private static final StructType OUTPUT_TYPE = + new StructType( + new StructField[] { + new StructField( + "rewritten_manifests_count", + DataTypes.IntegerType, + true, + Metadata.empty()), + new StructField( + "added_manifests_count", + DataTypes.IntegerType, + true, + Metadata.empty()) + }); + + protected RewriteManifestProcedure(TableCatalog tableCatalog) { + super(tableCatalog); + } + + @Override + public ProcedureParameter[] parameters() { + return PARAMETERS; + } + + @Override + public StructType outputType() { + return OUTPUT_TYPE; + } + + @Override + public InternalRow[] call(InternalRow args) { + Identifier tableIdent = toIdentifier(args.getString(0), PARAMETERS[0].name()); + String where = args.isNullAt(1) ? null : args.getString(1); + String options = args.isNullAt(2) ? null : args.getString(2); + + Table table = loadSparkTable(tableIdent).getTable(); + HashMap dynamicOptions = new HashMap<>(); + ProcedureUtils.putAllOptions(dynamicOptions, options); + FileStoreTable fileStoreTable = (FileStoreTable) table.copy(dynamicOptions); + + // 1. read the latest snapshot and its data manifests + Snapshot latestSnapshot = fileStoreTable.store().snapshotManager().latestSnapshot(); + if (latestSnapshot == null) { + return new InternalRow[] {newInternalRow(0, 0)}; + } + ManifestList manifestList = fileStoreTable.store().manifestListFactory().create(); + List currentManifests = manifestList.readDataManifests(latestSnapshot); + if (currentManifests.isEmpty()) { + return new InternalRow[] {newInternalRow(0, 0)}; + } + + // 2. filter manifests by the optional where clause (partition-stats pruning) + PartitionPredicate partitionPredicate = resolvePartitionPredicate(tableIdent, table, where); + if (partitionPredicate != null) { + currentManifests = filterManifests(currentManifests, partitionPredicate); + if (currentManifests.isEmpty()) { + return new InternalRow[] {newInternalRow(0, 0)}; + } + } + + // 3. globally sort manifest entries and write them back as new manifest files + List newManifests = rewriteManifests(fileStoreTable, currentManifests); + + // 4. commit the rewritten manifests (optimistic concurrency with retry) + try (BatchTableCommit commit = fileStoreTable.newBatchWriteBuilder().newCommit()) { + commit.replaceManifests(currentManifests, newManifests); + } catch (Exception e) { + throw new RuntimeException(e); + } + + // rewritten_manifests_count: total manifests produced by the rewrite + // added_manifests_count: net change in manifest count vs before rewrite (may be negative + // when the rewrite also consolidates small manifest files) + int rewrittenCount = newManifests.size(); + int addedCount = newManifests.size() - currentManifests.size(); + return new InternalRow[] {newInternalRow(rewrittenCount, addedCount)}; + } + + /** + * Parse the {@code where} SQL string into a {@link PartitionPredicate}, validating that it only + * references partition columns. Returns {@code null} when {@code where} is blank. + */ + private PartitionPredicate resolvePartitionPredicate( + Identifier tableIdent, Table table, String where) { + if (StringUtils.isNullOrWhitespaceOnly(where)) { + return null; + } + FileStoreTable fileStoreTable = (FileStoreTable) table; + DataSourceV2Relation relation = createRelation(tableIdent); + Expression condition = ExpressionUtils.resolveFilter(spark(), relation, where); + checkArgument( + ExpressionUtils.isValidPredicate( + spark(), condition, fileStoreTable.partitionKeys().toArray(new String[0])), + "Only partition predicate is supported, your predicate is %s, but partition keys are %s", + condition, + fileStoreTable.partitionKeys()); + Predicate predicate = + ExpressionUtils.convertConditionToPaimonPredicate( + condition, + ((LogicalPlan) relation).output(), + table.rowType(), + false) + .getOrElse(null); + + // the predicate references fields by their index in the full row type; map them to the + // partition type index so PartitionPredicate can test against partitionStats / BinaryRow + // partitions (which only contain partition columns) + List partitionKeys = fileStoreTable.partitionKeys(); + int[] fieldIdxToPartitionIdx = + fileStoreTable.schema().fields().stream() + .mapToInt(f -> partitionKeys.indexOf(f.name())) + .toArray(); + Predicate partitionPredicate = + PredicateBuilder.transformFieldMapping(predicate, fieldIdxToPartitionIdx) + .orElse(null); + + RowType partitionType = fileStoreTable.store().partitionType(); + return PartitionPredicate.fromPredicate(partitionType, partitionPredicate); + } + + /** Keep only manifests whose partition stats may match the predicate. */ + private List filterManifests( + List manifests, PartitionPredicate predicate) { + return manifests.stream() + .filter( + m -> { + SimpleStats stats = m.partitionStats(); + return predicate.test( + m.numAddedFiles() + m.numDeletedFiles(), + stats.minValues(), + stats.maxValues(), + stats.nullCounts()); + }) + .collect(Collectors.toList()); + } + + private List rewriteManifests( + FileStoreTable table, List currentManifests) { + List partitionFieldTypes = table.store().partitionType().getFieldTypes(); + ManifestFileMetaSerializer metaSerializer = new ManifestFileMetaSerializer(); + + // serialize ManifestFileMeta to byte[] so they can travel through the RDD + List serializedMetas = new ArrayList<>(currentManifests.size()); + for (ManifestFileMeta meta : currentManifests) { + try { + serializedMetas.add(metaSerializer.serializeToBytes(meta)); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + + int numPartitions = computeNumPartitions(table, currentManifests); + + JavaSparkContext jsc = new JavaSparkContext(spark().sparkContext()); + + // The table is captured by the closure (not broadcast) so that Spark serializes it with + // Java serialization (ClosureCleaner always uses JavaSerializer). This is required because + // the table's FileIO (e.g. HadoopFileIO) holds a SerializableConfiguration whose transient + // Configuration is restored only via Java writeObject/readObject — Kryo (used by broadcast + // when spark.serializer=KryoSerializer) would skip those callbacks and leave a null conf. + FileStoreTable tableRef = table; + + // Step 1: read all entries (ADD and DELETE) and pair them with their sort key + JavaPairRDD entryRDD = + jsc.parallelize(serializedMetas, Math.min(serializedMetas.size(), numPartitions)) + .flatMapToPair( + (PairFlatMapFunction) + manifestBytes -> { + ManifestFileMetaSerializer ser = + new ManifestFileMetaSerializer(); + ManifestFileMeta manifest = + ser.deserializeFromBytes(manifestBytes); + ManifestFile manifestFile = + tableRef.store().manifestFileFactory().create(); + ManifestEntrySerializer entrySer = + new ManifestEntrySerializer(); + + List> pairs = + new ArrayList<>(); + for (ManifestEntry entry : + manifestFile.read( + manifest.fileName(), + manifest.fileSize())) { + ManifestEntrySortKey key = + new ManifestEntrySortKey( + entry.partition(), + entry.bucket(), + entry.file().level(), + entry.fileName(), + partitionFieldTypes); + pairs.add( + new Tuple2<>( + key, + entrySer.serializeToBytes(entry))); + } + return pairs.iterator(); + }); + + // Step 2: global sort. The sort key is (partition, bucket, level, fileName), so + // the ADD and DELETE of the same file are adjacent. + JavaPairRDD sortedRDD = + entryRDD.sortByKey(true, numPartitions); + + // Step 3: within each partition, stream the sorted entries, cancel ADD/DELETE pairs per + // manifest entry, and write surviving entries to a single manifest file. Because the sort + // key + // is (partition, bucket, level, fileName), entries of the same file are consecutive; the + // per-key buffer holds at most one ADD and one DELETE, so memory is negligible. Each task + // produces at most one manifest, sized roughly to the manifest target file size. + List serializedResult = + sortedRDD + .mapPartitions( + (FlatMapFunction< + Iterator>, + byte[]>) + iter -> { + FileStoreTable t = tableRef; + FileStorePathFactory pathFactory = + t.store().pathFactory(); + Path manifestPath = pathFactory.newManifestFile(); + ManifestFile manifestFile = + t.store().manifestFileFactory().create(); + ManifestEntryWriter writer = + manifestFile.createManifestEntryWriter( + manifestPath); + ManifestEntrySerializer entrySer = + new ManifestEntrySerializer(); + + // per-current-key buffer: at most one ADD and one + // DELETE for the same Identifier + ManifestEntrySortKey currentKey = null; + ManifestEntry bufferedAdd = null; + ManifestEntry bufferedDelete = null; + try { + while (iter.hasNext()) { + Tuple2 pair = + iter.next(); + ManifestEntrySortKey key = pair._1; + ManifestEntry entry = + entrySer.deserializeFromBytes(pair._2); + + if (currentKey != null + && currentKey.compareTo(key) != 0) { + // key changed: flush the previous group + ManifestEntry survived = + mergeGroup( + bufferedAdd, + bufferedDelete); + if (survived != null) { + writer.write(survived); + } + bufferedAdd = null; + bufferedDelete = null; + } + + currentKey = key; + if (entry.kind() == FileKind.ADD) { + if (bufferedAdd != null) { + throw new IllegalStateException( + "Duplicate ADD entry for " + + entry.identifier()); + } + bufferedAdd = entry; + } else { + if (bufferedDelete != null) { + throw new IllegalStateException( + "Duplicate DELETE entry for " + + entry.identifier()); + } + bufferedDelete = entry; + } + } + + // flush the last group + if (currentKey != null) { + ManifestEntry survived = + mergeGroup(bufferedAdd, bufferedDelete); + if (survived != null) { + writer.write(survived); + } + } + } finally { + writer.close(); + } + + if (writer.recordCount() == 0) { + // nothing survived — delete the empty file and + // emit nothing + manifestFile.delete(writer.path().getName()); + return Collections.emptyList().iterator(); + } + ManifestFileMeta newMeta = writer.result(); + ManifestFileMetaSerializer ser = + new ManifestFileMetaSerializer(); + return Collections.singletonList( + ser.serializeToBytes(newMeta)) + .iterator(); + }) + .collect(); + + List newManifests = new ArrayList<>(serializedResult.size()); + for (byte[] bytes : serializedResult) { + try { + newManifests.add(metaSerializer.deserializeFromBytes(bytes)); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + return newManifests; + } + + /** + * Merge the ADD and DELETE entries of the same file (same Identifier) in an order-independent + * way. Returns the surviving ADD if only an ADD is present, {@code null} if the ADD and DELETE + * cancel each other out, or the DELETE if only a DELETE is present (kept so it can match an ADD + * in a previous manifest). + */ + private static ManifestEntry mergeGroup(ManifestEntry addEntry, ManifestEntry deleteEntry) { + if (addEntry != null && deleteEntry != null) { + // ADD + DELETE cancel out + return null; + } + if (addEntry != null) { + return addEntry; + } + return deleteEntry; + } + + /** + * Estimate the number of output manifests (and thus the sort parallelism) from the input + * manifests. The surviving entry count is {@code added - deleted} (ADD/DELETE pairs cancel), + * the average entry size is {@code totalSize / (added + deleted)}, and the estimated output + * size is {@code avgEntrySize * survivingEntries}. The parallelism is the floor of that over + * the manifest target file size, so each task produces roughly one manifest of the target size. + */ + private int computeNumPartitions(FileStoreTable table, List manifests) { + long targetSizeBytes = table.coreOptions().manifestTargetSize().getBytes(); + long totalSizeBytes = 0L; + long addedEntries = 0L; + long deletedEntries = 0L; + for (ManifestFileMeta manifest : manifests) { + totalSizeBytes += manifest.fileSize(); + addedEntries += manifest.numAddedFiles(); + deletedEntries += manifest.numDeletedFiles(); + } + + if (totalSizeBytes <= 0) { + throw new IllegalStateException( + "Cannot compute parallelism: total manifest size is " + totalSizeBytes); + } + long totalEntries = addedEntries + deletedEntries; + if (totalEntries <= 0) { + throw new IllegalStateException( + "Cannot compute parallelism: total manifest entries is " + totalEntries); + } + long survivingEntries = addedEntries - deletedEntries; + if (survivingEntries < 0) { + throw new IllegalStateException( + "Cannot compute parallelism: surviving entries (added - deleted) is " + + survivingEntries); + } + double avgEntrySizeBytes = (double) totalSizeBytes / totalEntries; + long estimatedOutputSizeBytes = (long) (avgEntrySizeBytes * survivingEntries); + // floor division: prefer fewer, slightly-over-target manifests over more, under-target + // ones. For example, 35M estimated with 8M target -> 4 manifests (~8.75M each), not 5 + // (~7M each). At least one task is always produced. + return (int) Math.max(1, estimatedOutputSizeBytes / targetSizeBytes); + } + + @Override + public String description() { + return "This procedure rewrites and globally sorts manifest entries."; + } + + public static ProcedureBuilder builder() { + return new Builder() { + @Override + public RewriteManifestProcedure doBuild() { + return new RewriteManifestProcedure(tableCatalog()); + } + }; + } +} diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RewriteManifestProcedureTest.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RewriteManifestProcedureTest.scala new file mode 100644 index 000000000000..637e73fa6b40 --- /dev/null +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RewriteManifestProcedureTest.scala @@ -0,0 +1,363 @@ +/* + * 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.procedure + +import org.apache.paimon.spark.PaimonSparkTestBase + +import org.apache.spark.sql.streaming.StreamTest +import org.assertj.core.api.Assertions + +import scala.jdk.CollectionConverters._ + +/** Test rewrite manifest procedure. See [[RewriteManifestProcedure]]. */ +class RewriteManifestProcedureTest extends PaimonSparkTestBase with StreamTest { + + test("Paimon Procedure: rewrite manifest cleans delete entries and keeps data intact") { + // A small manifest target file size forces several manifest files, so that the global sort + // actually reshuffles entries across files. + spark.sql(s""" + |CREATE TABLE T (id INT, value STRING, dt STRING) + |TBLPROPERTIES ( + | 'bucket' = '2', + | 'bucket-key' = 'id', + | 'write-only' = 'true', + | 'file.format' = 'avro', + | 'manifest.target-file-size' = '1KB' + |) + |PARTITIONED BY (dt) + |""".stripMargin) + + // Insert partitions in non-sorted order across multiple commits so that manifest entries are + // interleaved and out of partition order. + spark.sql(s"INSERT INTO T VALUES (1, 'a', '2024-01-03')") + spark.sql(s"INSERT INTO T VALUES (2, 'b', '2024-01-01')") + spark.sql(s"INSERT INTO T VALUES (3, 'c', '2024-01-02')") + spark.sql(s"INSERT INTO T VALUES (4, 'd', '2024-01-01')") + spark.sql(s"INSERT INTO T VALUES (5, 'e', '2024-01-03')") + + // Partition-level overwrites produce ADD/DELETE pairs that rewrite_manifest should cancel, + // without dropping the other partitions. + spark.sql(s"INSERT OVERWRITE T PARTITION (dt = '2024-01-03') VALUES (1, 'a2')") + spark.sql(s"INSERT OVERWRITE T PARTITION (dt = '2024-01-01') VALUES (2, 'b2'), (4, 'd2')") + + Thread.sleep(10000) + + val expectedCount = spark.sql("SELECT count(*) FROM T").collectAsList().get(0).getLong(0) + val expectedIdSum = spark.sql("SELECT sum(id) FROM T").collectAsList().get(0).getLong(0) + + // before rewrite there should be some delete entries + val beforeDeleted = + spark + .sql("SELECT sum(num_deleted_files) FROM `T$manifests`") + .collectAsList() + .get(0) + .getLong(0) + Assertions.assertThat(beforeDeleted).isGreaterThan(0L) + + val beforeManifests = + spark.sql("SELECT count(*) FROM `T$manifests`").collectAsList().get(0).getLong(0) + + val result = + spark.sql("CALL sys.rewrite_manifest(table => 'T')").collectAsList().get(0) + + // rewritten_manifests_count > 0 and added_manifests_count == after - before + val rewrittenCount = result.getInt(0) + val addedCount = result.getInt(1) + Assertions.assertThat(rewrittenCount).isGreaterThan(0) + Assertions.assertThat(addedCount).isEqualTo(rewrittenCount - beforeManifests.toInt) + + // after rewrite all delete entries must be cleaned + val afterDeleted = + spark + .sql("SELECT sum(num_deleted_files) FROM `T$manifests`") + .collectAsList() + .get(0) + .getLong(0) + Assertions.assertThat(afterDeleted).isEqualTo(0L) + + // data must be intact + Assertions + .assertThat(spark.sql("SELECT count(*) FROM T").collectAsList().get(0).getLong(0)) + .isEqualTo(expectedCount) + Assertions + .assertThat(spark.sql("SELECT sum(id) FROM T").collectAsList().get(0).getLong(0)) + .isEqualTo(expectedIdSum) + } + + test("Paimon Procedure: rewrite manifest orders manifest partitions globally") { + spark.sql(s""" + |CREATE TABLE T2 (id INT, dt STRING) + |TBLPROPERTIES ( + | 'bucket' = '-1', + | 'write-only' = 'true', + | 'manifest.target-file-size' = '1KB', + | 'file.format' = 'avro' + |) + |PARTITIONED BY (dt) + |""".stripMargin) + + // insert partitions in reverse order across commits + spark.sql(s"INSERT INTO T2 VALUES (1, '2024-01-03')") + spark.sql(s"INSERT INTO T2 VALUES (2, '2024-01-02')") + spark.sql(s"INSERT INTO T2 VALUES (3, '2024-01-01')") + spark.sql(s"INSERT INTO T2 VALUES (4, '2024-01-03')") + spark.sql(s"INSERT INTO T2 VALUES (5, '2024-01-01')") + + Thread.sleep(10000) + + val expectedCount = spark.sql("SELECT count(*) FROM T2").collectAsList().get(0).getLong(0) + + spark.sql("CALL sys.rewrite_manifest(table => 'T2')") + + // After global sort, manifest partition ranges must not overlap: sorting manifests by their + // min partition must yield a sequence where each manifest's min is >= the previous manifest's + // max. This is robust regardless of the row order Spark returns for the system table. + val ranges = spark + .sql("SELECT min_partition_stats, max_partition_stats FROM `T2$manifests`") + .collectAsList() + .asScala + .filter(r => !r.isNullAt(0)) + .map(r => (r.getString(0), r.getString(1))) + .sortBy(_._1) + + for (i <- 1 until ranges.length) { + Assertions.assertThat(ranges(i)._1.compareTo(ranges(i - 1)._2) >= 0).isTrue + } + + // data intact + Assertions + .assertThat(spark.sql("SELECT count(*) FROM T2").collectAsList().get(0).getLong(0)) + .isEqualTo(expectedCount) + } + + test("Paimon Procedure: rewrite manifest with where only rewrites matching manifests") { + spark.sql(s""" + |CREATE TABLE T3 (id INT, dt STRING) + |TBLPROPERTIES ( + | 'bucket' = '-1', + | 'write-only' = 'true', + | 'manifest.target-file-size' = '1KB' + |) + |PARTITIONED BY (dt) + |""".stripMargin) + + // insert into multiple partitions across commits + spark.sql(s"INSERT INTO T3 VALUES (1, '2024-01-01')") + spark.sql(s"INSERT INTO T3 VALUES (2, '2024-01-02')") + spark.sql(s"INSERT INTO T3 VALUES (3, '2024-01-03')") + spark.sql(s"INSERT INTO T3 VALUES (4, '2024-01-01')") + + Thread.sleep(10000) + + val expectedCount = spark.sql("SELECT count(*) FROM T3").collectAsList().get(0).getLong(0) + + // record manifest file names before rewrite + val allManifestsBefore = + spark.sql("SELECT file_name FROM `T3$manifests`").collectAsList().asScala.map(_.getString(0)) + + // rewrite only manifests that may match dt = '2024-01-01' + spark.sql("CALL sys.rewrite_manifest(table => 'T3', where => 'dt = \"2024-01-01\"')") + + val allManifestsAfter = + spark.sql("SELECT file_name FROM `T3$manifests`").collectAsList().asScala.map(_.getString(0)) + + // some manifests should have been rewritten (new file names appear) + val newManifests = allManifestsAfter.filter(!allManifestsBefore.contains(_)) + Assertions.assertThat(newManifests.nonEmpty).isTrue + + // but not all manifests are rewritten — at least one original manifest survives + val survivingOriginals = allManifestsBefore.filter(allManifestsAfter.contains(_)) + Assertions.assertThat(survivingOriginals.nonEmpty).isTrue + + // data intact + Assertions + .assertThat(spark.sql("SELECT count(*) FROM T3").collectAsList().get(0).getLong(0)) + .isEqualTo(expectedCount) + Assertions + .assertThat( + spark + .sql("SELECT count(*) FROM T3 WHERE dt = '2024-01-01'") + .collectAsList() + .get(0) + .getLong(0)) + .isEqualTo(2) + } + + test("Paimon Procedure: rewrite manifest with where matching nothing returns zero") { + spark.sql(s""" + |CREATE TABLE T4 (id INT, dt STRING) + |TBLPROPERTIES ( + | 'bucket' = '-1', + | 'write-only' = 'true', + | 'manifest.target-file-size' = '1KB', + | 'file.format' = 'avro' + |) + |PARTITIONED BY (dt) + |""".stripMargin) + + spark.sql(s"INSERT INTO T4 VALUES (1, '2024-01-01')") + spark.sql(s"INSERT INTO T4 VALUES (2, '2024-01-02')") + Thread.sleep(10000) + + val manifestsBefore = + spark.sql("SELECT file_name FROM `T4$manifests`").collectAsList().asScala.map(_.getString(0)) + + // where matches no partition — nothing should be rewritten + val result = + spark + .sql("CALL sys.rewrite_manifest(table => 'T4', where => 'dt = \"1999-01-01\"')") + .collectAsList() + .get(0) + Assertions.assertThat(result.getInt(0)).isEqualTo(0) + Assertions.assertThat(result.getInt(1)).isEqualTo(0) + + // all original manifests survive untouched + val manifestsAfter = + spark.sql("SELECT file_name FROM `T4$manifests`").collectAsList().asScala.map(_.getString(0)) + Assertions.assertThat(manifestsAfter).isEqualTo(manifestsBefore) + + // data intact + Assertions + .assertThat(spark.sql("SELECT count(*) FROM T4").collectAsList().get(0).getLong(0)) + .isEqualTo(2) + } + + test("Paimon Procedure: rewrite manifest with range where rewrites matching partitions") { + spark.sql(s""" + |CREATE TABLE T5 (id INT, dt STRING) + |TBLPROPERTIES ( + | 'bucket' = '-1', + | 'write-only' = 'true', + | 'manifest.target-file-size' = '1KB', + | 'file.format' = 'avro' + |) + |PARTITIONED BY (dt) + |""".stripMargin) + + spark.sql(s"INSERT INTO T5 VALUES (1, '2024-01-01')") + spark.sql(s"INSERT INTO T5 VALUES (2, '2024-01-02')") + spark.sql(s"INSERT INTO T5 VALUES (3, '2024-01-03')") + spark.sql(s"INSERT INTO T5 VALUES (4, '2024-01-04')") + Thread.sleep(10000) + + val expectedCount = spark.sql("SELECT count(*) FROM T5").collectAsList().get(0).getLong(0) + + // range where: only Jan 1-2 + spark.sql( + "CALL sys.rewrite_manifest(table => 'T5', where => 'dt >= \"2024-01-01\" AND dt <= \"2024-01-02\"')") + + // data intact across all partitions + Assertions + .assertThat(spark.sql("SELECT count(*) FROM T5").collectAsList().get(0).getLong(0)) + .isEqualTo(expectedCount) + Assertions + .assertThat( + spark + .sql("SELECT count(*) FROM T5 WHERE dt >= '2024-01-01' AND dt <= '2024-01-02'") + .collectAsList() + .get(0) + .getLong(0)) + .isEqualTo(2) + } + + test("Paimon Procedure: rewrite manifest on unpartitioned table") { + spark.sql(s""" + |CREATE TABLE T6 (id INT, value STRING) + |TBLPROPERTIES ( + | 'bucket' = '-1', + | 'write-only' = 'true', + | 'manifest.target-file-size' = '1KB', + | 'file.format' = 'avro' + |) + |""".stripMargin) + + spark.sql(s"INSERT INTO T6 VALUES (1, 'a')") + spark.sql(s"INSERT INTO T6 VALUES (2, 'b')") + Thread.sleep(10000) + + val expectedCount = spark.sql("SELECT count(*) FROM T6").collectAsList().get(0).getLong(0) + + spark.sql("CALL sys.rewrite_manifest(table => 'T6')") + + // data intact + Assertions + .assertThat(spark.sql("SELECT count(*) FROM T6").collectAsList().get(0).getLong(0)) + .isEqualTo(expectedCount) + } + + test("Paimon Procedure: rewrite manifest with non-partition where throws") { + spark.sql(s""" + |CREATE TABLE T7 (id INT, dt STRING) + |TBLPROPERTIES ( + | 'bucket' = '-1', + | 'write-only' = 'true', + | 'manifest.target-file-size' = '1KB', + | 'file.format' = 'avro' + |) + |PARTITIONED BY (dt) + |""".stripMargin) + + spark.sql(s"INSERT INTO T7 VALUES (1, '2024-01-01')") + Thread.sleep(10000) + + // where on a non-partition column (id) must fail + Assertions + .assertThatThrownBy( + () => spark.sql("CALL sys.rewrite_manifest(table => 'T7', where => 'id = 1')")) + .hasMessageContaining("Only partition predicate is supported") + } + + test("Paimon Procedure: rewritten manifest sizes are within the target bound") { + // Use a small target file size and enough data so that multiple manifests are produced. + // Each task writes a single non-rolling manifest, so every output manifest should be at most + // a modest multiple of the target size (allowing headroom for a single partition that + // slightly exceeds the range estimate). + val targetSize = 1024L + spark.sql(s""" + |CREATE TABLE T8 (id INT, value STRING, dt STRING) + |TBLPROPERTIES ( + | 'bucket' = '4', + | 'bucket-key' = 'id', + | 'write-only' = 'true', + | 'file.format' = 'avro', + | 'manifest.target-file-size' = '${targetSize}B' + |) + |PARTITIONED BY (dt) + |""".stripMargin) + + // Insert enough rows across enough partitions to produce multiple output manifests. + for (dt <- 0 until 20) { + val values = (0 until 20).map(i => s"(${dt * 20 + i}, '${"x" * 50}', '2024-01-${dt + 1}')") + spark.sql(s"INSERT INTO T8 VALUES ${values.mkString(", ")}") + } + Thread.sleep(10000) + + spark.sql("CALL sys.rewrite_manifest(table => 'T8')") + + val sizes = + spark.sql("SELECT file_size FROM `T8$manifests`").collectAsList().asScala.map(_.getLong(0)) + Assertions.assertThat(sizes.nonEmpty).isTrue + // Every rewritten manifest should be <= 3x target size (generous bound for the last manifest + // of a task and range-partitioner skew). The key invariant: no manifest should be wildly + // oversized (e.g. 10x), which would indicate the parallelism estimate is wrong. + for (size <- sizes) { + Assertions.assertThat(size).isLessThanOrEqualTo(targetSize * 3) + } + } +} From 79fdb22a9879c345c33e876ec7ace52843e5c415 Mon Sep 17 00:00:00 2001 From: zhongyujiang Date: Fri, 14 Aug 2026 17:02:38 +0800 Subject: [PATCH 2/5] [spark]: fix output counts and rename filtered manifests variable rewritten_manifests_count should be the number of input manifests rewritten (currentManifests), not the output. added_manifests_count should be the number of new manifests produced (newManifests), not the net delta. Also rename the post-filter variable to manifestsToRewrite for clarity. Co-Authored-By: Claude --- .../procedure/RewriteManifestProcedure.java | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java index 8f1b32e8add4..bb2374f1b264 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java @@ -146,29 +146,29 @@ public InternalRow[] call(InternalRow args) { } // 2. filter manifests by the optional where clause (partition-stats pruning) + List manifestsToRewrite = currentManifests; PartitionPredicate partitionPredicate = resolvePartitionPredicate(tableIdent, table, where); if (partitionPredicate != null) { - currentManifests = filterManifests(currentManifests, partitionPredicate); - if (currentManifests.isEmpty()) { + manifestsToRewrite = filterManifests(manifestsToRewrite, partitionPredicate); + if (manifestsToRewrite.isEmpty()) { return new InternalRow[] {newInternalRow(0, 0)}; } } // 3. globally sort manifest entries and write them back as new manifest files - List newManifests = rewriteManifests(fileStoreTable, currentManifests); + List newManifests = rewriteManifests(fileStoreTable, manifestsToRewrite); // 4. commit the rewritten manifests (optimistic concurrency with retry) try (BatchTableCommit commit = fileStoreTable.newBatchWriteBuilder().newCommit()) { - commit.replaceManifests(currentManifests, newManifests); + commit.replaceManifests(manifestsToRewrite, newManifests); } catch (Exception e) { throw new RuntimeException(e); } - // rewritten_manifests_count: total manifests produced by the rewrite - // added_manifests_count: net change in manifest count vs before rewrite (may be negative - // when the rewrite also consolidates small manifest files) - int rewrittenCount = newManifests.size(); - int addedCount = newManifests.size() - currentManifests.size(); + // rewritten_manifests_count: number of manifests that were rewritten (input) + // added_manifests_count: number of new manifests produced by the rewrite (output) + int rewrittenCount = manifestsToRewrite.size(); + int addedCount = newManifests.size(); return new InternalRow[] {newInternalRow(rewrittenCount, addedCount)}; } From a78e1995d4317d0513454c857ac8782cf6660a4c Mon Sep 17 00:00:00 2001 From: zhongyujiang Date: Fri, 14 Aug 2026 17:49:45 +0800 Subject: [PATCH 3/5] [spark]: fix test assertion for corrected output count semantics Update the assertion to match the corrected semantics: rewritten_manifests_count == input manifest count, added_manifests_count == output manifest count (> 0). Co-Authored-By: Claude --- .../spark/procedure/RewriteManifestProcedureTest.scala | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RewriteManifestProcedureTest.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RewriteManifestProcedureTest.scala index 637e73fa6b40..43c066d1634c 100644 --- a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RewriteManifestProcedureTest.scala +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RewriteManifestProcedureTest.scala @@ -76,11 +76,12 @@ class RewriteManifestProcedureTest extends PaimonSparkTestBase with StreamTest { val result = spark.sql("CALL sys.rewrite_manifest(table => 'T')").collectAsList().get(0) - // rewritten_manifests_count > 0 and added_manifests_count == after - before + // rewritten_manifests_count == before (all manifests rewritten without where) + // added_manifests_count > 0 (new manifests produced) val rewrittenCount = result.getInt(0) val addedCount = result.getInt(1) - Assertions.assertThat(rewrittenCount).isGreaterThan(0) - Assertions.assertThat(addedCount).isEqualTo(rewrittenCount - beforeManifests.toInt) + Assertions.assertThat(rewrittenCount).isEqualTo(beforeManifests.toInt) + Assertions.assertThat(addedCount).isGreaterThan(0) // after rewrite all delete entries must be cleaned val afterDeleted = From 960fbc25f8c4d13e81c68c8d5885a33660145103 Mon Sep 17 00:00:00 2001 From: zhongyujiang Date: Sun, 16 Aug 2026 17:25:46 +0800 Subject: [PATCH 4/5] [spark]: fix testReplaceManifestsThrowsOnConcurrentCompact for apache The test expected compactManifest to change manifest file names, but with default settings the merge may be a no-op (same manifests, no new snapshot). Use a small manifest.target-file-size and manifest.full-compaction-file-size to force compactManifest to actually rewrite the manifests, making the stale manifestsBefore absent from the current base. Co-Authored-By: Claude --- .../apache/paimon/operation/FileStoreCommitTest.java | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java index 46877ae17c24..29462fe6e8d2 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreCommitTest.java @@ -50,6 +50,7 @@ import org.apache.paimon.operation.commit.ConflictDetection; import org.apache.paimon.operation.commit.ManifestEntryChanges; import org.apache.paimon.operation.commit.RetryCommitResult; +import org.apache.paimon.options.MemorySize; import org.apache.paimon.options.Options; import org.apache.paimon.predicate.PredicateBuilder; import org.apache.paimon.schema.Schema; @@ -2462,7 +2463,14 @@ public void testReplaceManifestsPreservesConcurrentDelta() throws Exception { @Test public void testReplaceManifestsThrowsOnConcurrentCompact() throws Exception { - TestFileStore store = createStore(false, 2); + // use a small manifest target size so compactManifest actually rewrites the manifests + // (changes file names), which makes the stale manifestsBefore absent from the current base + Map options = new HashMap<>(); + options.put(CoreOptions.MANIFEST_TARGET_FILE_SIZE.key(), "1KB"); + options.put( + CoreOptions.MANIFEST_FULL_COMPACTION_FILE_SIZE.key(), + MemorySize.ofBytes(1).toString()); + TestFileStore store = createStore(false, 2, CoreOptions.ChangelogProducer.NONE, options); for (int i = 0; i < 3; i++) { store.commitData(generateDataList(10), gen::getPartition, kv -> 0); } From 83953e5cfb0ede468d4723a521eecfab9b57fec6 Mon Sep 17 00:00:00 2001 From: zhongyujiang Date: Sun, 16 Aug 2026 22:55:54 +0800 Subject: [PATCH 5/5] [spark]: use ManifestAvroWriter instead of ManifestEntryWriter Apache paimon master uses ManifestAvroWriter (not ManifestEntryWriter). Replace createManifestEntryWriter(Path) with createAvroWriter(Path), and use result().isEmpty() + abort() for empty manifest cleanup instead of recordCount() + path().getName(). Co-Authored-By: Claude --- .../procedure/RewriteManifestProcedure.java | 17 ++++++++--------- 1 file changed, 8 insertions(+), 9 deletions(-) diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java index bb2374f1b264..b1ece9b72bc1 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/RewriteManifestProcedure.java @@ -21,11 +21,11 @@ import org.apache.paimon.Snapshot; import org.apache.paimon.fs.Path; import org.apache.paimon.manifest.FileKind; +import org.apache.paimon.manifest.ManifestAvroWriter; import org.apache.paimon.manifest.ManifestEntry; import org.apache.paimon.manifest.ManifestEntrySerializer; import org.apache.paimon.manifest.ManifestEntrySortKey; import org.apache.paimon.manifest.ManifestFile; -import org.apache.paimon.manifest.ManifestFile.ManifestEntryWriter; import org.apache.paimon.manifest.ManifestFileMeta; import org.apache.paimon.manifest.ManifestFileMetaSerializer; import org.apache.paimon.manifest.ManifestList; @@ -316,9 +316,8 @@ private List rewriteManifests( Path manifestPath = pathFactory.newManifestFile(); ManifestFile manifestFile = t.store().manifestFileFactory().create(); - ManifestEntryWriter writer = - manifestFile.createManifestEntryWriter( - manifestPath); + ManifestAvroWriter writer = + manifestFile.createAvroWriter(manifestPath); ManifestEntrySerializer entrySer = new ManifestEntrySerializer(); @@ -379,13 +378,13 @@ private List rewriteManifests( writer.close(); } - if (writer.recordCount() == 0) { - // nothing survived — delete the empty file and - // emit nothing - manifestFile.delete(writer.path().getName()); + List results = writer.result(); + if (results.isEmpty()) { + // nothing survived — clean up and emit nothing + writer.abort(); return Collections.emptyList().iterator(); } - ManifestFileMeta newMeta = writer.result(); + ManifestFileMeta newMeta = results.get(0); ManifestFileMetaSerializer ser = new ManifestFileMetaSerializer(); return Collections.singletonList(