【发布时间】:2016-11-18 18:10:17
【问题描述】:
我一直在 Spark Streaming 中使用 mapWithState API,但是关于 StateSpec.function 有两点不清楚:
假设我的功能是:
def trackStateForKey(batchTime: Time,
key: Long,
newValue: Option[JobData],
currentState: State[JobData]): Option[(Long, JobData)]
为什么新值是
Option[T]类型?据我所见,它总是为我定义的,而且由于该方法应该以新状态调用,我真的不明白为什么它可以是可选的。-
返回值是什么意思?我试图在文档和源代码中找到一些指针,但没有一个描述它的用途。既然我正在使用
state.remove()和state.update()修改键的状态,为什么我必须对返回值做同样的事情?在我当前的实现中,如果我删除密钥,我会返回
None,如果我更新它,我会返回Some(newState),但我不确定这是否正确。
【问题讨论】:
标签: scala apache-spark spark-streaming