【问题标题】:Apache Spark - foreach Vs foreachPartition When to use What?Apache Spark - foreach Vs foreachPartition 什么时候使用?
【发布时间】:2015-08-09 16:06:23
【问题描述】:

我想知道与foreach 方法相比,考虑到我按顺序流经RDD 的情况,由于更高级别的并行性,foreachPartition 是否会产生更好的性能对累加器变量进行一些求和。

【问题讨论】:

    标签: java scala foreach apache-spark


    【解决方案1】:

    foreachforeachPartitions 是操作。

    foreach(function): 单位

    用于调用具有副作用的操作的通用函数。对于每个 RDD 中的元素,它调用传递的函数。 这是 通常用于操作累加器或写入外部 商店。

    注意:在foreach() 之外修改除累加器以外的变量可能会导致未定义的行为。详情请见Understanding closures

    example

    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
    

    foreachPartition(函数): 单位

    类似于 foreach() ,但不是为每个调用函数 元素,它为每个分区调用它。功能应该可以 接受一个迭代器。这比foreach() 更有效,因为 它减少了函数调用的次数(就像 mapPartitions() 一样)。

    foreachPartition 的用法示例:


    • 示例 1:对于每个分区,您要使用一个数据库连接(每个分区块在内部),然后这是一个使用 scala 的示例用法。
    /** * 使用 foreach 分区插入数据库。 * * @param sqlDatabaseConnectionString * @param sqlTableName */ def insertToTable(sqlDatabaseConnectionString: String, sqlTableName: String): Unit = { //numPartitions = 您可以计划提供的同时数据库连接数 datframe.repartition(numofpartitionsyouwant) val tableHeader: String = dataFrame.columns.mkString(",") dataFrame.foreachPartition { 分区 => // 注意:每个分区一个连接(更好的方法是使用连接池) val sqlExecutorConnection: Connection = DriverManager.getConnection(sqlDatabaseConnectionString) //使用 1000 的批处理大小,因为某些数据库不能使用超过 1000 的批处理大小,例如:Azure sql partition.grouped(1000).foreach { 组=> val insertString: scala.collection.mutable.StringBuilder = new scala.collection.mutable.StringBuilder() group.foreach { 记录 => insertString.append("('" + record.mkString(",") + "'),") } sqlExecutorConnection.createStatement() .executeUpdate(f"插入 [$sqlTableName] ($tableHeader) 值" + insertString.stripSuffix(",")) } sqlExecutorConnection.close() // 关闭连接,这样连接就不会耗尽。 } }
    • 示例 2:

    在 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 生产者是异步的,并且 在发送前大量缓冲数据。


    累加器采样 sn-p 来玩弄它...通过它 你可以测试一下性能

    测试(“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) }

    结论:

    foreachPartition 分区上的操作很明显它会是 比foreach更好的优势

    经验法则:

    foreachPartition 访问成本高时应使用 诸如数据库连接或 kafka 生产者等资源将初始化 每个分区一个,而不是每个元素一个(foreach)。当它 来到蓄电池,您可以通过上述测试来衡量性能 方法,在累加器的情况下也应该工作得更快..

    另外...参见map vs mappartitions,它具有相似的概念,但它们是转换。

    【讨论】:

    • 一个很棒的解释能否请您添加 foreach 分区将比 foreach 慢的场景(假设是累加器),因为在这种情况下 foreachpartition 将在内部调用 foreach。
    • @RamGhadiyaram 我们是否在 JAVA 中提供了类似的功能。当我尝试在每个分区上使用 grouped() 时,它没有显示任何可用的此类方法。我正在使用 Spark 2.1.0
    • AFAIK scala 可用。所以它在java中不可用l你可以做正常的批处理操作。我的意思是你可以做类似的of
    • @RamGhadiyaram ,有 30 个分区和 30 个内核,需要将 15GB 数据复制到 cassandra ,在运行 SparkJob 时,我只有一个处理器承担所有负载,其他执行器无法参与处理。顺便说一句,我在 hdfs 中使用 parquet 文件格式保存,你能帮我吗
    • 打印分区长度。如果它是 1(因为一个分区正在加载)然后尝试重新分区,然后执行 foeach 分区。
    【解决方案2】:

    foreach 在多个节点上自动运行循环。

    但是,有时您想在每个节点上执行一些操作。例如,建立与数据库的连接。您不能只建立一个连接并将其传递给foreach 函数:连接只在一个节点上建立。

    因此,使用foreachPartition,您可以在运行循环之前连接到每个节点上的数据库。

    【讨论】:

    • 这仍然不是每个节点,而是每个分区。可以有比节点更多的分区。如果您需要每个节点的连接(更可能是每个 JVM 或 YARN 容器),您需要一些其他解决方案。
    • @user2456600 。您知道如何为每个执行器的每个 jvm 设置一个类吗?
    • 如果使用 Scala,一个选项是在对象或类中使用惰性 val,它会在第一次被引用时在 JVM 中初始化。但这也有缺点,如果每个执行程序使用多个线程,则必须小心它指向的对象是线程安全的。此外,很难将运行时初始化参数(如配置)传递给初始化。
    【解决方案3】:

    foreachforeachPartitions 之间确实没有太大区别。在幕后,foreach 所做的只是使用提供的函数调用迭代器的foreachforeachPartition 只是让您有机会在迭代器循环之外做一些事情,通常是一些昂贵的事情,比如启动数据库连接或类似的事情。因此,如果您没有任何事情可以为每个节点的迭代器执行一次并在整个过程中重复使用,那么我建议使用foreach 以提高清晰度并降低复杂性。

    【讨论】:

      【解决方案4】:

      foreachPartition 仅在您遍历按分区聚合的数据时才有用。

      一个很好的例子是处理每个用户的点击流。每次完成用户的事件流时,您都希望清除计算缓存,但将其保留在同一用户的记录之间,以便计算一些用户行为洞察。

      【讨论】:

        【解决方案5】:

        foreachPartition 并不意味着它是每个节点的活动,而是针对每个分区执行的,与节点数量相比,您可能拥有大量分区,在这种情况下,您的性​​能可能会下降。如果您打算在节点级别进行活动,here 解释的解决方案可能很有用,尽管它未经我测试

        【讨论】:

        • 我使用了类似的代码,使用 foreachPartition 将数据插入 Oracle。性能极其缓慢。
        猜你喜欢
        • 1970-01-01
        • 2018-08-14
        • 1970-01-01
        • 1970-01-01
        • 2022-06-25
        • 2015-12-12
        • 2016-09-09
        • 2015-05-14
        • 2018-10-28
        相关资源
        最近更新 更多