【问题标题】:Concurrency in spark aggregate methodspark聚合方法中的并发
【发布时间】:2016-04-16 22:49:46
【问题描述】:

我是 Spark 和 Map Reduce 的初学者,据我了解 spark aggregate(ByKey) 方法遵循 map reduce 模式,我希望有人帮我确认它是否正确。

  1. 第一个函数参数“segFunc”获取每个键的数据 并为每个键并行运行。就像地图中的 map() 减少。
  2. 第二个函数参数“combFun”收集数据 对于每个键,即使跨分区,它也不是并行运行的 并且系统保证了这个联合收割机的同步 所有键之间的功能。就像 map reduce 中的 combiner()。

请指正,非常感谢。

【问题讨论】:

    标签: java scala concurrency apache-spark mapreduce


    【解决方案1】:

    它遵循 map/reduce 模式,但你的 map/reduce 模式错误。

    第一阶段将并行运行,并为每个键创建一条记录(这些记录将保存在内存中或溢出到磁盘,具体取决于 Spark 中的可用资源还是保存到 Hadoop 中的磁盘)

    然后下一阶段将(或至少可以)也并行运行 - 每个键。之前创建的数据将被获取并合并,因此每个键的数据将到达单个目的地(reducer)

    获取阶段称为洗牌

    Hadoop 中的组合器正在执行类似 reduce 的行为并在映射阶段发出部分结果(朝向 reducer)

    【讨论】:

    • 非常感谢,@arnon,所以在 map-reduce 中合并每个键的数据的第二阶段可以并行运行,然后在 spark 中的 aggregateByKey 方法的“combFun”函数参数的情况下,如果我想得到所有数据的集合,那么我必须使用synchronizedList而不是普通列表,对吗?
    • 再次感谢 Arnon,在阅读了一些内容后,我意识到我混淆了“组合器”和“组合()函数”的概念,每个桶(分区)的组合器可以并行运行,但是组合器运行 combine() 方法一次,与 mapper 和 reducer 相同。
    • 回到我的问题,因为函数参数只在其单独的实例中运行,除了结果之外什么都不会在它们之间共享,函数是线程安全的,它的同步是有保证的。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-01-11
    • 2013-08-24
    • 2016-11-22
    • 1970-01-01
    • 1970-01-01
    • 2019-04-07
    • 2018-01-06
    相关资源
    最近更新 更多