【问题标题】:Convert JavaPairRDD to Dataframe in Spark Java API在 Spark Java API 中将 JavaPairRDD 转换为 Dataframe
【发布时间】:2017-05-24 22:59:59
【问题描述】:

我正在使用带有 Java 7 的 Spark 1.6

我有一对 RDD:

JavaPairRDD<String, String> filesRDD = sc.wholeTextFiles(args[0]);

我想用架构将它转换成DataFrame

看来我必须先将pairRDD转换为RowRDD。

那么如何从 PairRDD 创建 RowRdd 呢?

【问题讨论】:

    标签: java apache-spark spark-dataframe rdd java-pair-rdd


    【解决方案1】:

    对于 Java 7,您需要定义一个映射函数

    public static final Function<Tuple2<String, String>,Row> mappingFunc = (tuple) -> {
        return RowFactory.create(tuple._1(),tuple._2());
    };
    

    现在你可以调用这个函数来获取JavaRDD&lt;Row&gt;

    JavaRDD<Row> rowRDD = filesRDD.map(mappingFunc);
    

    Java 8 就像

    JavaRDD<Row> rowRDD = filesRDD.map(tuple -> RowFactory.create(tuple._1(),tuple._2()));
    

    从 JavaPairRDD 获取 Dataframe 的另一种方法是

    DataFrame df = sqlContext.createDataset(JavaPairRDD.toRDD(filesRDD), Encoders.tuple(Encoders.STRING(),Encoders.STRING())).toDF();
    

    【讨论】:

      【解决方案2】:

      以下是实现此目的的一种方法。

          //Read whole files
          JavaPairRDD<String, String> pairRDD = sparkContext.wholeTextFiles(path);
      
          //create a structType for creating the dataframe later. You might want to
          //do this in a different way if your schema is big/complicated. For the sake of this
          //example I took a simple one.
          StructType structType = DataTypes
                  .createStructType(
                          new StructField[]{
                                  DataTypes.createStructField("id", DataTypes.StringType, true)
                                  , DataTypes.createStructField("name", DataTypes.StringType, true)});
      
      
          //create an RDD<Row> from pairRDD
          JavaRDD<Row> rowJavaRDD = pairRDD.values().flatMap(new FlatMapFunction<String, Row>() {
              public Iterable<Row> call(String s) throws Exception {
                  List<Row> rows = new ArrayList<Row>();
                  for (String line : s.split("\n")) {
                      String[] values = line.split(",");
                      Row row = RowFactory.create(values[0], values[1]);
                      rows.add(row);
                  }
                  return rows;
              }
          });
      
      
          //Create Dataframe.
          sqlContext.createDataFrame(rowJavaRDD, structType);
      

      我使用的示例数据
      文件 1:

      1, john  
      2, steve
      

      文件2:

      3, Mike  
      4, Mary  
      

      df.show() 的输出:

      +---+------+
      | id|  name|
      +---+------+
      |  1|  john|
      |  2| steve|
      |  3|  Mike|
      |  4|  Mary|
      +---+------+
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2017-05-09
        • 2017-03-25
        • 2016-01-23
        • 2017-03-17
        • 2016-01-05
        • 2022-01-08
        • 1970-01-01
        相关资源
        最近更新 更多