[toc]
接上文分析,要将timer改成基于rocksdb,其实就是要对存储timer的set和queue提供基于rocksdb的存储方案。以下我们基于flink1.7版本源码分析
registerProcessingTimeTimer
1 | public void registerProcessingTimeTimer(N namespace, long time) { |
可以和看到1.4版本中的基本逻辑是一致的,只是存储方式变化了,下面我们就来分析一下新的存储方式是怎么实现的。在存储的选择上依然有Heap和RocksDB两个方式
1 | switch (priorityQueueStateType) { |
从timerService的需求来看我们可以看到这样的几个需求:
- 能够每次poll出最近需要触发的timer,实际上是需要维护一个小顶堆
- 能够对每一个key的timer去重
针对这两个需求,总体来看基于Heap的实现是通过基于数组实现了一个二叉堆,具体实现类为HeapPriorityQueue, 然后针对去重的功能又继承该PQ,通过一个hashmap数组,数组的每一个元素代表一个KG的一组不重复timer,同时这组timer内部也维护了timer在二叉堆中存储的下标,方便deleteTimer时的快速删除。
基于Heap的实现
HeapPriorityQueueElement,AbstractHeapPriorityQueue,HeapPriorityQueue,HeapPriorityQueueSet
- 存储基于数组,通过
HeapPriorityQueueElement记录自己所在的index,可以达到快速删除的目的 - 数组中存储的是一个二叉树,数组的起始位置是从1开始,为了使一些热点方法做更少的计算
在实现上是采用template的设计模式,主要实现逻辑交由子类来实现:
- addInternal
- removeInternal
- getHeadElementIndex
addInternal
1 | public boolean add(@Nonnull T toAdd) { |
添加一个timer至数组中,返回值false表示队首的元素没有改变,true则表示改变了或者不确定
1 | protected void addInternal(@Nonnull T element) { |
1 | private int increaseSizeByOne() { |
1 | // 将新加入的元素存储到相应的idx处,并且记录该元素在queue中的位置 |
1 | private void siftUp(int idx) { |
1 | // 比较两个值的优先级 |
removeInternal
- 抽取第一个timer用以触发
- 用户删除某个timer的行为
删除的方式是通过idx下标来实现快速删除的,这也就是HeapPriorityQueueElement中记录idx的作用
1 | protected T removeInternal(int removeIdx) { |
1 | private void adjustElementAtIndex(T element, int index) { |
1 | private void siftDown(int idx) { |
以上的操作大致上是一个二叉堆的增删的调整过程,涉及的具体算法可以查阅下文末的资料。
以下来分析rocksdb存储的实现
KeyGroupPartitionedPriorityQueue
基于rocksdb的存储是通过这个KeyGroupPartitionedPriorityQueue类来实现的,这个类中通过一个内存优先级队列,也就是上文中提到的内部实现的HeapPriorityQueue,用以存储所有KG的timer,而每一个分组的timer是如何存储呢?
1 | for (int i = 0; i < keyGroupedHeaps.length; i++) { |
在这里的构造函数可以看到,其实是通过orderedCacheFactory,从字面意思看是一个有序的缓存,也就是为每一个KG创建一个有序的缓存类,并将其添加到优先级队列中,这里的subHeap也是一个可比较的类,相当于去取这两个subHeap的堆顶的元素拿出来比较下就可以知道这两个subheap的排序方式了。
比如poll的逻辑,首先先从HeapPQ中挑出堆顶(一个subPQ),然后再从这个PQ中取出堆顶就是要触发的timer了,而这个subPQ就是真是数据(timer)存储的地方了。
1 | public T poll() { |
RocksDBCachingPriorityQueueSet
这个是上节中subHeap的实现类
1 | private void checkRefillCacheFromStore() { |
1 | public E peek() { |
1 | public E poll() { |
1 | public boolean add(@Nonnull E toAdd) { |
1 | public boolean remove(@Nonnull E toRemove) { |
Iterator中seekHint的作用:
1 | private RocksBytesIterator(@Nonnull RocksIteratorWrapper iterator) { |
再说checkpoint
在doc中作者也提到之所以要做这个feature除了因为timer过多会导致OOM等问题,还有一个原因是因为timer的属性虽然和keyed state很类似,但是代码管理以及checkpoint的方式都是单独的一块逻辑,并且checkpoint的持久化过程还是同步的(因为是以raw keyedstate的方式去进行的),再修改之后,每次注册的timeservice都会注册到kvstatInfo中,将checkpoint的逻辑统一到statebackend中并且实现了异步化。
https://blog.csdn.net/u010224394/article/details/8834969
https://github.com/apache/flink/pull/6159
https://docs.google.com/document/d/1XbhJRbig5c5Ftd77d0mKND1bePyTC26Pz04EvxdA7Jc/edit