【问题标题】:Spark Java DataFrame - ClassCastException when reading multiple files into multiple datasetsSpark Java DataFrame - 将多个文件读入多个数据集时出现 ClassCastException
【发布时间】:2018-03-17 18:08:22
【问题描述】:

我正在尝试将数据从单独的文件读取到单独的 RDD 中,然后我将其转换为 DataFrame(使用 Java api)。

在使用单个 POJO 仅处理一个数据集时,我没有遇到任何问题,但是当我尝试读取映射到不同 POJO 的附加数据集时,我就开始遇到此问题:

18/03/17 00:58:28 ERROR Executor: Exception in task 0.0 in stage 1.0 (TID 1)
java.lang.ClassCastException: TestMain$PageView cannot be cast to TestMain$BlacklistedPage
    at org.apache.spark.api.java.JavaRDD$$anonfun$filter$1.apply(JavaRDD.scala:78)
    at org.apache.spark.api.java.JavaRDD$$anonfun$filter$1.apply(JavaRDD.scala:78)
    at scala.collection.Iterator$$anon$13.hasNext(Iterator.scala:463)
    at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408)
    at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:234)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:228)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$25.apply(RDD.scala:827)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$25.apply(RDD.scala:827)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:323)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:287)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
    at org.apache.spark.scheduler.Task.run(Task.scala:108)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:338)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

以下是一些测试代码,似乎可以复制我遇到的问题(使用并行数据而不是 textFile 输入)。我正在使用 Spark 2.2.1。我滥用的 SparkSession 有什么问题吗?

import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;

import java.io.Serializable;
import java.util.Arrays;
import java.util.Objects;

public class TestMain {
    public static void main(String[] args) {
        SparkSession spark = SparkSession.builder().master("local").appName("demo").getOrCreate();

        JavaSparkContext sparkContext = JavaSparkContext.fromSparkContext(spark.sparkContext());

        JavaRDD<String> rawBlacklisted = sparkContext.parallelize(Arrays.asList("af .sy", "af 2009"));
        JavaRDD<String> raw = sparkContext.parallelize(Arrays.asList("ab .sy 100", "af 2009 10", "aa title 5"));

        JavaRDD<BlacklistedPage> blackListedPages = rawBlacklisted.map(BlacklistedPage::parse).filter(Objects::nonNull);
        JavaRDD<PageView> rawPageViews = raw.map(PageView::parse).filter(Objects::nonNull);

        Dataset<Row> first = spark.createDataFrame(blackListedPages, BlacklistedPage.class);
        Dataset<Row> second = spark.createDataFrame(rawPageViews, PageView.class);

        first.show(10);
        second.show(10);
    }

    public static class BlacklistedPage implements Serializable {
        private String domainCode;
        private String pageTitle;

        static  BlacklistedPage parse(String line) {
            String[] data = line.split(" ");
            if (data.length < 2) {
                return null;
            }
            return new BlacklistedPage(data[0], data[1]);
        }

        BlacklistedPage(String domainCode, String pageTitle) {
            this.domainCode = domainCode;
            this.pageTitle = pageTitle;
        }

        // getters and setters omitted for clarity
    }

    public static class PageView implements Serializable {
        private String domainCode;
        private String pageTitle;
        private Integer viewCount;

        static PageView parse(String line) {
            String[] data = line.split(" ");

            if (data.length < 3) {
                return null;
            }

            return new PageView(data[0], data[1], Integer.parseInt(data[2]));
        }

        PageView(String domainCode, String pageTitle, Integer viewCount) {
            this.domainCode = domainCode;
            this.pageTitle = pageTitle;
            this.viewCount = viewCount;
        }

        // getters and setters omitted for clarity
    }
}

【问题讨论】:

    标签: java apache-spark spark-dataframe apache-spark-dataset


    【解决方案1】:

    您好,我刚刚修改了您的部分代码,现在它似乎工作正常:

        JavaRDD<PageView> rawPageViews = raw.map(PageView::parse).filter(new Function<TestMain.PageView, Boolean>() {
    
            /**
             * 
             */
            private static final long serialVersionUID = 1L;
    
            @Override
            public Boolean call(PageView arg0) throws Exception {
                // TODO Auto-generated method stub
                boolean nullcatcher = true;
                if (arg0==null||arg0.equals(null)){
                    nullcatcher = false;
                }
                return nullcatcher;
            }
        });
    

    【讨论】:

    • 我想知道在过滤器调用上使用Objects::nonNull 是否有什么问题。我还发现用以下内容替换我的过滤器似乎有效:JavaRDD&lt;PageView&gt; rawPageViews = raw.map(PageView::parse).filter(p -&gt; p != null);
    猜你喜欢
    • 2020-02-02
    • 1970-01-01
    • 2020-01-29
    • 2017-06-23
    • 1970-01-01
    • 2012-06-28
    • 1970-01-01
    • 2021-01-05
    • 2023-01-20
    相关资源
    最近更新 更多