【发布时间】:2020-04-07 23:00:10
【问题描述】:
我尝试使用以下代码在数据集上调用 groupByKey:
SparkSession SPARK_SESSION = new SparkSession(new SparkContext("local", "app"));
JavaSparkContext JAVA_SPARK_CONTEXT = new JavaSparkContext(SPARK_SESSION.sparkContext());
@Data
@NoArgsConstructor
@AllArgsConstructor
class Chunk implements Serializable {
private Integer id;
private String letters;
}
class JavaAggregator extends Aggregator<Chunk, String, String> {
@Override
public String zero() {
return "";
}
@Override
public String reduce(String b, Chunk a) {
return b + a.getLetters();
}
@Override
public String merge(String b1, String b2) {
return b1 + b2;
}
@Override
public String finish(String reduction) {
return reduction;
}
@Override
public Encoder<String> bufferEncoder() {
return Encoders.bean(String.class);
}
@Override
public Encoder<String> outputEncoder() {
return Encoders.bean(String.class);
}
}
List<Chunk> chunkList = List.of(
new Chunk(1, "a"), new Chunk(2, "1"), new Chunk(3, "-*-"),
new Chunk(1, "b"), new Chunk(2, "2"), new Chunk(3, "-**-"),
new Chunk(1, "c"), new Chunk(2, "3"), new Chunk(3, "-***-"));
Dataset<Row> df = SPARK_SESSION.createDataFrame(JAVA_SPARK_CONTEXT.parallelize(chunkList), Chunk.class);
Dataset<Chunk> ds = df.as(Encoders.bean(Chunk.class));
KeyValueGroupedDataset<Integer, Chunk> grouped = ds.groupByKey((Function1<Chunk, Integer>) v -> v.getId(), Encoders.bean(Integer.class));
但我得到 Exception 说 java.lang.AssertionError: 断言失败
at scala.Predef$.assert(Predef.scala:208)
at org.apache.spark.sql.catalyst.encoders.ExpressionEncoder$.javaBean(ExpressionEncoder.scala:87)
at org.apache.spark.sql.Encoders$.bean(Encoders.scala:142)
at org.apache.spark.sql.Encoders.bean(Encoders.scala)
我不是 scala 内部的专家,我很难说出代码有什么问题,因为一些编译代码会抛出异常,并且断言消息“断言失败”不是很有帮助。但也许我在代码中所做的一些根本性错误导致了这个异常?
【问题讨论】:
-
Spark 的版本是多少?我的猜测是断言是由于缺少 getter 和 setter。
-
这应该不是问题,因为我使用 lombok 的 @Data 注释来创建 getter/setter。但我发现问题出在编码器的使用上。我应该使用 Encoders.INT() 而不是 Encoders.bean(Integer.class)。我知道如果 Encoders.INT() 存在,那么我应该使用它,但还不明白为什么 Encoders.bean(Integer.class) 不起作用。无论如何,代码正在运行。感谢 Jacek 的建议 :)
-
您能否回答您的问题(甚至通过复制和粘贴您的评论)并接受以使 Spark 用户的生活更轻松? :)
标签: java apache-spark apache-spark-sql