【问题标题】:SPARK 2.1.1: Active batches are queued up when used sliding windowSPARK 2.1.1:活动批次在使用滑动窗口时排队
【发布时间】:2017-06-30 10:44:18
【问题描述】:

我正在 Apache Spark 2.1.1 中开发一个流式应用程序,滑动窗口的持续时间为 X 秒。并滑动 Y 秒。(我为 X 和 Y 尝试了不同的值)。我正在使用.createDirectStream 读取来自 kafka(27 个分区)的消息。我的集群配置是,

  • 执行器内核总数=135

  • 节点=3

  • 每个执行者的核心 = 5

  • 执行器内存 = 10G

当我尝试以 X=5 秒的窗口持续时间执行作业时。并滑动 Y = 1 秒。没有活动批次排队。根据我的逻辑,这是因为在这种情况下,(# of total cores =135) 和 (# of task created= # kafka partitions * windowed batches= 27*5=135)。

现在在窗口持续时间 X=120 秒的其他情况下。并滑动 y=1 秒,(创建的任务数 = 27*120= 3240),这远远超过可用内核(135)。因此 观察到巨大的队列,并且在调度中需要大量时间。

我的解释正确吗?如果是,那有什么补救措施?

【问题讨论】:

    标签: apache-spark cluster-computing spark-streaming


    【解决方案1】:

    尽管不查看集群很难 100% 确定,但您的解释似乎是正确的。任务调度开销可能会对低延迟应用程序造成影响。

    解决此问题的一种可能方法是改用基于窗口的缩减,例如 reduceByKeyAndWindow,尤其是在您的用例允许的情况下采用可逆函数的重载:reduceByKeyAndWindow(func, invFunc, windowLength, slideInterval, [numTasks])

    另请注意,可以使用此调用强制减少分区数量:https://github.com/apache/spark/blob/master/streaming/src/main/scala/org/apache/spark/streaming/dstream/PairDStreamFunctions.scala#L303

    【讨论】:

    • 感谢马斯格。对于我的用例,我使用 120 秒持续时间和 1 秒滑动的滑动窗口,然后是一个 SQL 查询来查找最大值和最小值以及包含的 groupBy 子句。您能否建议,我如何使用reduceByKeyAndWindow 函数实现 SQL 功能?
    猜你喜欢
    • 1970-01-01
    • 2018-10-04
    • 1970-01-01
    • 1970-01-01
    • 2018-11-09
    • 2016-10-07
    • 2019-10-21
    • 1970-01-01
    • 2020-08-25
    相关资源
    最近更新 更多