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 =