【问题标题】:Spark mapWithState API explanationSpark mapWithState API 解释
【发布时间】: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)]
  1. 为什么新值是Option[T] 类型?据我所见,它总是为我定义的,而且由于该方法应该以新状态调用,我真的不明白为什么它可以是可选的。

  2. 返回值是什么意思?我试图在文档和源代码中找到一些指针,但没有一个描述它的用途。既然我正在使用 state.remove()state.update() 修改键的状态,为什么我必须对返回值做同样的事情?

    在我当前的实现中,如果我删除密钥,我会返回None,如果我更新它,我会返回Some(newState),但我不确定这是否正确。

【问题讨论】:

    标签: scala apache-spark spark-streaming


    【解决方案1】:

    为什么新值是Option[T] 类型?据我所见,它是 总是为我定义,因为该方法应该被调用 有了一个新的状态,我真的不明白为什么会这样 可选。

    这是一个Option[T],因为如果您使用StateSpec.timeout 设置超时,例如:

    StateSpec.function(spec _).timeout(Milliseconds(5000))
    

    那么一旦函数超时传入的值将是None 并且State[T] 上的isTimingOut 方法将返回true。这是有道理的,因为状态的超时并不意味着指定键的新值已经到达,并且通常比将null 传递给T 更安全(无论如何这对原语都不起作用)作为您希望用户在Option[T] 上安全操作。

    您可以在 Sparks 实现中看到这一点:

    // Get the timed out state records, call the mapping function on each and collect the
    // data returned
    if (removeTimedoutData && timeoutThresholdTime.isDefined) {
      newStateMap.getByTime(timeoutThresholdTime.get).foreach { case (key, state, _) =>
        wrappedState.wrapTimingOutState(state)
        val returned = mappingFunction(batchTime, key, None, wrappedState) // <-- This.
        mappedData ++= returned
        newStateMap.remove(key)
      }
    }
    

    返回值是什么意思?我试图在 文档和源代码,但没有一个描述它是什么 用于。因为我正在使用 state.remove() 修改键的状态 和 state.update(),为什么我必须对 return 做同样的事情 价值观?

    返回值是沿火花图传递中间状态的一种方式。例如,假设我想更新我的状态,但还要在我的管道中使用 intermediate 数据执行一些操作,例如:

    dStream
      .mapWithState(stateSpec)
      .map(optionIntermediateResult.map(_ * 2))
      .foreachRDD( /* other stuff */)
    

    该返回值正是使我能够继续对所述数据进行操作的原因。如果你不关心中间结果而只想要完整的状态,那么输出None 就可以了。

    编辑:

    我写了一个blog post(在这个问题之后),试图对 API 进行深入的解释。

    【讨论】:

    • 谢谢,这就解释了。
    • 博文链接失效
    猜你喜欢
    • 1970-01-01
    • 2018-01-24
    • 1970-01-01
    • 2016-08-01
    • 1970-01-01
    • 1970-01-01
    • 2016-12-24
    • 2016-11-25
    • 1970-01-01
    相关资源
    最近更新 更多