【发布时间】: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