【流式数据处理】键控状态与 State TTL
内容提要
本文讨论了Flink中的Keyed State及其管理,介绍了五种状态类型(ValueState、ListState、MapState、ReducingState、AggregatingState)的适用场景和特点。同时分析了HashMapStateBackend与EmbeddedRocksDBStateBackend的优缺点,以及状态的TTL配置和清理策略。最后,提供了状态大小估算的方法,指出滑动窗口和无限MapState可能导致的状态膨胀问题。
关键要点
-
Keyed State 绑定在 keyBy 之后,KeyGroup 个数等于 maxParallelism。
-
五种状态类型:ValueState、ListState、MapState、ReducingState、AggregatingState,各自适用场景不同。
-
HashMapStateBackend 和 EmbeddedRocksDBStateBackend 的优缺点:HashMap 适合小状态,RocksDB 适合大状态和增量 checkpoint。
-
状态的 TTL 配置和清理策略影响 checkpoint 和 RocksDB compaction 的效果。
-
状态大小估算方法:活跃 key 数、并发窗口数和单条状态序列化后字节数的乘积。
-
滑动窗口和无限 MapState 是状态膨胀的主要来源,需注意管理。
延伸解读
状态类型的选择与应用场景
Flink 提供的五种 Keyed State 类型各有其适用场景。ValueState 适合简单计数,而 ListState 则适合会话缓冲,但需注意列表长度无界可能导致状态膨胀。MapState 适合动态维度的复杂场景,ReducingState 和 AggregatingState 则用于聚合计算。选择合适的状态类型可以有效提升流式处理的性能与资源利用率。
状态膨胀的管理策略
滑动窗口和无限 MapState 是导致状态膨胀的主要原因。为了避免状态膨胀,建议在设计时合理配置 TTL 和清理策略,定期评估状态大小,并根据实际情况调整并行度。此外,使用状态大小估算公式可以帮助开发者在初期阶段预判资源需求,避免后期的性能问题。
选择合适的 StateBackend
在选择 HashMapStateBackend 和 EmbeddedRocksDBStateBackend 时,需要考虑状态的大小和性能需求。HashMap 适合小状态且延迟低,而 RocksDB 则适合大状态和增量 checkpoint。开发者应根据具体的应用场景和资源限制,选择合适的 StateBackend,以确保系统的稳定性和高效性。
延伸问答
Flink中的Keyed State有哪些类型?
Flink中的Keyed State包括ValueState、ListState、MapState、ReducingState和AggregatingState。
HashMapStateBackend和EmbeddedRocksDBStateBackend的优缺点是什么?
HashMapStateBackend适合小状态,读写延迟低,但大状态会导致OOM;EmbeddedRocksDBStateBackend适合大状态和增量checkpoint,但读写延迟较高,调参复杂。
如何配置Keyed State的TTL?
可以使用StateTtlConfig配置TTL,包括设置过期时间、更新类型和可见性等参数。
状态膨胀的主要来源是什么?
状态膨胀的主要来源是滑动窗口和无限MapState。
如何估算Flink中的状态大小?
状态大小可以通过活跃key数、并发窗口数和单条状态序列化后字节数的乘积来估算。
Keyed State的TTL清理策略有哪些?
Keyed State的TTL清理策略包括增量清理、全量快照时清理和RocksDB compaction过滤器。