【问题标题】:Long and consistent wait between tasks in spark streaming job火花流作业中任务之间的长时间且一致的等待
【发布时间】:2018-11-17 17:30:52
【问题描述】:

我有一个在 Mesos 上运行的 spark 流式传输作业。 它的所有批次都需要完全相同的时间,而且这个时间比预期的要长得多。 这些作业从 kafka 中提取数据,处理数据并将其插入到 cassandra 中,然后再次返回到 kafka 中,进入不同的主题。

每批(如下)有 3 个作业,其中 2 个从 kafka 拉取,处理并插入到 cassandra,另一个从 kafka 拉取,处理并推回 kafka。

我在 spark UI 中检查了批次,发现它们都花费了相同的时间(4 秒)但向下钻取更多,它们实际上每个处理不到一秒,但它们都有相同的时间间隔(大约 4秒)。 添加更多执行器或更多处理能力看起来不会有什么不同。

Details of batch: Processing time = 12s & total delay = 1.2 s ??

所以我深入研究了批处理的每个作业(它们都需要完全相同的时间 = 4 秒,即使它们正在执行不同的处理):

他们都需要 4 秒来运行他们的阶段之一(从 kafka 读取的阶段)。 现在我深入到其中一个的阶段(它们都非常相似):

为什么要等待?整个事情实际上只需要0.5s就可以运行,它只是在等待。是在等卡夫卡吗?

有没有人经历过类似的事情? 我可能编码错误或配置不正确?

编辑:

这是触发此行为的最小代码。这让我觉得它一定是某种设置。

object Test {

  def main(args: Array[String]) {

    val sparkConf = new SparkConf(true)
    val streamingContext = new StreamingContext(sparkConf, Seconds(5))

    val kafkaParams = Map[String, String](
      "bootstrap.servers" -> "####,####,####",
      "group.id" -> "test"
    )

    val stream = KafkaUtils.createDirectStream[String, Array[Byte], StringDecoder, DefaultDecoder](
      streamingContext, kafkaParams, Set("test_topic")
    )

    stream.map(t => "LEN=" + t._2.length).print()

    streamingContext.start()
    streamingContext.awaitTermination()
  }
}

即使所有的执行者都在同一个节点(spark.executor.cores=2 spark.cores.max=2),问题依然存在,和之前一样正好是4秒:One mesos executor

即使主题没有消息(0 条记录的批次),每批次的火花流式处理也需要 4 秒。

我能够解决此问题的唯一方法是设置cores=1cores.max=1,以便它只创建一个要执行的任务。

此任务的位置为NODE_LOCAL。因此,似乎当NODE_LOCAL 执行是即时的,但当Locality 为ANY 时,连接到kafka 需要4 秒。所有机器都在同一个 10Gb 网络中。知道为什么会这样吗?

【问题讨论】:

    标签: apache-spark apache-kafka spark-streaming mesos


    【解决方案1】:

    问题出在 spark.locality.wait 上,this link 给了我这个想法

    它的默认值是 3 秒,并且在 spark 流中处理的每个批次都花费了整个时间。

    使用 Mesos (--conf spark.locality.wait=0) 提交作业时,我已将其设置为 0 秒,现在一切都按预期运行。

    【讨论】:

      猜你喜欢
      • 2017-06-18
      • 1970-01-01
      • 1970-01-01
      • 2020-06-30
      • 2021-11-07
      • 2019-03-12
      • 1970-01-01
      • 2015-10-19
      • 1970-01-01
      相关资源
      最近更新 更多