【问题标题】:Cast Spark dataframe’s schemaCast Spark 数据框架构
【发布时间】:2021-12-15 23:34:38
【问题描述】:

我有一个具有以下架构的数据框:

StructType currentSchema = new StructType(new StructField[]{
    new StructField("age", DataTypes.StringType, false, Metadata.empty()),
    new StructField("grade", DataTypes.StringType, false, Metadata.empty()),
    new StructField("dateOfBirth", DataTypes.StringType, false, Metadata.empty())
});

我想立即将它(不指定每一列)转换为以下架构:

StructType newSchema = new StructType(new StructField[]{
     new StructField("age", DataTypes.IntegerType, false, Metadata.empty()),
     new StructField("grade", DataTypes.IntegerType, false, Metadata.empty()),
     new StructField("dateOfBirth", DataTypes.DateType, false, Metadata.empty())
});

有什么办法可以做到df.convert(newSchema)这样的操作吗?

【问题讨论】:

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


    【解决方案1】:

    由于 DataFrame 是不可变的,您必须创建新的 DataFrame 来更改架构。为此,请执行以下方法之一:

    我:

     Dataset<Row> ndf = df.select(col("age").cast(DataTypes.IntegerType),
                                  col("grade").cast(DataTypes.IntegerType),
                                  col("dateOfBirth").cast(DataTypes.DateType));        
     ndf.printSchema();
    

    二:

    或(我只为年龄列做了):

    Dataset<Row> ndf = df.withColumn("new_age", df.col("age").cast(DataTypes.IntegerType)).drop("age");
    ndf.printSchema();
    

    III:

    最后但同样重要的是,使用 map 函数同时进行操作和更改类型:

    Dataset<Row> df2 = df.map(new MapFunction<Row, Row>() {
                @Override
                public Row call(Row row) throws Exception {
                    return RowFactory.create((int)row.getString(0),
                                             (int)row.getString(1),
                                             (date)row.getString(2));
                }
            }, RowEncoder.apply(newSchema));
    
    df2.printSchema();
    

    在此方法中,如果转换 (int) 不起作用,请使用 Integer.Parse 代替。

    【讨论】:

    • 我喜欢第三个选项的想法,但有没有办法动态地做到这一点?不需要指定每列的类型和未知的列号?
    • 我强烈推荐第三个。您可以选择更改行值以及数据类型。例如想要转换年龄并添加一个新列。您可以在map 函数中同时更改类型。实际上,在第一次转换操作中更改列的类型。
    • 听起来不错,你能在答案上编辑一下,我会接受吗?
    • 我认为你自己编辑它的分数很好。但是你认为我必须在这篇文章中改变什么?
    【解决方案2】:

    一种方法是让 spark 将所有列转换为您期望的新类型。我不确定它是否适用于所有类型的转换,但它适用于许多情况:

    List<Column> columns = Arrays
        .stream(newSchema.fields())
        .map(field -> col(field.name()).cast(field.dataType()))
        .collect(Collectors.toList());
    
    Dataset<Row> newResult = result.select(columns.toArray(new Column[0]));
    
    

    另一种方法是依靠 spark 将模式应用于 csv 文件的方式,但这需要将数据写入磁盘,因此我不推荐该选项。

    result.write().csv("somewhere");
    Dataset<Row> newResult = spark.read().schema(newSchema).csv("somewhere");
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-02-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-09-21
      • 1970-01-01
      • 2015-08-26
      • 2017-12-23
      相关资源
      最近更新 更多