【问题标题】:why there are so many tasks in my spark streaming job为什么我的 Spark Streaming 工作中有这么多任务
【发布时间】:2016-11-21 20:37:33
【问题描述】:

我想知道为什么我的 spark 流作业中有这么多任务号?它变得越来越大......

运行3.2小时后,增长到120020……运行一天后,增长到100万……为什么?

【问题讨论】:

  • 你的工作是做什么的?你能添加“流媒体”标签吗?

标签: apache-spark spark-streaming


【解决方案1】:

我强烈建议您检查参数spark.streaming.blockInterval,这是一个非常重要的参数。默认为 0.5 秒,即每 0.5 秒创建一个任务。

所以也许您可以尝试将 spark.streaming.blockInterval 增加到 1 分钟或 10 分钟,然后任务数应该会减少。

我的直觉只是因为你的消费者和生产者一样快,所以随着时间的推移,越来越多的任务被积累起来供进一步消费。

这可能是由于您的 Spark 集群无法处理如此大的批次。也可能与检查点间隔时间有关,可能你设置的太大或太小。也可能与您的ParallelismPartitionsData Locality等设置有关。

祝你好运

阅读本文

Tuning Spark Streaming for Throughput

.

How-to: Tune Your Apache Spark Jobs (Part 1)

.

How-to: Tune Your Apache Spark Jobs (Part 2)

【讨论】:

  • 你好,我觉得不正常的是“跳过的任务数”越来越多。实际处理任务数是常数。但是“跳过任务”的数量越来越多。我不明白哪些任务被跳过了……
  • 我遇到了同样的错误,但不太确定其原因,可能是内存不足或 RDD 拆分错误。我不确定……
  • 我猜这可能取决于rdd的血统。新的流式任务的血统没有被切断,“跳过的任务”是之前的血统。所以新任务记住的任务会越来越多。一旦一个新任务失败,所有原始“跳过的任务”都将再次执行。
  • 所以我想也许 checkpoint() 会解决这个问题。你知道我应该在哪里添加检查点
  • 对不起,我不知道,我现在也在检查点。我当前的检查点设置是 1 秒,并试图优化它:P
【解决方案2】:

SparkUI 功能意味着某些阶段依赖项可能已被计算或未计算,但由于它们的输出已经可用而被跳过。因此它们显示为skipped

请不要使用might,这意味着在工作完成之前Spark 不确定是否需要返回并重新计算最初跳过的一些阶段。

[1]https://github.com/apache/spark/blob/master/core/src/main/scala/org/apache/spark/ui/jobs/JobProgressListener.scala#L189

【讨论】:

    【解决方案3】:

    流式应用程序的本质是随着时间的推移为每批数据运行相同的进程。看起来您正在尝试以 1 秒的批处理间隔运行,并且每个间隔可能会产生多个作业。您在 3.2 小时内显示了 585 个工作,而不是 120020。但是,您的处理看起来也像是在 1 秒内完成。我想您的日程安排延迟非常高。我猜这是批处理间隔太小的症状。

    【讨论】:

    • 时间间隔为2分钟,我认为我的工作运行正常。
    • 而工作开始时,“所有任务号”并没有那么多,只有几百个。 3.2小时后增长到120020,半天后增长到1000000+。
    • 我无法理解 spark 流任务中的“跳过的任务”是什么,为什么它们被跳过了?以及为什么跳过的任务数越来越多
    猜你喜欢
    • 2016-10-12
    • 2015-12-20
    • 1970-01-01
    • 2013-10-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-25
    相关资源
    最近更新 更多