【问题标题】:How execution will happen in this spark program?这个 spark 程序将如何执行?
【发布时间】:2019-06-05 15:01:46
【问题描述】:

我遇到了下面的例子:

lines = sc.textFile("some_file.txt") //line_1

lineswithFriday = lines.filter(lambda line: line.startwith("Friday")) //line_2

lineswithFriday.first(); //line_3

它也说

spark 仅扫描文件,直到找到以开头的第一行 friday。它不需要遍历整个文件。

我的问题是:这是否意味着 spark 会在内存中一一加载每一行,看看它是否以 Friday 开头,如果是则停在那里?

假设line_1 基于核心和输入块创建了三个分区。 line_2 将通过每个内核上的单独工作线程进行计算。 在line_3 上,一旦任何工作人员找到以Friday 开头的行,它就会停在那里?

【问题讨论】:

  • >在 line_3 上,一旦任何工作找到从星期五开始的第一行,它就会停在那里?执行者将完成他们的任务,他们不会在任务执行过程中停下来。在返回的结果集中,第一次出现。更多细节请尝试:lineswithFriday.first().explain,详细执行计划请尝试:lineswithFriday.first().explain(True)

标签: apache-spark pyspark rdd


【解决方案1】:

first() 和 take(n) 如果单独使用可以优化。

Spark 中没有关于“进程间通信”的机制允许执行程序提前终止,因为继续处理结果将被认为是多余的。从架构上讲,这会导致各种问题。

【讨论】:

  • 你的意思是即使一个线程找到以Friday开头的行的实例,它会继续扫描整个文件?如果这是我从谷歌spark scans the file only until it finds the first line starting with friday. It does not need to go through entire file. 引用的正确陈述不正确?
  • 第二点:如果每个内核上运行三个线程/进程,如何从三个线程返回的结果中确定第一个结果?
  • 您不会扫描文件,而是扫描文件的分区。一旦进行中的任务必须完成。但是对于没有过滤器的情况是正确的。是优化。
  • 结果将与分区ID一起返回给驱动程序。第一个将返回具有最低分区 id 的列表中的第一个值。它还能如何工作?
  • 这都是关于并行性的。您的优化想法将意味着某种顺序处理。这不是 Spark。
猜你喜欢
  • 2017-09-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-09-17
  • 1970-01-01
  • 2017-02-07
  • 2021-11-17
相关资源
最近更新 更多