【发布时间】:2015-08-09 16:06:23
【问题描述】:
我想知道与foreach 方法相比,考虑到我按顺序流经RDD 的情况,由于更高级别的并行性,foreachPartition 是否会产生更好的性能对累加器变量进行一些求和。
【问题讨论】:
标签: java scala foreach apache-spark
我想知道与foreach 方法相比,考虑到我按顺序流经RDD 的情况,由于更高级别的并行性,foreachPartition 是否会产生更好的性能对累加器变量进行一些求和。
【问题讨论】:
标签: java scala foreach apache-spark
foreach 和 foreachPartitions 是操作。
用于调用具有副作用的操作的通用函数。对于每个 RDD 中的元素,它调用传递的函数。 这是 通常用于操作累加器或写入外部 商店。
注意:在foreach() 之外修改除累加器以外的变量可能会导致未定义的行为。详情请见Understanding closures。
scala> val accum = sc.longAccumulator("My Accumulator")
accum: org.apache.spark.util.LongAccumulator = LongAccumulator(id: 0, name: Some(My Accumulator), value: 0)
scala> sc.parallelize(Array(1, 2, 3, 4)).foreach(x => accum.add(x))
...
10/09/29 18:41:08 INFO SparkContext: Tasks finished in 0.317106 s
scala> accum.value
res2: Long = 10
类似于
foreach(),但不是为每个调用函数 元素,它为每个分区调用它。功能应该可以 接受一个迭代器。这比foreach()更有效,因为 它减少了函数调用的次数(就像mapPartitions() 一样)。
foreachPartition 的用法示例:
在 sparkstreaming (dstreams) 和 kafka 生产者中使用 foreachPartition
dstream.foreachRDD { rdd =>
rdd.foreachPartition { partitionOfRecords =>
// only once per partition You can safely share a thread-safe Kafka //producer instance.
val producer = createKafkaProducer()
partitionOfRecords.foreach { message =>
producer.send(message)
}
producer.close()
}
}
注意:如果你想避免这种为每个分区创建一次生产者的方式,更好的方法是使用广播生产者
sparkContext.broadcast因为 Kafka 生产者是异步的,并且 在发送前大量缓冲数据。
测试(“Foreach - 火花”){ 导入 spark.implicits._ var accum = sc.longAccumulator sc.parallelize(Seq(1,2,3)).foreach(x => accum.add(x)) 断言(accum.value == 6L) } test("Foreach 分区 - Spark") { 导入 spark.implicits._ var accum = sc.longAccumulator sc.parallelize(Seq(1,2,3)).foreachPartition(x => x.foreach(accum.add(_))) 断言(accum.value == 6L) }累加器采样 sn-p 来玩弄它...通过它 你可以测试一下性能
foreachPartition分区上的操作很明显它会是 比foreach更好的优势
foreachPartition访问成本高时应使用 诸如数据库连接或 kafka 生产者等资源将初始化 每个分区一个,而不是每个元素一个(foreach)。当它 来到蓄电池,您可以通过上述测试来衡量性能 方法,在累加器的情况下也应该工作得更快..
另外...参见map vs mappartitions,它具有相似的概念,但它们是转换。
【讨论】:
foreach 在多个节点上自动运行循环。
但是,有时您想在每个节点上执行一些操作。例如,建立与数据库的连接。您不能只建立一个连接并将其传递给foreach 函数:连接只在一个节点上建立。
因此,使用foreachPartition,您可以在运行循环之前连接到每个节点上的数据库。
【讨论】:
foreach 和 foreachPartitions 之间确实没有太大区别。在幕后,foreach 所做的只是使用提供的函数调用迭代器的foreach。 foreachPartition 只是让您有机会在迭代器循环之外做一些事情,通常是一些昂贵的事情,比如启动数据库连接或类似的事情。因此,如果您没有任何事情可以为每个节点的迭代器执行一次并在整个过程中重复使用,那么我建议使用foreach 以提高清晰度并降低复杂性。
【讨论】:
foreachPartition 仅在您遍历按分区聚合的数据时才有用。
一个很好的例子是处理每个用户的点击流。每次完成用户的事件流时,您都希望清除计算缓存,但将其保留在同一用户的记录之间,以便计算一些用户行为洞察。
【讨论】:
foreachPartition 并不意味着它是每个节点的活动,而是针对每个分区执行的,与节点数量相比,您可能拥有大量分区,在这种情况下,您的性能可能会下降。如果您打算在节点级别进行活动,here 解释的解决方案可能很有用,尽管它未经我测试
【讨论】: