【问题标题】:Apache Spark + Java: "java.lang.AssertionError: assertion failed" in ExpressionEncoderApache Spark + Java:ExpressionEncoder 中的“java.lang.AssertionError:断言失败”
【发布时间】: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


【解决方案1】:

我发现问题出在编码器的使用上。我应该使用 Encoders.INT() 而不是 Encoders.bean(Integer.class)

【讨论】:

    猜你喜欢
    • 2021-01-20
    • 2018-12-26
    • 1970-01-01
    • 2019-02-02
    • 2022-08-17
    • 2011-05-27
    • 1970-01-01
    • 2017-08-07
    相关资源
    最近更新 更多