【问题标题】:How to use foreachPartition in Spark?如何在 Spark 中使用 foreachPartition?
【发布时间】:2020-07-25 05:17:36
【问题描述】:

如何在 Spark Java 中使用以下函数?找遍了整个互联网,但找不到合适的例子。

public void foreachPartition(scala.Function1<scala.collection.Iterator<T>,scala.runtime.BoxedUnit> f)

我唯一知道的是它对进程batch of data有好处,也就是所谓的BoxedUnit

如何获得datasetbatch IDBoxedUnit 来批量处理数据?

谁能告诉我如何实现这个方法?

【问题讨论】:

  • 您是否尝试将您的 RDD 转换为 JavaRDD,然后应用 foreachpartition?以 VoidFunction 作为参数。
  • @ShemTov 它与 voidfunction 无关,我想按原样使用它。很想知道如何获取 batchId 来处理来自 Dataset 的一批数据。
  • 我猜你可以自己实现 batchId,使用累加器 - 从零开始,foreachPartition 的每个批次都将它增加一。它会在所有执行程序之间保持同步并为您提供有效的 batchId。

标签: java apache-spark spark-structured-streaming


【解决方案1】:

我认为您对 BoxedUnit 的含义有错误的印象,因此坚持在 Java 中使用 Scala 接口,由于暴露给 Java 的 Scala 中隐藏的复杂性的数量过于复杂。 scala.Function1&lt;scala.collection.Iterator&lt;T&gt;, scala.runtime.BoxedUnit&gt;(Iterator[T]) =&gt; Unit 的实现——一个接受 Iterator[T] 并返回 Unit 类型的 Scala 函数。 Scala 中的Unit 相当于Java 的voidBoxedUnitUnit 的盒装版本 - 它是一个堆对象,在其 UNIT 成员中保存单例单元值,并且是几乎从未在 Scala 程序中出现的实现细节。如果数据集是 DataFrame,那么 T 将是 org.apache.spark.sql.Row,并且您需要在 Row 对象的集合上处理 Scala 迭代器。

要在 Java 中定义 scala.Function1&lt;scala.collection.Iterator&lt;Row&gt;, scala.runtime.BoxedUnit&gt;,您需要创建 AbstractFunction1&lt;scala.collection.Iterator&lt;Row&gt;, scala.runtime.BoxedUnit&gt; 的实例并覆盖其 apply() 方法,您必须返回 BoxedUnit.UNIT。您还需要使其可序列化,因此您通常声明自己的类,该类继承自AbstractFunction1 并实现Serializable。你也可以通过暴露一个不同的、对 Java 更友好的抽象方法来对其进行 Java 化:

import org.apache.spark.sql.Row;
import scala.runtime.AbstractFunction1;
import scala.runtime.BoxedUnit;
import scala.collection.JavaConverters;
import java.util.Iterator;

class MyPartitionFunction<T> extends AbstractFunction1<scala.collection.Iterator<T>, BoxedUnit>
   implements Serializable {
   @Override
   public BoxedUnit apply(scala.collection.Iterator<T> iterator) {
      call(JavaConverters.asJavaIteratorConverter(iterator).asJava());
      return BoxedUnit.UNIT;
   }

   public abstract void call(Iterator<T> iterator);
}

df.foreachPartition(new MyPartitionFunction<Row>() {
   @Override
   public void call(Iterator<Row> iterator) {
      for (Row row : iterator) {
         // do something with the row
      }
   }
});

这是相当数量的实现复杂性,这就是为什么有 Java 特定版本采用 ForeachPartitionFunction&lt;T&gt; 而上面的代码变为:

import org.apache.spark.sql.Row;
import org.apache.spark.api.java.function.ForeachPartitionFunction;
import java.util.Iterator;

df.foreachPartition(new ForeachPartitionFunction<Row>() {
   public void call(Iterator<Row> iterator) throws Exception {
      for (Row row : iterator) {
         // do something with the row
      }
   }
}

功能与Scala接口提供的完全相同,只是Apache Spark为您进行迭代器转换,它还为您提供了一个不需要您的友好Java类导入和实现 Scala 类型。

也就是说,我认为您对 Spark 的工作原理有些误解。您不需要使用foreachPartition 批量处理流数据。这是由 Spark 的流引擎自动为您完成的。您编写指定 transformationsaggregations 的流式查询,然后随着新数据从流中到达而逐步应用。

foreachPartitionforeach 的一种形式,保留用于一些特殊的批处理情况,例如,当您需要在处理函数中进行一些昂贵的对象实例化并且对每一行都进行时会产生巨大的开销。使用foreachPartition,您的处理函数每个分区只调用一次,因此您可以实例化一次昂贵的对象,然后遍历分区的数据。这会减少处理时间,因为您只需执行一次昂贵的操作。

但是,您甚至不能在流媒体源上调用foreach()foreachPartition(),因为这会导致AnalysisException。相反,您必须使用DataStreamWriterforeach()foreachBatch() 方法。 DataStreamWriter.foreach() 采用 ForeachWriter 的实例,而 DataStreamWriter.foreachBatch() 采用接收数据集和批次 ID 的 void 函数。 ForeachWriter 在其 open() 方法中接收一个纪元 ID。同样,foreachBatch() 具有功能相同的 Scala 和 Java 两种风格,因此如果您要使用 Java 编写,请使用 Java 特定的风格。

【讨论】:

  • lliev,scala.Function1 类型是否等同于 MyPartitionFunction 类型?
  • 我的理解是否正确,Dataset 分区的 BoxedUnit 值是在调用 foreachPartition 然后调用以下 apply 方法时在内部生成的? UNIT 不是我手上可以做的事情,而是供内部参考的。正确吗?
  • foreachPartition 什么都不返回,它会作为 Action 工作吗?
  • @Maria, BoxedUnit不是数据集的属性。它是单例模式的一种实现——只有一个BoxedUnit 实例,它可以通过静态UNIT 成员获得。该实例由 Scala 运行时创建。 Scala是一种函数式编程语言,每个函数都必须返回一个值,所以Unit是一个特殊值,当函数不返回时隐式返回,即相当于Java中的void。这与 Spark 无关——它只是 Scala 中的一个实现细节,暴露给 Java 程序以实现互操作性。
  • apply() 方法在函数对象被调用时被调用。因此,每当您调用foreachParitition 方法时,驱动程序都会序列化一堆MyPartitionFunction 实例并将它们发送到执行程序,然后执行程序调用apply() 方法,将一个迭代器传递给相应分区中的所有数据。同样,apply() 方法适用于 Scala 的工作方式。 foreachPartition 的等效 Java 特定版本通过 ForeachPartitionFunction 实例的 call() 方法接收迭代器。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-06-20
  • 2018-08-12
  • 2018-05-29
  • 2015-08-09
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多