【问题标题】:Save as Parquet file in spark java在 spark java 中另存为 Parquet 文件
【发布时间】:2018-02-05 13:47:46
【问题描述】:

我是 Spark 的新手。我正在尝试在本地模式(Windows)下使用 spark java 将 csv 文件保存为镶木地板。我收到了这个错误。

原因:org.apache.spark.SparkException:写入行时任务失败

我参考了其他线程并禁用了火花推测

set("spark.speculation", "false")

我仍然得到错误。我在 csv 中仅使用两列进行测试,但仍然出现在这个问题中。

输入:

portfolio_id;portfolio_code
1000042;CHNTESTPF04
1000042;CHNTESTPF04
1000042;CHNTESTPF04
1000042;CHNTESTPF04
1000042;CHNTESTPF04
1000042;CHNTESTPF04
1000042;CHNTESTPF04

我的代码:

JavaPairRDD<Integer, String> rowJavaRDD = pairRDD.mapToPair(new PairFunction<String, Integer, String>() {
    private Double[] splitStringtoDoubles(String s){
       String[] splitVals = s.split(";");
       Double[] vals = new Double[splitVals.length];
       for(int i= 0; i < splitVals.length; i++){
           vals[i] = Double.parseDouble(splitVals[i]);
       }
       return vals;
    }

    @Override
    public Tuple2<Integer, String> call(String arg0) throws Exception {
        // TODO Auto-generated method stub
        return null;
    }
});


SQLContext SQLContext;
SQLContext = new org.apache.spark.sql.SQLContext(sc);

Dataset<Row> fundDF = SQLContext.createDataFrame(rowJavaRDD.values(), funds.class);
fundDF.printSchema();

fundDF.write().parquet("C:/test");

请帮助我在这里缺少的东西。

【问题讨论】:

  • 请把完整的错误和堆栈跟踪放在你的问题中。
  • 我通过在 Tuple2() 中添加函数拆分解决了错误,如下所示:public void run(String t, String u){ public Tuple2 call(String rec) { String[] 标记 = rec.split(";"); String[] vals = new String[tokens.length]; for(int i= 0; i (tokens[0], tokens[1]); } });
  • @Ans8 请将您的解决方案放入答案中并接受它,这样它就会从“未回答”部分消失。

标签: java apache-spark parquet


【解决方案1】:

请找到我的 Java/Spark 代码 1) 加载 CSV indo Spark 数据集 2) 将数据集保存到镶木地板

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.SparkSession;

SparkSession spark = SparkSession
                    .builder()
                    .appName("csv2parquet")
                    .config("spark.sql.warehouse.dir", "/file:/tmp")
                    .master("local")
                    .getOrCreate();

final String dir = "test.csv";


Dataset<Row> ds = spark.read().option("header", true).option("inferSchema", true).csv(dir);

final String parquetFile = "test.parquet";
final String codec = "parquet";

ds.write().format(codec).save(parquetFile);

spark.stop();

将此添加到您的 pom 中

<dependency>
                <groupId>org.apache.hadoop</groupId>
                <artifactId>hadoop-mapreduce-client-core</artifactId>
                <version>2.8.1</version>
</dependency>

【讨论】:

    【解决方案2】:

    在这里回答,所以它从@Glennie Helles Sindholt 所说的未回答部分开始。很抱歉延迟发布此内容

    我通过在 Tuple2() 中添加函数拆分解决了这个错误:

        public void run(String t, String u)
    
        { 
    
         public Tuple2<String,String> call(String rec){ 
            String[] tokens = rec.split(";"); 
            String[] vals = new String[tokens.length]; 
            for(int i= 0; i < tokens.length; i++)
           { 
                vals[i] =tokens[i]; 
            } 
    
          return new Tuple2<String, String>(tokens[0], tokens[1]); 
    
         } 
       });
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-09-06
      • 2016-07-08
      • 2018-09-24
      • 2017-08-19
      • 2018-10-29
      • 2016-12-15
      • 2017-08-07
      相关资源
      最近更新 更多