mapGroupsWithState中setTimeoutTimestamp()的作用及状态过期规则相关问题
Spark 流式计算mapGroupsWithState相关问题解答
问题1:mapGroupsWithState()中的setTimeoutTimestamp()具体工作原理是什么?
setTimeoutTimestamp()是Structured Streaming中mapGroupsWithState算子配套的事件时间维度超时控制接口,核心工作逻辑如下:
- 该方法针对单个分组状态生效,调用时需要传入一个毫秒级时间戳参数,可以理解为给当前分组设置了一个超时的时间节点:如果到了这个时间节点之后,该分组还没有新的数据流入,就会被标记为超时。
- 超时判定是在每批次数据处理完成、全局水位线(watermark)更新后统一触发的:Spark会对比每个分组设置的超时时间戳和当前最新的全局水位线,只要分组的超时时间戳 ≤ 当前水位线,就会判定该分组超时。
- 每当分组有新数据流入处理时,你可以重新调用
setTimeoutTimestamp()更新超时时间,相当于给活跃的分组延长状态有效期,避免正常业务的状态被误判超时。
问题2:当运行时长超出setTimeoutTimestamp()指定的时长后,是否会自动删除对应的状态?
不会默认自动删除,具体逻辑如下:
- 分组被判定超时后,Spark会主动触发你在
mapGroupsWithState中定义的超时处理分支,状态的清理需要你在该分支中主动调用state.remove()方法执行。 - 如果你不在超时分支中手动删除状态,哪怕分组已经超时,对应的状态会一直保存在状态存储中,持续占用存储资源,直到任务重启(非持久化状态场景)或者你主动触发清理逻辑。
可以参考如下代码示例理解处理逻辑:
// 初始化时指定使用事件时间超时策略 mapGroupsWithState(GroupStateTimeout.EventTimeTimeout()) { (key: String, events: Iterator[BizEvent], state: GroupState[GroupStateData]) => if (state.hasTimedOut) { // 超时处理分支,主动删除状态 val stateData = state.get state.remove() // 处理超时后的业务逻辑,返回对应结果 buildTimeoutResult(key, stateData) } else { // 正常处理新流入的数据 val eventList = events.toList val maxEventTime = eventList.map(_.eventTimeMs).max // 更新当前分组的超时时间:最大事件时间+2小时 state.setTimeoutTimestamp(maxEventTime + 2 * 3600 * 1000) // 更新状态逻辑 val newState = updateState(state.getOption, eventList) state.update(newState) buildNormalResult(key, newState) } }
内容的提问来源于stack exchange,提问作者Ayushi Dewangan
相关产品推荐
相关产品推荐

