【发布时间】:2018-07-25 12:23:55
【问题描述】:
Flink de/serialise operator state 的频率如何?每次获取/更新还是基于检查点?状态后端有影响吗?
我怀疑对于具有不同键(数百万)和每个键每秒数千个事件的键控流的情况,反序列化可能是一个大问题。我说的对吗?
【问题讨论】:
标签: apache-flink flink-streaming
Flink de/serialise operator state 的频率如何?每次获取/更新还是基于检查点?状态后端有影响吗?
我怀疑对于具有不同键(数百万)和每个键每秒数千个事件的键控流的情况,反序列化可能是一个大问题。我说的对吗?
【问题讨论】:
标签: apache-flink flink-streaming
你的假设是正确的。这取决于状态后端。
在 JVM 堆上存储状态的后端(MemoryStateBackend 和 FSStateBackend)不会为常规读/写访问序列化状态,而是将其作为对象保存在堆上。虽然这会导致访问速度非常快,但您显然受限于 JVM 堆的大小,并且还可能面临垃圾收集问题。采用检查点时,对象会被序列化并持久化,以便在发生故障时能够恢复。
相比之下,RocksDBStateBackend 将所有状态存储为嵌入式 RocksDB 实例中的字节数组。因此,它对每次读/写访问的密钥状态进行反序列化。您可以通过选择适当的状态原语来控制“多少”状态被序列化,即ValueState、ListState、MapState 等。
例如,ValueState 总是作为一个整体反序列化/序列化,而 MapState.get(key) 仅序列化键(用于查找)并反序列化键的返回值。因此,您应该使用MapState<String, String> 而不是ValueState<HashMap<String, String>>。类似的考虑也适用于其他状态原语。
RocksDBStateBackend 通过将文件复制到持久文件系统来检查其状态。因此,在采取检查点时不会涉及额外的序列化。
【讨论】: