【问题标题】:How to set a TTL for all items in a MapState in Flink?如何为 Flink 中的 MapState 中的所有项目设置 TTL?
【发布时间】:2021-02-19 06:58:23
【问题描述】:
我想在给定的时间戳清除 MapState 中的所有条目。
我正在考虑两种方法:
- 将
cleanup timestamp 保存在ValueState 中,为cleanup timestamp 注册一个定时器,当定时器触发时清除MapState。尽管cleanup timestamp 可能相同,但添加到MapState 的每个项目都会发生这种情况。我依靠 Flink 对计时器进行重复数据删除。
- 根据 (
cleanup timestamp - current timestamp) 计算 TTL,使用 StateTtlConfig 为 MapState 设置 TTL
哪种方法更好(性能、准确性等)?
StateTtlConfig 是否适用于均匀时间处理?
【问题讨论】:
标签:
join
streaming
apache-flink
flink-streaming
【解决方案1】:
如果您打算同时清除MapState 中的所有条目,那么我不会使用StateTtlConfig,因为Flink 将花费8 个字节来存储每个映射条目的计时器。这是很多不必要的存储开销。
使用StateTtlConfig,只能根据处理时间来指定状态到期。
另外请记住,StateTtlConfig 不能添加到现有状态描述符或从现有状态描述符中删除。
【解决方案2】:
如果您想在给定的时间戳清除所有记录,那么选项 1 应该更具有明确的性能。一般来说,在 flink 状态上使用 TTL 会增加额外的开销,因为每次访问 MapState 中的 key 时都会检查它。这取决于您如何准确使用该状态(您计划在那里存储多少记录,您访问它的频率以及状态的类型),但我会说,如果您只想将所有记录删除给定时间戳,那么选项 1 是更好的主意。
【解决方案3】:
我认为使用 timerService 来注册一个计时器会给你带来更好的性能。但是事件时间计时器只会在水印进入时触发,你也可以使用当前的来调度和合并这些计时器与下一个水印。
val coalescedTime = ctx.timerService.currentWatermark + 1
ctx.timerService.registerEventTimeTimer(coalescedTime)