mapGroupsWithState中eventTimeTimeout与processingTimeTimeout区别及状态过期设置
1.
mapGroupsWithState 两种超时策略的核心区别 二者的本质差异是判断超时的时间锚点不同:
processingTimeTimeout:以执行计算的节点本地系统时间为锚点。只要当前系统时间与该分组键最后一次被处理的时间差达到超时阈值,就会触发超时清理,完全不感知数据本身携带的业务事件时间,也不受水位线(watermark)影响。eventTimeTimeout:以数据携带的业务事件时间为锚点,依赖流的水位线推进触发超时。只有当全局水位线超过「该分组键最后一次接收数据的事件时间 + 超时阈值」时,才会触发超时清理,可适配处理迟到数据的场景。
举个实际例子方便理解:假设超时阈值设为10分钟,某分组键最后一次接收的数据事件时间为12:00,当前全局水位线为11:55,当前节点系统时间为12:11。此时processingTimeTimeout会判定该键超时直接清状态,eventTimeTimeout则不会触发超时,该键的状态依然可以被后续到达的、事件时间符合水位线要求的迟到数据更新。
2. 10分钟状态过期的实现方案
完全可以实现你要求的效果,两种超时策略都可以满足,根据你的业务场景选择即可:
基于事件时间的10分钟过期
适用于你希望按数据实际产生时间判断过期的场景:
- 首先给输入流设置水位线,示例逻辑:
ds.withWatermark("event_time_col", "0 second"),第二个参数为允许的额外迟到时间,可根据业务调整。 - 在
mapGroupsWithState中指定超时策略为GroupStateTimeout.EventTimeTimeout()。 - 每次处理该分组键的数据时,调用
state.setTimeoutTimestamp(当前数据事件时间戳 + 10 * 60 * 1000)设置超时时间。
当水位线超过你设置的超时时间戳后,该键的状态会被标记为超时,你可以在处理逻辑中通过state.hasTimedOut()判断后清理状态,后续晚到10分钟以上的同键数据会被识别为新分组,状态从头开始维护。
基于处理时间的10分钟过期
适用于你希望按数据到达系统的时间判断过期的场景:
- 不需要设置水位线,直接在
mapGroupsWithState中指定超时策略为GroupStateTimeout.ProcessingTimeTimeout()。 - 每次处理该分组键的数据时,调用
state.setTimeoutDuration(10 * 60 * 1000)设置超时时长。
只要该键10分钟内没有新数据进入处理流程,状态就会被触发超时清理,后续到达的同键数据从头维护状态。
内容的提问来源于stack exchange,提问作者Ayushi Dewangan
相关产品推荐
相关产品推荐

