【问题标题】:Spark: How can DataFrame be Dataset[Row] if DataFrame's have a schemaSpark:如果 DataFrame 有模式,DataFrame 如何成为 Dataset[Row]
【发布时间】:2017-02-16 07:49:06
【问题描述】:

This article 声称 Spark 中的 DataFrame 等价于 Dataset[Row],但 this blog post 表明 DataFrame 具有架构。

以将RDD转换为DataFrame的博客文章中的示例:如果DataFrameDataset[Row]相同,那么将RDD转换为DataFrame应该很简单

val rddToDF = rdd.map(value => Row(value))

但它反而表明它是这个

val rddStringToRowRDD = rdd.map(value => Row(value))
val dfschema = StructType(Array(StructField("value",StringType)))
val rddToDF = sparkSession.createDataFrame(rddStringToRowRDD,dfschema)
val rDDToDataSet = rddToDF.as[String]

很明显,数据框实际上是由行和模式组成的数据集。

【问题讨论】:

    标签: scala apache-spark apache-spark-sql apache-spark-dataset


    【解决方案1】:

    在 Spark 2.0 中,代码中有: type DataFrame = Dataset[Row]

    它是Dataset[Row],只是因为定义。

    Dataset 也有模式,您可以使用printSchema() 函数打印它。通常 Spark 会推断模式,因此您不必自己编写它 - 但是它仍然存在;)

    您也可以使用createTempView(name) 并在 SQL 查询中使用它,就像 DataFrames 一样。

    换句话说,Dataset = DataFrame from Spark 1.5 + encoder,将行转换为您的类。在 Spark 2.0 中合并类型后,DataFrame 只是Dataset[Row] 的别名,因此没有指定编码器。

    关于转换:rdd.map() 也返回 RDD,它从不返回 DataFrame。你可以这样做:

    // Dataset[Row]=DataFrame, without encoder
    val rddToDF = sparkSession.createDataFrame(rdd)
    // And now it has information, that encoder for String should be used - so it becomes Dataset[String]
    val rDDToDataSet = rddToDF.as[String]
    
    // however, it can be shortened to:
    val dataset = sparkSession.createDataset(rdd)
    

    【讨论】:

      【解决方案2】:

      注意(除了T Gaweda 的答案)每个Row (Row.schema) 都有一个关联的架构。但是,在将其集成到 DataFrame(或 Dataset[Row])中之前,不会设置此架构

      scala> Row(1).schema
      res12: org.apache.spark.sql.types.StructType = null
      
      scala> val rdd = sc.parallelize(List(Row(1)))
      rdd: org.apache.spark.rdd.RDD[org.apache.spark.sql.Row] = ParallelCollectionRDD[5] at parallelize at <console>:28
      scala> spark.createDataFrame(rdd,schema).first
      res15: org.apache.spark.sql.Row = [1]
      scala> spark.createDataFrame(rdd,schema).first.schema
      res16: org.apache.spark.sql.types.StructType = StructType(StructField(a,IntegerType,true))
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2017-04-17
        • 2019-02-16
        • 2017-04-03
        • 1970-01-01
        • 1970-01-01
        • 2018-10-30
        • 1970-01-01
        相关资源
        最近更新 更多