Apache Flink类型及序列化研读&生产应用|得物技术
目录
一、背景
二、简单理论阐述(基于Flink 1.13)
1. 支持的数据类型
2. TypeInformation
3. 何时需要数据类型获取
4. 数据类型的自动推断
三、开发实践
1.Flink代码作业
1.1 如何显式指定数据类型
1.2 自定义数据类型&自定义序列化器
1.2.1 POJO类
1.2.2 自定义TypeInformation
1.2.3 自定义TypeSerializer
2. Flink SQL自定义函数
2.1 如何显式指定数据类型
2.1.1 指定accumulatorType
2.1.2 指定outputType
2.1.3 指定intputType
2.1.4 根据inputType动态调整outType
或者accumulatorType
2.2 自定义DataType
四、结语
一
背景
二
简单理论阐述(基于Flink 1.13)
支持的数据类型
Java Tuples and Scala Case Classes Java POJOs Primitive Types Regular Classes Values Hadoop Writables Special Types
TypeInformation
Basic types:所有的Java类型以及包装类:void,String,Date,BigDecimal,and BigInteger等。 Primitive arrays以及Object arrays Composite types Flink Java Tuples(Flink Java API的一部分):最多25个字段,不支持空字段 Scala case classes(包括Scala Tuples):不支持null字段 Row:具有任意数量字段并支持空字段的Tuples POJO 类:JavaBeans Auxiliary types (Option,Either,Lists,Maps,…) Generic types:Flink内部未维护的类型,这种类型通常是由Kryo序列化。
org.apache.flink.api.common.typeutils.TypeSerializer#serialize org.apache.flink.api.common.typeutils.TypeSerializer#deserialize(org.apache.flink.core.memory.DataInputView)
何时需要数据类型获取
@Experimentalpublic <OUT> DataStreamSource<OUT> fromSource(Source<OUT, ?, ?> source,WatermarkStrategy<OUT> timestampsAndWatermarks,String sourceName,TypeInformation<OUT> typeInfo) {final TypeInformation<OUT> resolvedTypeInfo =getTypeInfo(source, sourceName, Source.class, typeInfo);return new DataStreamSource<>(this,checkNotNull(source, "source"),checkNotNull(timestampsAndWatermarks, "timestampsAndWatermarks"),checkNotNull(resolvedTypeInfo),checkNotNull(sourceName));}
/** Returns the {@code TypeInformation} for the elements of the input. */public TypeInformation<IN> getInputType() {return input.getOutputType();}
@Overridepublic <T> ValueState<T> getState(ValueStateDescriptor<T> stateProperties) {KeyedStateStore keyedStateStore = checkPreconditionsAndGetKeyedStateStore(stateProperties);stateProperties.initializeSerializerUnlessSet(getExecutionConfig());return keyedStateStore.getState(stateProperties);}public void initializeSerializerUnlessSet(ExecutionConfig executionConfig) {if (serializerAtomicReference.get() == null) {checkState(typeInfo != null, "no serializer and no type info");// try to instantiate and set the serializerTypeSerializer<T> serializer = typeInfo.createSerializer(executionConfig);// use cas to assure the singletonif (!serializerAtomicReference.compareAndSet(null, serializer)) {LOG.debug("Someone else beat us at initializing the serializer.");}}}
数据类型的自动推断
public <R> SingleOutputStreamOperator<R> flatMap(FlatMapFunction<T, R> flatMapper) {TypeInformation<R> outType = TypeExtractor.getFlatMapReturnTypes((FlatMapFunction)this.clean(flatMapper), this.getType(), Utils.getCallLocationName(), true);return this.flatMap(flatMapper, outType);}
/*** A utility for reflection analysis on classes, to determine the return type of implementations of* transformation functions.** <p>NOTES FOR USERS OF THIS CLASS: Automatic type extraction is a hacky business that depends on a* lot of variables such as generics, compiler, interfaces, etc. The type extraction fails regularly* with either {@link MissingTypeInfo} or hard exceptions. Whenever you use methods of this class,* make sure to provide a way to pass custom type information as a fallback.*/
@PublicEvolvingpublic static <IN, OUT> TypeInformation<OUT> getUnaryOperatorReturnType(Function function,Class<?> baseClass,int inputTypeArgumentIndex,int outputTypeArgumentIndex,int[] lambdaOutputTypeArgumentIndices,TypeInformation<IN> inType,String functionName,boolean allowMissing) {Preconditions.checkArgument(inType == null || inputTypeArgumentIndex >= 0,"Input type argument index was not provided");Preconditions.checkArgument(outputTypeArgumentIndex >= 0, "Output type argument index was not provided");Preconditions.checkArgument(lambdaOutputTypeArgumentIndices != null,"Indices for output type arguments within lambda not provided");// explicit result type has highest precedenceif (function instanceof ResultTypeQueryable) {return ((ResultTypeQueryable<OUT>) function).getProducedType();}// perform extractiontry {final LambdaExecutable exec;try {exec = checkAndExtractLambda(function);} catch (TypeExtractionException e) {throw new InvalidTypesException("Internal error occurred.", e);}if (exec != null) {// parameters must be accessed from behind, since JVM can add additional parameters// e.g. when using local variables inside lambda function// paramLen is the total number of parameters of the provided lambda, it includes// parameters added through closurefinal int paramLen = exec.getParameterTypes().length;final Method sam = TypeExtractionUtils.getSingleAbstractMethod(baseClass);// number of parameters the SAM of implemented interface has; the parameter indexing// applies to this rangefinal int baseParametersLen = sam.getParameterTypes().length;final Type output;if (lambdaOutputTypeArgumentIndices.length > 0) {output =TypeExtractionUtils.extractTypeFromLambda(baseClass,exec,lambdaOutputTypeArgumentIndices,paramLen,baseParametersLen);} else {output = exec.getReturnType();TypeExtractionUtils.validateLambdaType(baseClass, output);}return new TypeExtractor().privateCreateTypeInfo(output, inType, null);} else {if (inType != null) {validateInputType(baseClass, function.getClass(), inputTypeArgumentIndex, inType);}return new TypeExtractor().privateCreateTypeInfo(baseClass,function.getClass(),outputTypeArgumentIndex,inType,null);}} catch (InvalidTypesException e) {if (allowMissing) {return (TypeInformation<OUT>)new MissingTypeInfo(functionName != null ? functionName : function.toString(), e);} else {throw e;}}}
首先判断该算子是否实现了ResultTypeQueryable接口,本质上就是用户是否显式指定了数据类型,例如我们熟悉的Kafka source就实现了该方法,当使用了JSONKeyValueDeserializationSchema,就显式指定了类型,用户自定义Schema就需要自己指定。
public class KafkaSource<OUT>implements Source<OUT, KafkaPartitionSplit, KafkaSourceEnumState>,ResultTypeQueryable<OUT>//deserializationSchema 是需要用户自己定义的。@Overridepublic TypeInformation<OUT> getProducedType() {return deserializationSchema.getProducedType();}//JSONKeyValueDeserializationSchema@Overridepublic TypeInformation<ObjectNode> getProducedType() {return getForClass(ObjectNode.class);}
未实现ResultTypeQueryable接口,就会通过反射的方法获取ReturnType,判断逻辑大概是从是否是Java 8 lambda方法开始判断的。获取到返回类型后再通过new TypeExtractor()).privateCreateTypeInfo(output,inType,(TypeInformation)null)封装成Flink内部能识别的数据类型;大致分为2类,泛型类型变量TypeVariable以及非泛型类型变量。这个封装的过程也是非常重要的,推断的数据类型是Flink内部封装好的类型,序列化基本都很高效,如果不是, 就会推断为GenericTypeInfo走Kryo等序列化方式。如感兴趣,可以看下这块的源码,在此不再赘述。
三
开发实践
Flink代码作业
如何显式指定数据类型
这个简单了,几乎所有的source、Keyby、算子等都暴露了指定TypeInformation<OUT> typeInfo的构造方法,以下简单列举几个:
source
@Experimentalpublic <OUT> DataStreamSource<OUT> fromSource(Source<OUT, ?, ?> source, WatermarkStrategy<OUT> timestampsAndWatermarks, String sourceName, TypeInformation<OUT> typeInfo) {TypeInformation<OUT> resolvedTypeInfo = this.getTypeInfo(source, sourceName, Source.class, typeInfo);return new DataStreamSource(this, (Source)Preconditions.checkNotNull(source, "source"), (WatermarkStrategy)Preconditions.checkNotNull(timestampsAndWatermarks, "timestampsAndWatermarks"), (TypeInformation)Preconditions.checkNotNull(resolvedTypeInfo), (String)Preconditions.checkNotNull(sourceName));}
map
public <R> SingleOutputStreamOperator<R> map(MapFunction<T, R> mapper, TypeInformation<R> outputType) {return transform("Map", outputType, new StreamMap<>(clean(mapper)));}
自定义Operator
@PublicEvolvingpublic <R> SingleOutputStreamOperator<R> transform(String operatorName, TypeInformation<R> outTypeInfo, OneInputStreamOperator<T, R> operator) {return this.doTransform(operatorName, outTypeInfo, SimpleOperatorFactory.of(operator));}
keyBy
public <K> KeyedStream<T, K> keyBy(KeySelector<T, K> key, TypeInformation<K> keyType) {Preconditions.checkNotNull(key);Preconditions.checkNotNull(keyType);return new KeyedStream(this, (KeySelector)this.clean(key), keyType);}
状态后端
public ValueStateDescriptor(String name, TypeInformation<T> typeInfo) {super(name, typeInfo, (Object)null);}
自定义数据类型&自定义序列化器
POJO类
@Datapublic class BroadcastConfig implements Serializable {public String config_type;public String date;public String media_id;public String account_id;public String label_id;public long start_time;public long end_time;public int interval;public String msg;public BroadcastConfig() {}}
HashMap<String, TypeInformation<?>> pojoFieldName = new HashMap<>();pojoFieldName.put("config_type", Types.STRING);pojoFieldName.put("date", Types.STRING);pojoFieldName.put("media_id", Types.STRING);pojoFieldName.put("account_id", Types.STRING);pojoFieldName.put("label_id", Types.STRING);pojoFieldName.put("start_time", Types.LONG);pojoFieldName.put("end_time", Types.LONG);pojoFieldName.put("interval", Types.INT);pojoFieldName.put("msg", Types.STRING);return Types.POJO(BroadcastConfig.class,pojoFieldName);
自定义TypeInformation
public class BinaryRowDataTypeInfo extends TypeInformation<BinaryRowData> {private static final long serialVersionUID = 4786289562505208256L;private final int numFields;private final Class<BinaryRowData> clazz;private final TypeSerializer<BinaryRowData> serializer;public BinaryRowDataTypeInfo(int numFields) {this.numFields=numFields;this.clazz=BinaryRowData.class;serializer= new BinaryRowDataSerializer(numFields);}@Overridepublic boolean isBasicType() {return false;}@Overridepublic boolean isTupleType() {return false;}@Overridepublic int getArity() {return numFields;}@Overridepublic int getTotalFields() {return numFields;}@Overridepublic Class<BinaryRowData> getTypeClass() {return this.clazz;}@Overridepublic boolean isKeyType() {return false;}@Overridepublic TypeSerializer<BinaryRowData> createSerializer(ExecutionConfig config) {return serializer;}@Overridepublic String toString() {return "BinaryRowDataTypeInfo<" + clazz.getCanonicalName() + ">";}@Overridepublic boolean equals(Object obj) {if (obj instanceof BinaryRowDataTypeInfo) {BinaryRowDataTypeInfo that = (BinaryRowDataTypeInfo) obj;return that.canEqual(this)&& this.numFields==that.numFields;} else {return false;}}@Overridepublic int hashCode() {return Objects.hash(this.clazz,serializer.hashCode());}@Overridepublic boolean canEqual(Object obj) {return obj instanceof BinaryRowDataTypeInfo;}}
自定义TypeSerializer
public class Roaring64BitmapTypeSerializer extends TypeSerializer<Roaring64Bitmap> {/*** Sharable instance of the Roaring64BitmapTypeSerializer.*/public static final Roaring64BitmapTypeSerializer INSTANCE = new Roaring64BitmapTypeSerializer();private static final long serialVersionUID = -8544079063839253971L;@Overridepublic boolean isImmutableType() {return false;}@Overridepublic TypeSerializer<Roaring64Bitmap> duplicate() {return this;}@Overridepublic Roaring64Bitmap createInstance() {return new Roaring64Bitmap();}@Overridepublic Roaring64Bitmap copy(Roaring64Bitmap from) {Roaring64Bitmap copiedMap = new Roaring64Bitmap();from.forEach(copiedMap::addLong);return copiedMap;}@Overridepublic Roaring64Bitmap copy(Roaring64Bitmap from, Roaring64Bitmap reuse) {from.forEach(reuse::addLong);return reuse;}@Overridepublic int getLength() {return -1;}@Overridepublic void serialize(Roaring64Bitmap record, DataOutputView target) throws IOException {record.serialize(target);}@Overridepublic Roaring64Bitmap deserialize(DataInputView source) throws IOException {Roaring64Bitmap navigableMap = new Roaring64Bitmap();navigableMap.deserialize(source);return navigableMap;}@Overridepublic Roaring64Bitmap deserialize(Roaring64Bitmap reuse, DataInputView source) throws IOException {reuse.deserialize(source);return reuse;}@Overridepublic void copy(DataInputView source, DataOutputView target) throws IOException {Roaring64Bitmap deserialize = this.deserialize(source);copy(deserialize);}@Overridepublic boolean equals(Object obj) {if (obj == this) {return true;} else if (obj != null && obj.getClass() == Roaring64BitmapTypeSerializer.class) {return true;} else {return false;}}@Overridepublic int hashCode() {return this.getClass().hashCode();}@Overridepublic TypeSerializerSnapshot<Roaring64Bitmap> snapshotConfiguration() {return new Roaring64BitmapTypeSerializer.Roaring64BitmapSerializerSnapshot();}public static final class Roaring64BitmapSerializerSnapshotextends SimpleTypeSerializerSnapshot<Roaring64Bitmap> {public Roaring64BitmapSerializerSnapshot() {super(() -> Roaring64BitmapTypeSerializer.INSTANCE);}}}
Flink SQL自定义函数
如何显式指定数据类型
@Overridepublic Optional<RexNode> convert(CallExpression call, ConvertContext context) {FunctionDefinition functionDefinition = call.getFunctionDefinition();// built-in functions without implementation are handled separatelyif (functionDefinition instanceof BuiltInFunctionDefinition) {final BuiltInFunctionDefinition builtInFunction =(BuiltInFunctionDefinition) functionDefinition;if (!builtInFunction.getRuntimeClass().isPresent()) {return Optional.empty();}}TypeInference typeInference =functionDefinition.getTypeInference(context.getDataTypeFactory());if (typeInference.getOutputTypeStrategy() == TypeStrategies.MISSING) {return Optional.empty();}switch (functionDefinition.getKind()) {case SCALAR:case TABLE:List<RexNode> args =call.getChildren().stream().map(context::toRexNode).collect(Collectors.toList());final BridgingSqlFunction sqlFunction =BridgingSqlFunction.of(context.getDataTypeFactory(),context.getTypeFactory(),SqlKind.OTHER_FUNCTION,call.getFunctionIdentifier().orElse(null),functionDefinition,typeInference);return Optional.of(context.getRelBuilder().call(sqlFunction, args));default:return Optional.empty();}}
指定accumulatorType
这是之前写的AbstractLastValueWithRetractAggFunction功能主要是为了实现具有local-global的逻辑的LastValue,提升作业性能。
accumulator对象:LastValueWithRetractAccumulator,可以看到该对象是一个非常复杂的对象,包含5个属性,还有List<Tuple2> 复杂嵌套,以及MapView等可以操作状态后端的对象,甚至有Object这种通用的对象。
public static class LastValueWithRetractAccumulator {public Object lastValue = null;public Long lastOrder = null;public List<Tuple2<Object, Long>> retractList = new ArrayList<>();public MapView<Object, List<Long>> valueToOrderMap = new MapView<>();public MapView<Long, List<Object>> orderToValueMap = new MapView<>();@Overridepublic boolean equals(Object o) {if (this == o) {return true;}if (!(o instanceof LastValueWithRetractAccumulator)) {return false;}LastValueWithRetractAccumulator that = (LastValueWithRetractAccumulator) o;return Objects.equals(lastValue, that.lastValue)&& Objects.equals(lastOrder, that.lastOrder)&& Objects.equals(retractList, that.retractList)&& valueToOrderMap.equals(that.valueToOrderMap)&& orderToValueMap.equals(that.orderToValueMap);}@Overridepublic int hashCode() {return Objects.hash(lastValue, lastOrder, valueToOrderMap, orderToValueMap, retractList);}}
public TypeInference getTypeInference(DataTypeFactory typeFactory) {return TypeInference.newBuilder().accumulatorTypeStrategy(callContext -> {List<DataType> dataTypes = callContext.getArgumentDataTypes();DataType argDataType;if (dataTypes.get(0).getLogicalType().getTypeRoot().getFamilies().contains(LogicalTypeFamily.CHARACTER_STRING)) {argDataType = DataTypes.STRING();} elseargDataType = DataTypeUtils.toInternalDataType(dataTypes.get(0));DataType accDataType = DataTypes.STRUCTURED(LastValueWithRetractAccumulator.class,DataTypes.FIELD("lastValue", argDataType.nullable()),DataTypes.FIELD("lastOrder", DataTypes.BIGINT()),DataTypes.FIELD("retractList", DataTypes.ARRAY(DataTypes.STRUCTURED(Tuple2.class,DataTypes.FIELD("f0", argDataType.nullable()),DataTypes.FIELD("f1", DataTypes.BIGINT()))).bridgedTo(List.class)),DataTypes.FIELD("valueToOrderMap",MapView.newMapViewDataType(argDataType.nullable(),DataTypes.ARRAY(DataTypes.BIGINT()).bridgedTo(List.class))),//todo:blink 使用SortedMapView 优化性能,开源使用MapView key天然字典升序,倒序遍历性能可能不佳DataTypes.FIELD("orderToValueMap",MapView.newMapViewDataType(DataTypes.BIGINT(),DataTypes.ARRAY(argDataType.nullable()).bridgedTo(List.class))));return Optional.of(accDataType);}).build();}
指定outputType
@Overridepublic TypeInference getTypeInference(DataTypeFactory typeFactory) {return TypeInference.newBuilder().outputTypeStrategy(callContext -> {List<DataType> dataTypes = callContext.getArgumentDataTypes();DataType argDataType;if (dataTypes.get(0).getLogicalType().getTypeRoot().getFamilies().contains(LogicalTypeFamily.CHARACTER_STRING)) {argDataType = DataTypes.STRING();} elseargDataType = DataTypeUtils.toInternalDataType(dataTypes.get(0));return Optional.of(argDataType);}).build();}
指定intputType
根据inputType动态调整outType或者accumulatorType
.outputTypeStrategy(callContext -> {List<DataType> dataTypes = callContext.getArgumentDataTypes();DataType argDataType;if (dataTypes.get(0).getLogicalType().getTypeRoot().getFamilies().contains(LogicalTypeFamily.CHARACTER_STRING)) {argDataType = DataTypes.STRING();} elseargDataType = DataTypeUtils.toInternalDataType(dataTypes.get(0));return Optional.of(argDataType);}
自定义DataType
/*** Data type of an arbitrary serialized type. This type is a black box within the table* ecosystem and is only deserialized at the edges.** <p>The raw type is an extension to the SQL standard.** <p>This method assumes that a {@link TypeSerializer} instance is present. Use {@link* #RAW(Class)} for automatically generating a serializer.** @param clazz originating value class* @param serializer type serializer* @see RawType*/public static <T> DataType RAW(Class<T> clazz, TypeSerializer<T> serializer) {return new AtomicDataType(new RawType<>(clazz, serializer));}
public TypeInference getTypeInference(DataTypeFactory typeFactory) {return TypeInference.newBuilder().accumulatorTypeStrategy(callContext -> {DataType type = DataTypes.RAW(Roaring64Bitmap.class,Roaring64BitmapTypeSerializer.INSTANCE);return Optional.of(type);}).outputTypeStrategy(callContext -> Optional.of(DataTypes.BIGINT())).build();}
四
结语
往期回顾
1. 深入理解Babel - 项目管理工具lerna解析|得物技术
2. 可视化流量录制规则探索和实践|得物技术
3. 在得物的小程序生态实践
4. 客服测试流水线编排设计思路和准入准出应用|得物技术
5. 深入剖析时序Prophet模型:工作原理与源码解析|得物技术
文 / 木木
关注得物技术,每周一、三、五更新技术干货
要是觉得文章对你有帮助的话,欢迎评论转发点赞~
未经得物技术许可严禁转载,否则依法追究法律责任。
“
扫码添加小助手微信
如有任何疑问,或想要了解更多技术资讯,请添加小助手微信: