【问题标题】:Spark SQL Java GenericRowWithSchema cannot be cast to java.lang.StringSpark SQL Java GenericRowWithSchema 无法转换为 java.lang.String
【发布时间】:2019-04-18 14:33:55
【问题描述】:

我有一个应用程序试图从集群目录中读取一组 csv,并使用 Spark 将它们写入 parquet 文件。

SparkSession sparkSession = createSession();
    JavaRDD<Row> entityRDD = sparkSession.read()
            .csv(dataCluster + "measures/measures-*.csv")
            .javaRDD()
            .mapPartitionsWithIndex(removeHeader, false)
            .map((Function<String, Measure>) s -> {
                String[] parts = s.split(COMMA);
                Measure measure = new Measure();
                measure.setCobDate(parts[0]);
                measure.setDatabaseId(Integer.valueOf(parts[1]));
                measure.setName(parts[2]);

                return measure;
            });

    Dataset<Row> entityDataFrame = sparkSession.createDataFrame(entityRDD, Measure.class);
    entityDataFrame.printSchema();

    //Create parquet file here
    String parquetDir = dataCluster + "measures/parquet/measures";
    entityDataFrame.write().mode(SaveMode.Overwrite).parquet(parquetDir);


    sparkSession.stop();

Measure 类是一个实现 Serializable 的简单 POJO。模式已打印,因此将 DataFrame 条目转换为 parquet 文件肯定有问题。 这是我得到的错误:

Lost task 2.0 in stage 1.0 (TID 3, redlxd00006.fakepath.com, executor 1): org.apache.spark.SparkException: Task failed while writing rows
        at org.apache.spark.sql.execution.datasources.FileFormatWriter$.org$apache$spark$sql$execution$datasources$FileFormatWriter$$executeTask(FileFormatWriter.scala:204)
        at org.apache.spark.sql.execution.datasources.FileFormatWriter$$anonfun$write$1$$anonfun$3.apply(FileFormatWriter.scala:129)
        at org.apache.spark.sql.execution.datasources.FileFormatWriter$$anonfun$write$1$$anonfun$3.apply(FileFormatWriter.scala:128)
        at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
        at org.apache.spark.scheduler.Task.run(Task.scala:99)
        at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:282)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
        at java.lang.Thread.run(Thread.java:745)
    Caused by: java.lang.ClassCastException: org.apache.spark.sql.catalyst.expressions.GenericRowWithSchema cannot be cast to java.lang.String
        at org.apache.spark.api.java.JavaPairRDD$$anonfun$toScalaFunction$1.apply(JavaPairRDD.scala:1040)
        at scala.collection.Iterator$$anon$11.next(Iterator.scala:409)
        at scala.collection.Iterator$$anon$11.next(Iterator.scala:409)
        at scala.collection.Iterator$$anon$11.next(Iterator.scala:409)
        at org.apache.spark.sql.execution.datasources.FileFormatWriter$SingleDirectoryWriteTask.execute(FileFormatWriter.scala:244)
        at org.apache.spark.sql.execution.datasources.FileFormatWriter$$anonfun$org$apache$spark$sql$execution$datasources$FileFormatWriter$$executeTask$3.apply(FileFormatWriter.scala:190)
        at org.apache.spark.sql.execution.datasources.FileFormatWriter$$anonfun$org$apache$spark$sql$execution$datasources$FileFormatWriter$$executeTask$3.apply(FileFormatWriter.scala:188)
        at org.apache.spark.util.Utils$.tryWithSafeFinallyAndFailureCallbacks(Utils.scala:1341)
        at org.apache.spark.sql.execution.datasources.FileFormatWriter$.org$apache$spark$sql$execution$datasources$FileFormatWriter$$executeTask(FileFormatWriter.scala:193)
        ... 8 more

最终我的意图是使用Spark SQL过滤和连接其他csv的数据,包含其他表数据,并将整个结果写入parquet。 我只发现了没有解决我的问题的 scala 相关问题。非常感谢任何帮助。

csv:

cob_date, database_id, name
20181115,56459865,name1
20181115,56652865,name6
20181115,56459845,name32
20181115,15645936,name3

【问题讨论】:

    标签: java apache-spark apache-spark-sql parquet mapr


    【解决方案1】:
    .map((Function<String, Measure>) s -> {
    

    看起来应该是这里

    .map((Function<Row, Measure>) s -> {
    

    【讨论】:

      【解决方案2】:

      按照 Serge 的建议添加 toDF() 并更新地图 lambda 解决了我的问题:

      SparkSession sparkSession = createSession();
      JavaRDD<Row> entityRDD = sparkSession.read()
           .csv(prismDataCluster + "measures/measures-*chop.csv")
           .toDF("cobDate","databaseId","name")
           .javaRDD()
           .mapPartitionsWithIndex(removeHeader, false)
           .map((Function<Row, Measure>) row -> {
                    Measure measure = new Measure();
                    measure.setCobDate(row.getString(row.fieldIndex("cobDate")));
                    measure.setDatabaseId(row.getString(row.fieldIndex("databaseId")));
                    measure.setName(row.getString(row.fieldIndex("name")));
      

      TVM。

      【讨论】:

        猜你喜欢
        • 2021-02-25
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2014-06-04
        • 1970-01-01
        • 1970-01-01
        • 2014-12-14
        • 1970-01-01
        相关资源
        最近更新 更多