我认为您对 BoxedUnit 的含义有错误的印象,因此坚持在 Java 中使用 Scala 接口,由于暴露给 Java 的 Scala 中隐藏的复杂性的数量过于复杂。 scala.Function1<scala.collection.Iterator<T>, scala.runtime.BoxedUnit> 是 (Iterator[T]) => Unit 的实现——一个接受 Iterator[T] 并返回 Unit 类型的 Scala 函数。 Scala 中的Unit 相当于Java 的void。 BoxedUnit 是 Unit 的盒装版本 - 它是一个堆对象,在其 UNIT 成员中保存单例单元值,并且是几乎从未在 Scala 程序中出现的实现细节。如果数据集是 DataFrame,那么 T 将是 org.apache.spark.sql.Row,并且您需要在 Row 对象的集合上处理 Scala 迭代器。
要在 Java 中定义 scala.Function1<scala.collection.Iterator<Row>, scala.runtime.BoxedUnit>,您需要创建 AbstractFunction1<scala.collection.Iterator<Row>, scala.runtime.BoxedUnit> 的实例并覆盖其 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<T> 而上面的代码变为:
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 的流引擎自动为您完成的。您编写指定 transformations 和 aggregations 的流式查询,然后随着新数据从流中到达而逐步应用。
foreachPartition 是foreach 的一种形式,保留用于一些特殊的批处理情况,例如,当您需要在处理函数中进行一些昂贵的对象实例化并且对每一行都进行时会产生巨大的开销。使用foreachPartition,您的处理函数每个分区只调用一次,因此您可以实例化一次昂贵的对象,然后遍历分区的数据。这会减少处理时间,因为您只需执行一次昂贵的操作。
但是,您甚至不能在流媒体源上调用foreach() 或foreachPartition(),因为这会导致AnalysisException。相反,您必须使用DataStreamWriter 的foreach() 或foreachBatch() 方法。 DataStreamWriter.foreach() 采用 ForeachWriter 的实例,而 DataStreamWriter.foreachBatch() 采用接收数据集和批次 ID 的 void 函数。 ForeachWriter 在其 open() 方法中接收一个纪元 ID。同样,foreachBatch() 具有功能相同的 Scala 和 Java 两种风格,因此如果您要使用 Java 编写,请使用 Java 特定的风格。