From c4ed19555886d4fa9c07ebc3817e87425ce70734 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BB=9F=E5=BC=8B?= Date: Thu, 20 Aug 2026 15:52:33 +0800 Subject: [PATCH] [core] Encapsulate split serialization protocols --- .../table/FallbackReadFileStoreTable.java | 17 +- .../paimon/table/source/IncrementalSplit.java | 57 +++-- .../paimon/table/source/QueryAuthSplit.java | 120 +++++++++- .../paimon/table/source/SplitSerializer.java | 213 +----------------- .../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 d22753a9a9ee..9b26b26f6932 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 OptionalLong mergedRowCount() { } 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 @@ private void assign(FallbackSplitImpl other) { } 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 5b672a4fefd5..b4bd8be490ce 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.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 String toString() { } 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 @@ private void writeObject(ObjectOutputStream objectOutputStream) throws IOExcepti 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 deletionFileSerializer = DeletionFile::deserialize; int beforeNumber = in.readInt(); - beforeFiles = new ArrayList<>(beforeNumber); + List beforeFiles = new ArrayList<>(beforeNumber); for (int i = 0; i < beforeNumber; i++) { beforeFiles.add(dataFileMetaSerializer.deserialize(in)); } - beforeDeletionFiles = DeletionFile.deserializeList(in, deletionFileSerializer); + List beforeDeletionFiles = + DeletionFile.deserializeList(in, deletionFileSerializer); int fileNumber = in.readInt(); - afterFiles = new ArrayList<>(fileNumber); + List afterFiles = new ArrayList<>(fileNumber); for (int i = 0; i < fileNumber; i++) { afterFiles.add(dataFileMetaSerializer.deserialize(in)); } - afterDeletionFiles = DeletionFile.deserializeList(in, deletionFileSerializer); + List 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 84646181e11f..3f6106d015ad 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 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 OptionalLong mergedRowCount() { } 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 @@ private void assign(QueryAuthSplit other) { } 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 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 readStringList(DataInputView in) throws IOException { + if (!in.readBoolean()) { + return null; + } + + int size = in.readInt(); + List 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 map) + throws IOException { + if (map == null) { + out.writeBoolean(false); + return; + } + + out.writeBoolean(true); + out.writeInt(map.size()); + for (Map.Entry entry : new TreeMap<>(map).entrySet()) { + writeString(out, entry.getKey()); + writeString(out, entry.getValue()); + } + } + + @Nullable + private static Map readNullableStringMap(DataInputView in) throws IOException { + if (!in.readBoolean()) { + return null; + } + + int size = in.readInt(); + Map 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 070956b811d3..d16e08cb7eed 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 static void serialize(Split split, DataOutputView out) throws IOException 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 static void serialize(Split split, DataOutputView out) throws IOException ((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 static Split deserialize(DataInputView in) throws IOException { 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 beforeFiles = readDataFiles(in); - FunctionWithIOException deletionFileSerializer = - DeletionFile::deserialize; - List beforeDeletionFiles = - DeletionFile.deserializeList(in, deletionFileSerializer); - List afterFiles = readDataFiles(in); - List 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 files, DataOutputView out) - throws IOException { - DataFileMetaSerializer serializer = new DataFileMetaSerializer(); - out.writeInt(files.size()); - for (DataFileMeta file : files) { - serializer.serialize(file, out); - } - } - - private static List readDataFiles(DataInputView in) throws IOException { - int size = in.readInt(); - List 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 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 readStringList(DataInputView in) throws IOException { - if (!in.readBoolean()) { - return null; - } - - int size = in.readInt(); - List 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 map) - throws IOException { - if (map == null) { - out.writeBoolean(false); - return; - } - - out.writeBoolean(true); - out.writeInt(map.size()); - for (Map.Entry entry : new TreeMap<>(map).entrySet()) { - writeString(out, entry.getKey()); - writeString(out, entry.getValue()); - } - } - - @Nullable - private static Map readNullableStringMap(DataInputView in) throws IOException { - if (!in.readBoolean()) { - return null; - } - - int size = in.readInt(); - Map 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 f52e8022ec0958b5157f1d17e842272854156fc6..cb057e42d72f52a32f41d4361c578d5908fa0c1c 100644 GIT binary patch delta 29 gcmZ3)zKmTYIKam