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