Github user sihuazhou commented on a diff in the pull request: https://github.com/apache/flink/pull/6196#discussion_r197329533 --- Diff: flink-core/src/main/java/org/apache/flink/api/common/typeutils/CompositeSerializer.java --- @@ -0,0 +1,204 @@ +package org.apache.flink.api.common.typeutils; + +import org.apache.flink.api.java.tuple.Tuple2; +import org.apache.flink.core.memory.DataInputView; +import org.apache.flink.core.memory.DataOutputView; +import org.apache.flink.util.Preconditions; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.List; +import java.util.stream.Collectors; +import java.util.stream.IntStream; + +/** + * Base class for composite serializers. + * + * <p>This class serializes a list of objects + * + * @param <T> type of custom serialized value + */ +@SuppressWarnings("unchecked") +public abstract class CompositeSerializer<T> extends TypeSerializer<T> { + private final List<TypeSerializer> originalSerializers; + + protected CompositeSerializer(List<TypeSerializer> originalSerializers) { + Preconditions.checkNotNull(originalSerializers); + this.originalSerializers = originalSerializers; + } + + protected abstract T composeValue(List values); + + protected abstract List decomposeValue(T v); + + protected abstract CompositeSerializer<T> createSerializerInstance(List<TypeSerializer> originalSerializers); + + private T composeValueInternal(List values) { + Preconditions.checkArgument(values.size() == originalSerializers.size()); + return composeValue(values); + } + + private List decomposeValueInternal(T v) { + List values = decomposeValue(v); + Preconditions.checkArgument(values.size() == originalSerializers.size()); + return values; + } + + private CompositeSerializer<T> createSerializerInstanceInternal(List<TypeSerializer> originalSerializers) { + Preconditions.checkArgument(originalSerializers.size() == originalSerializers.size()); --- End diff -- I think this check looks like a bug.
---