What is State
有状态计算 VS 无状态计算
无状态计算指的是数据进入Flink后经过算子时只需要对当前数据进行处理就能得到想要的结果
有状态计算就是需要和历史的一些状态或进行相关操作,才能计算出正确的结果
比如去重、滑动窗口
状态是计算过程中的数据信息,在容错恢复和 Checkpoint 中有重要的作用,流计算在本质上是 Incremental Processing,因此需要不断查询保持状态;
为了确保 Exactly- once 语义,需要数据能够写入到状态中;而持久化存储,能够保证在整个分布式系统运行失败或者挂掉的情况下做到 Exactly- once,这是状态的另外一个价值。
状态管理的目标:易用、高效、可靠
Kinds of state in Flink<br>状态类型
Operator State
可以用于所有的算子
一个operator对应一个state
并发改变时需要选择分配方式,内置:1.均匀分配 2.所有state合并后再分发给每个实例
需要你实现CheckPointedFunction或ListCheckPointed接口
只支持 List state、Union List state、Broadcast state
Keyed State
只能应用在KeyedSteam上
每个key 对应一个 state,一个operator处理多个key,会访问相应的多个state
并发改变时,state随着key在实例间迁移
通过RuntimeContext访问,需要operator是一个richFunction
支持ValuedState、ListState、Reducing State、MapState、Aggregating State
Fault Tolerance<br>状态容错
Check pointing<br>备份
Checkpoint 是 Flink 实现容错机制的核心,它周期性的记录计算过程中 Operator 的状态,并生成快照持久化存储,备份至远程的分布式系统中。
Barriers
Exactly Once & At Least Once
Asynchronous State Snapshots<br>异步check point
Checkpointing Algorithm(Chandy-Lamport)<br>checkpoint算法
Incremental Checkpointing
Save points
当 Flink 作业发生故障崩溃时,可以有选择的从 Checkpoint 中恢复,保证了计算的一致性。
State Processor API(read, write, and modify savepoints and checkpoints using Flink's batch DataSet API)
ValueState 单个值 update/get
MapState Map put/putAll/remove/contains/entries/iterator/keys/values
ListState List add/addAll/update/get
ReducingState 单个值 add/addAll/update/get
AggregatingState 单个值 add IN类型,get Out 类型