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 e5265e7ecc [core] Encapsulate split serialization protocols (#9316)
e5265e7ecc is described below

commit e5265e7ecc893f9b831b027cbb3b3c9c6c5562bd
Author: YeJunHao <[email protected]>
AuthorDate: Thu Aug 20 16:24:39 2026 +0800

    [core] Encapsulate split serialization protocols (#9316)
---
 .../paimon/table/FallbackReadFileStoreTable.java   |  17 +-
 .../paimon/table/source/IncrementalSplit.java      |  57 ++++--
 .../apache/paimon/table/source/QueryAuthSplit.java | 120 +++++++++++-
 .../paimon/table/source/SplitSerializer.java       | 213 +--------------------
 .../resources/compatibility/split-v1-incremental   | Bin 930 -> 934 bytes
 5 files changed, 174 insertions(+), 233 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
index d22753a9a9..9b26b26f69 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
@@ -308,11 +308,15 @@ public class FallbackReadFileStoreTable extends 
DelegatedFileStoreTable {
         }
 
         private void writeObject(ObjectOutputStream out) throws IOException {
-            serialize(new DataOutputViewStreamWrapper(out));
+            SplitSerializer.serialize(this, new 
DataOutputViewStreamWrapper(out));
         }
 
         private void readObject(ObjectInputStream in) throws IOException, 
ClassNotFoundException {
-            assign(deserialize(new DataInputViewStreamWrapper(in)));
+            Split split = SplitSerializer.deserialize(new 
DataInputViewStreamWrapper(in));
+            if (!(split instanceof FallbackSplitImpl)) {
+                throw new IOException("Deserialized split is not a 
FallbackSplitImpl: " + split);
+            }
+            assign((FallbackSplitImpl) split);
         }
 
         private void assign(FallbackSplitImpl other) {
@@ -321,15 +325,14 @@ public class FallbackReadFileStoreTable extends 
DelegatedFileStoreTable {
         }
 
         public void serialize(DataOutputView out) throws IOException {
-            SplitSerializer.serialize(this, out);
+            out.writeBoolean(isFallback);
+            SplitSerializer.serialize(split, out);
         }
 
         public static FallbackSplitImpl deserialize(DataInputView in) throws 
IOException {
+            boolean isFallback = in.readBoolean();
             Split split = SplitSerializer.deserialize(in);
-            if (!(split instanceof FallbackSplitImpl)) {
-                throw new IOException("Deserialized split is not a 
FallbackSplitImpl: " + split);
-            }
-            return (FallbackSplitImpl) split;
+            return new FallbackSplitImpl(split, isFallback);
         }
     }
 
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java
 
b/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java
index 5b672a4fef..b4bd8be490 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/IncrementalSplit.java
@@ -23,6 +23,7 @@ import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.io.DataFileMetaSerializer;
 import org.apache.paimon.io.DataInputView;
 import org.apache.paimon.io.DataInputViewStreamWrapper;
+import org.apache.paimon.io.DataOutputView;
 import org.apache.paimon.io.DataOutputViewStreamWrapper;
 import org.apache.paimon.utils.FunctionWithIOException;
 
@@ -191,7 +192,27 @@ public class IncrementalSplit implements Split {
     }
 
     private void writeObject(ObjectOutputStream objectOutputStream) throws 
IOException {
-        DataOutputViewStreamWrapper out = new 
DataOutputViewStreamWrapper(objectOutputStream);
+        serialize(new DataOutputViewStreamWrapper(objectOutputStream));
+    }
+
+    private void readObject(ObjectInputStream objectInputStream)
+            throws IOException, ClassNotFoundException {
+        assign(deserialize(new DataInputViewStreamWrapper(objectInputStream)));
+    }
+
+    protected void assign(IncrementalSplit other) {
+        snapshotId = other.snapshotId;
+        partition = other.partition;
+        bucket = other.bucket;
+        totalBuckets = other.totalBuckets;
+        beforeFiles = other.beforeFiles;
+        beforeDeletionFiles = other.beforeDeletionFiles;
+        afterFiles = other.afterFiles;
+        afterDeletionFiles = other.afterDeletionFiles;
+        isStreaming = other.isStreaming;
+    }
+
+    public void serialize(DataOutputView out) throws IOException {
         out.writeInt(VERSION);
         out.writeLong(snapshotId);
         serializeBinaryRow(partition, out);
@@ -216,39 +237,49 @@ public class IncrementalSplit implements Split {
         out.writeBoolean(isStreaming);
     }
 
-    private void readObject(ObjectInputStream objectInputStream)
-            throws IOException, ClassNotFoundException {
-        DataInputViewStreamWrapper in = new 
DataInputViewStreamWrapper(objectInputStream);
+    public static IncrementalSplit deserialize(DataInputView in) throws 
IOException {
         int version = in.readInt();
         if (version != VERSION) {
             throw new UnsupportedOperationException("Unsupported version: " + 
version);
         }
 
-        snapshotId = in.readLong();
-        partition = deserializeBinaryRow(in);
-        bucket = in.readInt();
-        totalBuckets = in.readInt();
+        long snapshotId = in.readLong();
+        BinaryRow partition = deserializeBinaryRow(in);
+        int bucket = in.readInt();
+        int totalBuckets = in.readInt();
 
         DataFileMetaSerializer dataFileMetaSerializer = new 
DataFileMetaSerializer();
         FunctionWithIOException<DataInputView, DeletionFile> 
deletionFileSerializer =
                 DeletionFile::deserialize;
 
         int beforeNumber = in.readInt();
-        beforeFiles = new ArrayList<>(beforeNumber);
+        List<DataFileMeta> beforeFiles = new ArrayList<>(beforeNumber);
         for (int i = 0; i < beforeNumber; i++) {
             beforeFiles.add(dataFileMetaSerializer.deserialize(in));
         }
 
-        beforeDeletionFiles = DeletionFile.deserializeList(in, 
deletionFileSerializer);
+        List<DeletionFile> beforeDeletionFiles =
+                DeletionFile.deserializeList(in, deletionFileSerializer);
 
         int fileNumber = in.readInt();
-        afterFiles = new ArrayList<>(fileNumber);
+        List<DataFileMeta> afterFiles = new ArrayList<>(fileNumber);
         for (int i = 0; i < fileNumber; i++) {
             afterFiles.add(dataFileMetaSerializer.deserialize(in));
         }
 
-        afterDeletionFiles = DeletionFile.deserializeList(in, 
deletionFileSerializer);
+        List<DeletionFile> afterDeletionFiles =
+                DeletionFile.deserializeList(in, deletionFileSerializer);
 
-        isStreaming = in.readBoolean();
+        boolean isStreaming = in.readBoolean();
+        return new IncrementalSplit(
+                snapshotId,
+                partition,
+                bucket,
+                totalBuckets,
+                beforeFiles,
+                beforeDeletionFiles,
+                afterFiles,
+                afterDeletionFiles,
+                isStreaming);
     }
 }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/QueryAuthSplit.java 
b/paimon-core/src/main/java/org/apache/paimon/table/source/QueryAuthSplit.java
index 84646181e1..3f6106d015 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/QueryAuthSplit.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/QueryAuthSplit.java
@@ -29,8 +29,13 @@ import javax.annotation.Nullable;
 import java.io.IOException;
 import java.io.ObjectInputStream;
 import java.io.ObjectOutputStream;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
 import java.util.OptionalLong;
+import java.util.TreeMap;
 
 /** A wrapper class for {@link Split} that adds query authorization 
information. */
 public class QueryAuthSplit implements Split {
@@ -71,11 +76,15 @@ public class QueryAuthSplit implements Split {
     }
 
     private void writeObject(ObjectOutputStream out) throws IOException {
-        serialize(new DataOutputViewStreamWrapper(out));
+        SplitSerializer.serialize(this, new DataOutputViewStreamWrapper(out));
     }
 
     private void readObject(ObjectInputStream in) throws IOException, 
ClassNotFoundException {
-        assign(deserialize(new DataInputViewStreamWrapper(in)));
+        Split split = SplitSerializer.deserialize(new 
DataInputViewStreamWrapper(in));
+        if (!(split instanceof QueryAuthSplit)) {
+            throw new IOException("Deserialized split is not a QueryAuthSplit: 
" + split);
+        }
+        assign((QueryAuthSplit) split);
     }
 
     private void assign(QueryAuthSplit other) {
@@ -84,14 +93,113 @@ public class QueryAuthSplit implements Split {
     }
 
     public void serialize(DataOutputView out) throws IOException {
-        SplitSerializer.serialize(this, out);
+        SplitSerializer.serialize(split, out);
+        writeAuthResult(out, authResult);
     }
 
     public static QueryAuthSplit deserialize(DataInputView in) throws 
IOException {
         Split split = SplitSerializer.deserialize(in);
-        if (!(split instanceof QueryAuthSplit)) {
-            throw new IOException("Deserialized split is not a QueryAuthSplit: 
" + split);
+        return new QueryAuthSplit(split, readAuthResult(in));
+    }
+
+    private static void writeAuthResult(
+            DataOutputView out, @Nullable TableQueryAuthResult authResult) 
throws IOException {
+        if (authResult == null) {
+            out.writeBoolean(false);
+            return;
+        }
+
+        out.writeBoolean(true);
+        writeStringList(out, authResult.filter());
+        writeStringMap(out, authResult.columnMasking());
+    }
+
+    @Nullable
+    private static TableQueryAuthResult readAuthResult(DataInputView in) 
throws IOException {
+        if (!in.readBoolean()) {
+            return null;
+        }
+        return new TableQueryAuthResult(readStringList(in), 
readNullableStringMap(in));
+    }
+
+    private static void writeStringList(DataOutputView out, @Nullable 
List<String> strings)
+            throws IOException {
+        if (strings == null) {
+            out.writeBoolean(false);
+            return;
+        }
+
+        out.writeBoolean(true);
+        out.writeInt(strings.size());
+        for (String string : strings) {
+            writeString(out, string);
         }
-        return (QueryAuthSplit) split;
+    }
+
+    @Nullable
+    private static List<String> readStringList(DataInputView in) throws 
IOException {
+        if (!in.readBoolean()) {
+            return null;
+        }
+
+        int size = in.readInt();
+        List<String> strings = new ArrayList<>(size);
+        for (int i = 0; i < size; i++) {
+            strings.add(readString(in));
+        }
+        return strings;
+    }
+
+    private static void writeStringMap(DataOutputView out, @Nullable 
Map<String, String> map)
+            throws IOException {
+        if (map == null) {
+            out.writeBoolean(false);
+            return;
+        }
+
+        out.writeBoolean(true);
+        out.writeInt(map.size());
+        for (Map.Entry<String, String> entry : new TreeMap<>(map).entrySet()) {
+            writeString(out, entry.getKey());
+            writeString(out, entry.getValue());
+        }
+    }
+
+    @Nullable
+    private static Map<String, String> readNullableStringMap(DataInputView in) 
throws IOException {
+        if (!in.readBoolean()) {
+            return null;
+        }
+
+        int size = in.readInt();
+        Map<String, String> map = new HashMap<>(size);
+        for (int i = 0; i < size; i++) {
+            map.put(readString(in), readString(in));
+        }
+        return map;
+    }
+
+    private static void writeString(DataOutputView out, @Nullable String 
string)
+            throws IOException {
+        if (string == null) {
+            out.writeInt(-1);
+            return;
+        }
+
+        byte[] bytes = string.getBytes(StandardCharsets.UTF_8);
+        out.writeInt(bytes.length);
+        out.write(bytes);
+    }
+
+    @Nullable
+    private static String readString(DataInputView in) throws IOException {
+        int length = in.readInt();
+        if (length < 0) {
+            return null;
+        }
+
+        byte[] bytes = new byte[length];
+        in.readFully(bytes);
+        return new String(bytes, StandardCharsets.UTF_8);
     }
 }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java 
b/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java
index 070956b811..d16e08cb7e 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/source/SplitSerializer.java
@@ -18,31 +18,15 @@
 
 package org.apache.paimon.table.source;
 
-import org.apache.paimon.catalog.TableQueryAuthResult;
-import org.apache.paimon.data.BinaryRow;
 import org.apache.paimon.globalindex.IndexedSplit;
-import org.apache.paimon.io.DataFileMeta;
-import org.apache.paimon.io.DataFileMetaSerializer;
 import org.apache.paimon.io.DataInputDeserializer;
 import org.apache.paimon.io.DataInputView;
 import org.apache.paimon.io.DataOutputView;
 import org.apache.paimon.io.DataOutputViewStreamWrapper;
 import org.apache.paimon.table.FallbackReadFileStoreTable;
-import org.apache.paimon.utils.FunctionWithIOException;
-
-import javax.annotation.Nullable;
 
 import java.io.ByteArrayOutputStream;
 import java.io.IOException;
-import java.nio.charset.StandardCharsets;
-import java.util.ArrayList;
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
-import java.util.TreeMap;
-
-import static org.apache.paimon.utils.SerializationUtils.deserializeBinaryRow;
-import static org.apache.paimon.utils.SerializationUtils.serializeBinaryRow;
 
 /**
  * Versioned binary serializer for non-system table {@link Split}s.
@@ -78,13 +62,13 @@ public class SplitSerializer {
 
         if (split instanceof QueryAuthSplit) {
             out.writeInt(QUERY_AUTH_SPLIT);
-            writeQueryAuthSplit((QueryAuthSplit) split, out);
+            ((QueryAuthSplit) split).serialize(out);
         } else if (split instanceof 
FallbackReadFileStoreTable.FallbackDataSplit) {
             out.writeInt(FALLBACK_DATA_SPLIT);
             ((FallbackReadFileStoreTable.FallbackDataSplit) 
split).serialize(out);
         } else if (split instanceof 
FallbackReadFileStoreTable.FallbackSplitImpl) {
             out.writeInt(FALLBACK_SPLIT);
-            writeFallbackSplit((FallbackReadFileStoreTable.FallbackSplitImpl) 
split, out);
+            ((FallbackReadFileStoreTable.FallbackSplitImpl) 
split).serialize(out);
         } else if (split instanceof IndexedSplit) {
             out.writeInt(INDEXED_SPLIT);
             ((IndexedSplit) split).serialize(out);
@@ -93,7 +77,7 @@ public class SplitSerializer {
             ((ChainSplit) split).serialize(out);
         } else if (split instanceof IncrementalSplit) {
             out.writeInt(INCREMENTAL_SPLIT);
-            writeIncrementalSplit((IncrementalSplit) split, out);
+            ((IncrementalSplit) split).serialize(out);
         } else if (split instanceof DataSplit) {
             out.writeInt(DATA_SPLIT);
             ((DataSplit) split).serialize(out);
@@ -122,204 +106,19 @@ public class SplitSerializer {
             case DATA_SPLIT:
                 return DataSplit.deserialize(in);
             case INCREMENTAL_SPLIT:
-                return readIncrementalSplit(in);
+                return IncrementalSplit.deserialize(in);
             case INDEXED_SPLIT:
                 return IndexedSplit.deserialize(in);
             case CHAIN_SPLIT:
                 return ChainSplit.deserialize(in);
             case QUERY_AUTH_SPLIT:
-                return readQueryAuthSplit(in);
+                return QueryAuthSplit.deserialize(in);
             case FALLBACK_DATA_SPLIT:
                 return 
FallbackReadFileStoreTable.FallbackDataSplit.deserialize(in);
             case FALLBACK_SPLIT:
-                return readFallbackSplit(in);
+                return 
FallbackReadFileStoreTable.FallbackSplitImpl.deserialize(in);
             default:
                 throw new IOException("Unsupported split type: " + type);
         }
     }
-
-    private static void writeIncrementalSplit(IncrementalSplit split, 
DataOutputView out)
-            throws IOException {
-        out.writeLong(split.snapshotId());
-        serializeBinaryRow(split.partition(), out);
-        out.writeInt(split.bucket());
-        out.writeInt(split.totalBuckets());
-        writeDataFiles(split.beforeFiles(), out);
-        DeletionFile.serializeList(out, split.beforeDeletionFiles());
-        writeDataFiles(split.afterFiles(), out);
-        DeletionFile.serializeList(out, split.afterDeletionFiles());
-        out.writeBoolean(split.isStreaming());
-    }
-
-    private static IncrementalSplit readIncrementalSplit(DataInputView in) 
throws IOException {
-        long snapshotId = in.readLong();
-        BinaryRow partition = deserializeBinaryRow(in);
-        int bucket = in.readInt();
-        int totalBuckets = in.readInt();
-        List<DataFileMeta> beforeFiles = readDataFiles(in);
-        FunctionWithIOException<DataInputView, DeletionFile> 
deletionFileSerializer =
-                DeletionFile::deserialize;
-        List<DeletionFile> beforeDeletionFiles =
-                DeletionFile.deserializeList(in, deletionFileSerializer);
-        List<DataFileMeta> afterFiles = readDataFiles(in);
-        List<DeletionFile> afterDeletionFiles =
-                DeletionFile.deserializeList(in, deletionFileSerializer);
-        boolean isStreaming = in.readBoolean();
-        return new IncrementalSplit(
-                snapshotId,
-                partition,
-                bucket,
-                totalBuckets,
-                beforeFiles,
-                beforeDeletionFiles,
-                afterFiles,
-                afterDeletionFiles,
-                isStreaming);
-    }
-
-    private static void writeQueryAuthSplit(QueryAuthSplit split, 
DataOutputView out)
-            throws IOException {
-        serialize(split.split(), out);
-        writeAuthResult(out, split.authResult());
-    }
-
-    private static QueryAuthSplit readQueryAuthSplit(DataInputView in) throws 
IOException {
-        Split split = deserialize(in);
-        TableQueryAuthResult authResult = readAuthResult(in);
-        return new QueryAuthSplit(split, authResult);
-    }
-
-    private static void writeFallbackSplit(
-            FallbackReadFileStoreTable.FallbackSplitImpl split, DataOutputView 
out)
-            throws IOException {
-        out.writeBoolean(split.isFallback());
-        serialize(split.wrapped(), out);
-    }
-
-    private static FallbackReadFileStoreTable.FallbackSplitImpl 
readFallbackSplit(DataInputView in)
-            throws IOException {
-        boolean isFallback = in.readBoolean();
-        Split split = deserialize(in);
-        return new FallbackReadFileStoreTable.FallbackSplitImpl(split, 
isFallback);
-    }
-
-    private static void writeAuthResult(
-            DataOutputView out, @Nullable TableQueryAuthResult authResult) 
throws IOException {
-        if (authResult == null) {
-            out.writeBoolean(false);
-            return;
-        }
-
-        out.writeBoolean(true);
-        writeStringList(out, authResult.filter());
-        writeStringMap(out, authResult.columnMasking());
-    }
-
-    @Nullable
-    private static TableQueryAuthResult readAuthResult(DataInputView in) 
throws IOException {
-        if (!in.readBoolean()) {
-            return null;
-        }
-        return new TableQueryAuthResult(readStringList(in), 
readNullableStringMap(in));
-    }
-
-    private static void writeDataFiles(List<DataFileMeta> files, 
DataOutputView out)
-            throws IOException {
-        DataFileMetaSerializer serializer = new DataFileMetaSerializer();
-        out.writeInt(files.size());
-        for (DataFileMeta file : files) {
-            serializer.serialize(file, out);
-        }
-    }
-
-    private static List<DataFileMeta> readDataFiles(DataInputView in) throws 
IOException {
-        int size = in.readInt();
-        List<DataFileMeta> files = new ArrayList<>(size);
-        DataFileMetaSerializer serializer = new DataFileMetaSerializer();
-        for (int i = 0; i < size; i++) {
-            files.add(serializer.deserialize(in));
-        }
-        return files;
-    }
-
-    private static void writeStringList(DataOutputView out, @Nullable 
List<String> strings)
-            throws IOException {
-        if (strings == null) {
-            out.writeBoolean(false);
-            return;
-        }
-
-        out.writeBoolean(true);
-        out.writeInt(strings.size());
-        for (String string : strings) {
-            writeString(out, string);
-        }
-    }
-
-    @Nullable
-    private static List<String> readStringList(DataInputView in) throws 
IOException {
-        if (!in.readBoolean()) {
-            return null;
-        }
-
-        int size = in.readInt();
-        List<String> strings = new ArrayList<>(size);
-        for (int i = 0; i < size; i++) {
-            strings.add(readString(in));
-        }
-        return strings;
-    }
-
-    private static void writeStringMap(DataOutputView out, @Nullable 
Map<String, String> map)
-            throws IOException {
-        if (map == null) {
-            out.writeBoolean(false);
-            return;
-        }
-
-        out.writeBoolean(true);
-        out.writeInt(map.size());
-        for (Map.Entry<String, String> entry : new TreeMap<>(map).entrySet()) {
-            writeString(out, entry.getKey());
-            writeString(out, entry.getValue());
-        }
-    }
-
-    @Nullable
-    private static Map<String, String> readNullableStringMap(DataInputView in) 
throws IOException {
-        if (!in.readBoolean()) {
-            return null;
-        }
-
-        int size = in.readInt();
-        Map<String, String> map = new HashMap<>(size);
-        for (int i = 0; i < size; i++) {
-            map.put(readString(in), readString(in));
-        }
-        return map;
-    }
-
-    private static void writeString(DataOutputView out, @Nullable String 
string)
-            throws IOException {
-        if (string == null) {
-            out.writeInt(-1);
-            return;
-        }
-
-        byte[] bytes = string.getBytes(StandardCharsets.UTF_8);
-        out.writeInt(bytes.length);
-        out.write(bytes);
-    }
-
-    @Nullable
-    private static String readString(DataInputView in) throws IOException {
-        int length = in.readInt();
-        if (length < 0) {
-            return null;
-        }
-
-        byte[] bytes = new byte[length];
-        in.readFully(bytes);
-        return new String(bytes, StandardCharsets.UTF_8);
-    }
 }
diff --git a/paimon-core/src/test/resources/compatibility/split-v1-incremental 
b/paimon-core/src/test/resources/compatibility/split-v1-incremental
index f52e8022ec..cb057e42d7 100644
Binary files 
a/paimon-core/src/test/resources/compatibility/split-v1-incremental and 
b/paimon-core/src/test/resources/compatibility/split-v1-incremental differ

Reply via email to