【问题标题】:Why does memory usage of spark worker increases with time?为什么 spark worker 的内存使用量会随着时间增加?
【发布时间】:2018-06-19 19:21:25
【问题描述】:

我有一个 Spark Streaming 应用程序正在运行,它使用 mapWithState 函数来跟踪 RDD 的状态。 该应用程序可以正常运行几分钟,但随后会因

而崩溃
org.apache.spark.shuffle.MetadataFetchFailedException: Missing an output location for shuffle 373

我观察到,即使我为 mapWithStateRDD 设置了超时,Spark 应用程序的内存使用量也会随着时间线性增加。请看下面的代码 sn -p 和内存使用情况 -

val completedSess = sessionLines
                    .mapWithState(StateSpec.function(trackStateFunction _)
                    .numPartitions(80)
                    .timeout(Minutes(5)))

如果每个 RDD 都有明确的超时,为什么内存应该随着时间线性增加?

我试过增加内存,但没关系。我错过了什么?

编辑 - 参考代码

def trackStateFunction(batchTime: Time, key: String, value: Option[String], state: State[(Boolean, List[String], Long)]): Option[(Boolean, List[String])] = {

  def updateSessions(newLine: String): Option[(Boolean, List[String])] = {
    val currentTime = System.currentTimeMillis() / 1000

    if (state.exists()) {
      val newLines = state.get()._2 :+ newLine

      //check if end of Session reached.
      // if yes, remove the state and return. Else update the state
      if (isEndOfSessionReached(value.getOrElse(""), state.get()._4)) {
        state.remove()
        Some(true, newLines)
      }
      else {
        val newState = (false, newLines, currentTime)
        state.update(newState)
        Some(state.get()._1, state.get()._2)
      }
    }
    else  {
      val newState = (false, List(value.get), currentTime)
      state.update(newState)
      Some(state.get()._1, state.get()._2)
    }
  }

  value match {
    case Some(newLine) => updateSessions(newLine)
    case _ if state.isTimingOut() => Some(true, state.get()._2)
    case _ => {
      println("Not matched to any expression")
      None
    }
  }
}

【问题讨论】:

  • 您有多少传入流量?多少内存/磁盘?我们需要更多信息。
  • 另外,你多久检查一次?
  • 我有一个由 4 个工作人员组成的集群(8 个内核,32 GB RAM,每个 128 GB SSD)。来自 Kinesis Stream 的传入流量为 10-15 MB/s。批处理间隔为 10 秒。检查点间隔为60s
  • 你在状态中存储了多少数据(也许共享代码)?有什么东西在不应该的时候保持静态吗?
  • @YuvalItzchakov 用代码 sn-p 更新了问题。它的灵感来自您的博文 :)

标签: scala apache-spark spark-streaming


【解决方案1】:

根据mapwithstate的信息: 国家规范 作为 RDD 的初始状态 - 您可以从某个存储加载初始状态,然后使用该状态启动流式传输作业。

分区数 - 键值状态 dstream 由键分区。如果你之前对状态的大小有一个很好的估计,你可以提供分区的数量来对它进行相应的分区。

分区器 - 您还可以提供自定义分区器。默认分区器是哈希分区器。如果你对 key space 有很好的了解,那么你可以提供一个自定义的 partitioner,它可以比默认的 hash partitioner 进行更高效的更新。

超时 - 这将确保其值在特定时间段内未更新的键将从状态中删除。这有助于清理旧密钥的状态。

因此,超时仅与一段时间后清理未更新的密钥有关。内存将运行满并最终阻塞,因为执行程序没有分配足够的内存。这给出了 MetaDataFetchFailed 异常。随着内存的增加,我希望你的意思是执行者。即使那样,增加执行程序的内存也可能不起作用,因为流仍在继续。使用 MapWithState 会话行将包含与输入 dstream 相同的记录数。所以解决这个问题就是让你的 dstream 更小。在流式上下文中,您可以设置一个批处理间隔,这很可能解决这个问题

val ssc = new StreamingContext(sc, Seconds(batchIntervalSeconds))

记得偶尔也做一个快照和一个检查点。快照将允许您将现在更早丢失的流中的信息用于其他计算。希望这有助于了解更多信息,请参阅:https://docs.cloud.databricks.com/docs/spark/1.6/examples/Streaming%20mapWithState.htmlhttp://asyncified.io/2016/07/31/exploring-stateful-streaming-with-apache-spark/

【讨论】:

    【解决方案2】:

    mapWithState 还将 mappedValues 存储在 RAM 中(请参阅 MapWithStateRDD),默认情况下 mapWithState 会在 RAM 中存储多达 20 个 MapWithStateRDD。

    简而言之,RAM 使用量与批处理间隔成正比,

    您可以尝试减少批处理间隔以减少 RAM 使用量。

    【讨论】:

      猜你喜欢
      • 2023-02-01
      • 1970-01-01
      • 1970-01-01
      • 2018-09-03
      • 1970-01-01
      • 2021-11-22
      • 1970-01-01
      • 2013-02-14
      相关资源
      最近更新 更多