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 cccc4456a9 [spark] Reduce driver memory usage when deserializing 
compact commit messages (#8897)
cccc4456a9 is described below

commit cccc4456a9e408eabe47ecf731411bee40f06b1c
Author: sanshi <[email protected]>
AuthorDate: Thu Jul 30 12:02:27 2026 +0800

    [spark] Reduce driver memory usage when deserializing compact commit 
messages (#8897)
---
 .../paimon/spark/procedure/CompactProcedure.java   | 42 ++++++++++++----------
 1 file changed, 23 insertions(+), 19 deletions(-)

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 6868bc2bcf..13a71d773c 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
@@ -381,11 +381,10 @@ public class CompactProcedure extends BaseProcedure {
 
         try (BatchTableCommit commit = writeBuilder.newCommit()) {
             CommitMessageSerializer serializer = new CommitMessageSerializer();
-            List<byte[]> serializedMessages = commitMessageJavaRDD.collect();
-            List<CommitMessage> messages = new 
ArrayList<>(serializedMessages.size());
-            for (byte[] serializedMessage : serializedMessages) {
-                messages.add(serializer.deserialize(serializer.getVersion(), 
serializedMessage));
-            }
+            List<byte[]> serializedMessages = new 
ArrayList<>(commitMessageJavaRDD.collect());
+            List<CommitMessage> messages =
+                    deserializeCommitMessagesAndReleaseSerializedBytes(
+                            serializer, serializedMessages);
             commit.commit(messages);
         } catch (Exception e) {
             throw new RuntimeException(e);
@@ -480,13 +479,10 @@ public class CompactProcedure extends BaseProcedure {
 
         try (TableCommitImpl commit = table.newCommit(commitUser)) {
             CommitMessageSerializer messageSerializerser = new 
CommitMessageSerializer();
-            List<byte[]> serializedMessages = commitMessageJavaRDD.collect();
-            List<CommitMessage> messages = new 
ArrayList<>(serializedMessages.size());
-            for (byte[] serializedMessage : serializedMessages) {
-                messages.add(
-                        messageSerializerser.deserialize(
-                                messageSerializerser.getVersion(), 
serializedMessage));
-            }
+            List<byte[]> serializedMessages = new 
ArrayList<>(commitMessageJavaRDD.collect());
+            List<CommitMessage> messages =
+                    deserializeCommitMessagesAndReleaseSerializedBytes(
+                            messageSerializerser, serializedMessages);
             commit.commit(messages);
         } catch (Exception e) {
             throw new RuntimeException(e);
@@ -578,14 +574,11 @@ public class CompactProcedure extends BaseProcedure {
                                                     return 
messagesBytes.iterator();
                                                 });
 
-                List<CommitMessage> messages = new ArrayList<>();
-                List<byte[]> serializedMessages = 
commitMessageJavaRDD.collect();
+                List<byte[]> serializedMessages = new 
ArrayList<>(commitMessageJavaRDD.collect());
                 try (TableCommitImpl commit = table.newCommit(commitUser)) {
-                    for (byte[] serializedMessage : serializedMessages) {
-                        messages.add(
-                                messageSerializerser.deserialize(
-                                        messageSerializerser.getVersion(), 
serializedMessage));
-                    }
+                    List<CommitMessage> messages =
+                            deserializeCommitMessagesAndReleaseSerializedBytes(
+                                    messageSerializerser, serializedMessages);
                     messages.addAll(
                             new 
DataEvolutionCompactionCommitPreparation(table, snapshot)
                                     .prepare(messages));
@@ -599,6 +592,17 @@ public class CompactProcedure extends BaseProcedure {
         }
     }
 
+    private static List<CommitMessage> 
deserializeCommitMessagesAndReleaseSerializedBytes(
+            CommitMessageSerializer serializer, List<byte[]> 
serializedMessages)
+            throws IOException {
+        List<CommitMessage> messages = new 
ArrayList<>(serializedMessages.size());
+        for (int i = 0; i < serializedMessages.size(); i++) {
+            byte[] serializedMessage = serializedMessages.set(i, null);
+            messages.add(serializer.deserialize(serializer.getVersion(), 
serializedMessage));
+        }
+        return messages;
+    }
+
     private Set<BinaryRow> getHistoryPartition(
             SnapshotReader snapshotReader, @Nullable Duration 
partitionIdleTime) {
         Set<Pair<BinaryRow, Long>> partitionInfo =

Reply via email to