0%

Flink:数据类型与序列化

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);
// TypeInformation<Pojo1> information = TypeInformation.of(Pojo1.class);

//Generic types have been disabled in the ExecutionConfig and type java.util.List is treated as a generic type.
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 List<String> names;
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);
// TypeInformation<Pojo1> information = TypeInformation.of(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就可以达到排序的效果,真实的数据不用移动。

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作为类型声明的统一入口,基本涵盖了常用类型。

类型擦除带来的问题

  1. Lambda函数的类型提取因为类型擦除导致Lambda函数的类型提取并不能总是有效的,有时候需要手动指定类型。
  2. 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中的通用序列化方式,缺点是序列化/反序列化效率相对较低。