Flink Data Types & Serialization 使用case class的坑 1 2 3 case class Event (id: Int ) { val lb = new ListBuffer [Int ] }
1 2 13:00:43,342 INFO org.apache.flink.api.java.typeutils.TypeExtractor - class org.myorg.quickstart.Event does not contain a setter for field id 13:00:43,343 INFO org.apache.flink.api.java.typeutils.TypeExtractor - Class class org.myorg.quickstart.Event cannot be used as a POJO type because not all fields are valid POJO fields, and must be processed as GenericType. Please read the Flink documentation on "Data Types & Serialization" for details of the effect on performance.
提示信息:找不到setter,对于POJO类型必须所有的字段必须要有setter和getter 命名是case class啊
再看生产环境的例子: 折腾了一下午
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 case class Event (uin: String , sPid: String , applyId: String , bankType: Long , transactionId: String , amount: Long , createTime: Long , bizType: Long , modifyTime: String , equ: ListBuffer [Int ] ) { def this (uin: String , sPid: String , applyId: String , bankType: Long , transactionId: String , amount: Long , createTime: Long , bizType: Long , modifyTime: String , equ: List [Int ] ) = { this (uin, sPid, applyId, bankType, transactionId, amount, createTime, bizType, modifyTime) if (equ != null ) { equities.appendAll(equ) } } private var _active = false val equities: ListBuffer [Int ] = new ListBuffer [Int ] def addEquity (id: Int ): Event = { equities.append(id) this } def equityString (separator: String ): String = { equities.mkString(separator) } def setActive (): Event = { _active = true this } def setActive (active: String ): Event = { _active = "1" .equals(active) this } def isActive () = { _active } }
使用的是flink 1.6版本的,case class识别出来了,但是equities没有传递到下一个算子中,始终没有值
老老实实的修改成普通类
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 class Event (_uin: String , _sPid: String , _applyId: String , _bankType: Long , _transactionId: String , _amount: Long , _createTime: Long , _bizType: Long , _modifyTime: String ) extends Serializable { private var equities: ListBuffer [Int ] = new ListBuffer [Int ] private var uin: String = _uin private var sPid: String = _sPid private var applyId: String = _applyId private var bankType: Long = _bankType private var transactionId: String = _transactionId private var amount: Long = _amount private var createTime: Long = _createTime private var bizType: Long = _bizType private var modifyTime: String = _modifyTime private var _active = false def this (uin: String , sPid: String , applyId: String , bankType: Long , transactionId: String , amount: Long , createTime: Long , bizType: Long , modifyTime: String , equ: List [Int ] ) = { this (uin, sPid, applyId, bankType, transactionId, amount, createTime, bizType, modifyTime) if (equ != null ) { equities.appendAll(equ) } } def getUin : String = uin def setUin (Uin : String ): Unit = { this .uin = Uin } def getSPid : String = sPid def setSPid (SPid : String ): Unit = { this .sPid = SPid } def getApplyId : String = applyId def setApplyId (ApplyId : String ): Unit = { this .applyId = ApplyId } def getBankType : Long = bankType def setBankType (BankType : Long ): Unit = { this .bankType = BankType } def getTransactionId : String = transactionId def setTransactionId (TransactionId : String ): Unit = { this .transactionId = TransactionId } def getAmount : Long = amount def setAmount (Amount : Long ): Unit = { this .amount = Amount } def getCreateTime : Long = createTime def setCreateTime (CreateTime : Long ): Unit = { this .createTime = CreateTime } def getBizType : Long = bizType def setBizType (BizType : Long ): Unit = { this .bizType = BizType } def getModifyTime : String = modifyTime def setModifyTime (ModifyTime : String ): Unit = { this .modifyTime = ModifyTime } def addEquity (id: Int ): Event = { equities.append(id) this } def equityString (separator: String ): String = { equities.mkString(separator) } def setActive (): Event = { _active = true this } def setActive (active: String ): Event = { _active = "1" .equals(active) this } def isActive () = { _active } def rights (split: String ) = { RightEvent (this , equities.mkString(split)) } def getEquities = equities def setEquities (equities: ListBuffer [Int ]) = { this .equities = equities } }
可以了
初步估计,序列化除了问题
flink 类型和序列化机制 flink 支持的数据类型
Java Tuples 跟 Scala Case 类 Java POJOs 基础类型(Primitive Types : int/long/string/char/short/boolean 等) 普通的类(非POJO) Values Hadoop Writable 特殊类型(Scala : Either, Option, Try; Java : List, Map) flink 支持的序列化
Tuple Row Pojo Avro Protobuf (via Kryo) Thrift (via Kryo) Kryo
flink 序列化性能 flink serialization performance results
可以看到 flink 内置的 Tuple 跟 Row 性能最好, POJO 次之, 一般 Tuple 跟 POJO使用的频率最高, 但是只有POJO 跟 Avro 支持 Schema 升级 POJO 一不小心就可能回退到 Kryo POJO 第一个要求是符合 Java Bean 规范, 但是目前(1.12.2) 还不支持特殊的类型(List, Map等) 为了避免 POJO序列化回退, 开发过程中可以开启
env.getConfig().disableGenericTypes(); 当POJO 中包含List/Map 处理方式
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 public class Pojo1 { public int id; public List<String> names } public static void main (String[] args) { ExecutionConfig config = new ExecutionConfig (); config.disableGenericTypes(); TypeInformation <Pojo1> information = Types.POJO(Pojo1.class); TypeSerializer <Pojo1> serializer = information.createSerializer(config); }
因为POJO包含了不支持的List该序列化, names 字段序列化会回退到 Kryo序列化
解决方案
1: 把 List 换成 Array
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 public static class Pojo1 { public int id; public String[] names; } public static void main (String[] args) { ExecutionConfig config = new ExecutionConfig (); config.disableGenericTypes(); TypeInformation <Pojo1> information = Types.POJO(Pojo1.class); TypeSerializer <Pojo1> serializer = information.createSerializer(config); }
2: 指定 POJO 字段 TypeInformation
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 public static void main (String[] args) { ExecutionConfig config = new ExecutionConfig (); config.disableGenericTypes(); Map <String, TypeInformation<?>> map = new HashMap <>(); map.put("id" , Types.INT); map.put("names" , Types.LIST(Types.STRING)); TypeInformation <Pojo1> information = Types.POJO(Pojo1.class, map); TypeSerializer <Pojo1> serializer = information.createSerializer(config); } ``` 3 : 自定义 TypeInfoFactory```java @TypeInfo(MyPojo1Factory.class) public static class Pojo1 { public int id; public List<String> names; } public static class MyPojo1Factory extends TypeInfoFactory <Pojo1> { @Override public TypeInformation <Pojo1> createTypeInfo (Type t, Map<String, TypeInformation<?>> genericParameters) { Map <String, TypeInformation<?>> map = new HashMap <>(); map.put("id" , Types.INT); map.put("names" , Types.LIST(Types.STRING)); TypeInformation <Pojo1> information = Types.POJO(Pojo1.class, map); return information; } } public static void main (String[] args) { ExecutionConfig config = new ExecutionConfig (); config.disableGenericTypes(); TypeInformation <Pojo1> information = Types.POJO(Pojo1.class); TypeSerializer <Pojo1> serializer = information.createSerializer(config); }
引用https://ci.apache.org/projects/flink/flink-docs-release-1.12/zh/dev/types_serialization.html https://flink.apache.org/news/2020/04/15/flink-serialization-tuning-vol-1.html
Flink的序列化 Flink实现了自己的序列化框架,并结合自身的内存模型,实现了对象的密集存储也高效操作。
Flink序列化框架
可以看出这种序列化方式存储密度是相当紧凑的。其中 int 占4字节,double 占8字节,POJO多个一个字节的header,PojoSerializer只负责将header序列化进去,并委托每个字段对应的serializer对字段进行序列化。 memory pool 内存池 memorySegment的数据结构,由两部分组成,一部分是存储key+pointer(完整二进制数据的指针以及定长的序列化后的key),第二部分是对象的二进制数据 如下图:
使用内存池管理内存和使用二进制存储数据的的好处: 避免oom,所有的运行时数据结构和算法只能通过内存池申请内存,保证了其使用的内存大小是固定的,不会因为运行时数据结构和算法而发生OOM。在内存吃紧的情况下,算法(sort/join等)会高效地将一大批内存块写到磁盘,之后再读回来。因此,OutOfMemoryErrors可以有效地被避免。 节省内存空间,Java 对象在存储上有很多额外的消耗,使用二进制可以避免。 高效的二进制操作 & 缓存友好的计算,第一,交换定长块(key+pointer)更高效,不用交换真实的数据也不用移动其他key和pointer。第二,这样做是缓存友好的,因为key都是连续存储在内存中的,可以大大减少 cache miss(cpu读取L1,L2,L3高速缓存速度高于读取主内存速度几个数量级,使用key+pointer极大提高缓存L1,L2,L3命中率) 注意:Flink 中,排序会先用 key 比大小,这样就可以直接用二进制的key比较而不需要反序列化出整个对象。因为key是定长的,如果key相同(或者没有提供二进制key),那就必须将真实的二进制数据反序列化出来,然后再做比较。之后,只需要交换key+pointer就可以达到排序的效果,真实的数据不用移动。
Flink SQL为什么与DataStream使用不同的类型系统? 1)该类型系统与SQL的兼容性不好。 2)无法控制Decimal类型的精度。 3)无法区分char和varchar类型。 4)物理类型和逻辑类型紧耦合。 5)物理类型是类型描述,而不是类型的序列化/反序列化器。
Flink SQL引入了新的LogicalTypes类型系统 TypeInformation类型系统是为DataStream/DataSet API设计的
DataType有两个职责: 1)声明逻辑类型LogicalType。 2)运行时逻辑转换类,允许为空。
类型推断:
Java的Flink应用使用反射机制获取Function的输入和输出类型。Scala使用Scala Macro类提取类型。
类型提取:
泛型的类型推断:
Java的泛型机制是在编译级别实现的。编译器生成的字节码在运行期间并不包含泛型的类型信息 使用TypeHint的匿名类来获取泛型的类型信息
Lambda函数的类型提取
Eclipse的JDT编译器会把Lambda函数的泛型签名等信息写入编译后的字节码中,而对于javac等常见的其他编译器,则不会这样做,因而Flink就无法获取具体类型信息了
(1)Java类型擦除的原因
1)避免JVM的重构。如果JVM将泛型类型延续到运行期,那么到运行期时JVM就需要进行大量的重构工作,提高了运行期的效率。 2)版本兼容。在编译期擦除可以更好地支持原生类型(Raw Type)。
(2)Java泛型类型擦除规则 1)如果是继承基类而来的泛型,就用getGenericSuperclass(), 转型为ParameterizedType来获得实际类型。 2)如果是实现接口而来的泛型,就用getGenericInterfaces(), 针对其中的元素转型为ParameterizedType来获得实际类型。 3)Java泛型在字节码中会被擦除,并不总是擦除为Object类型,而是擦除到上限类型。
显示类型:
Flink提供了等价的Types类 (org.apache.flink.api.common.typeinfo.Types),Types作为类型声明的统一入口,基本涵盖了常用类型。
类型擦除带来的问题
Lambda函数的类型提取因为类型擦除导致Lambda函数的类型提取并不能总是有效的,有时候需要手动指定类型。
Kryo的JavaSerializer在Flink下存在Bug,可能导致ClassNotFound异常推荐使用org.apache.flink.api.java.typeutils.runtime.kryo.JavaSerializer,而非com.esot-ericsoftware.kryo.serializers.JavaSerializer,以防止与Flink不兼容
SQL类型系统 Flink SQL中则使用DataType中的LogicalType类型系统来描述类型信息 LogicalType类型系统与SQL标准基本保持一致,同时增加了一些额外的信息,如是否可以为null等,目的是提高scala expression(标量表达式)的处理效率。
Flink SQL执行时,最终转换为了FlinkDataStream/DataSet应用,此时就需要TypeInfomation类型信息来实现序列化/反序列化,所以SQL逻辑类型LogicalType需要转换为TypeInfomation
1)org.apache.flink.types.Row:在Flink Planner中使用,是1.9版本之前FlinkSQL使用的Row结构,在SQL相关的算子、UDF函数、代码生成中都是使用该套Row结构。 2)org.apache.flink.table.dataformat.BaseRow及其子类:是在Blink Runtime和Blink Planner中使用的新的Row类型数据结构,在Blink算子、UDF函数和代码生成中使用此结构。
Blink Row总览
ColumnarRow
ColumnarRow是一种内存列式存储结构,每一列的抽象结构为ColumnVector。在当前的实现中,只支持堆上ColumnVector,堆外的ColumnVector尚不被支持。堆上ColumnVector本质上是使用Java原始类型数据保存一列的数据。Orc类型的列式存储使用了ColumnarRow。对于查询类的请求,使用列式存储能够提高CPU缓存命中率。CPU的数据预读取策略总是尝试将相邻的数据预读取到缓存中,因为列式存储形式中一列数据总是紧邻的,与行式数据相比,访问同一个字段的时候,CPU缓存命中率更高,因此CPU就无须浪费宝贵的事件周期去等待数据从内存加载,从而提高计算效率,如图5-8所示。
序列化
MapFunction使用了匿名内部类的方式实现,默认内部类会持有一个外部对象的引用this$0,如果外部对象不实现序列化接口,内部类的序列化会失败,所在Flink中使用ASM操作字节码将匿名内部类中的this$0设置为null。在FlinkDataStreamp的map、filter、keyBy等接口中都使用了ClosureCleaner#clean方法来设置this$0。
如果开发者在编写Flink应用过程中使用了自定义类型,并且又没有提供类型的注册和序列化/反序列化方法,Flink就无法对该类型进行该自定义序列化/反序列化。此时为了Flink的正常运行,对于这一类的数据类型,无法识别的类型就会交给Kryo进行序列化。Kryo可以对任意类型的Java对象进行序列化,是一种Java中的通用序列化方式,缺点是序列化/反序列化效率相对较低。