【问题标题】:How to add a schema to a Dataset in Spark?如何将模式添加到 Spark 中的数据集?
【发布时间】:2017-07-07 16:11:29
【问题描述】:

我正在尝试将文件加载到 spark 中。 如果我将普通的 textFile 加载到 Spark 中,如下所示:

val partFile = spark.read.textFile("hdfs://quickstart:8020/user/cloudera/partfile")

结果是:

partFile: org.apache.spark.sql.Dataset[String] = [value: string]

我可以在输出中看到一个数据集。但是如果我加载一个 Json 文件:

val pfile = spark.read.json("hdfs://quickstart:8020/user/cloudera/pjson")

结果是一个带有现成架构的数据框:

pfile: org.apache.spark.sql.DataFrame = [address: struct<city: string, state: string>, age: bigint ... 1 more field]

Json/parquet/orc 文件具有架构。所以我可以理解这是 Spark version:2x 的一个特性,它使事情变得更容易,因为在这种情况下我们直接获得了一个 DataFrame,而对于一个普通的 textFile,你会得到一个没有有意义的模式的数据集。 我想知道的是如何将模式添加到数据集,这是将 textFile 加载到 spark 中的结果。对于 RDD,有 case class/StructType 选项来添加模式并将其转换为 DataFrame。 谁能告诉我该怎么做?

【问题讨论】:

    标签: apache-spark


    【解决方案1】:

    当您使用textFile 时,文件的每一行都将是您数据集中的一个字符串行。要转换为带有 schema 的 DataFrame,可以使用toDF:

    val partFile = spark.read.textFile("hdfs://quickstart:8020/user/cloudera/partfile")
    
    import sqlContext.implicits._
    val df = partFile.toDF("string_column")
    

    在这种情况下,DataFrame 将有一个 StringType 类型的单列架构。

    如果您的文件包含更复杂的架构,您可以使用 csv 阅读器(如果文件是结构化 csv 格式):

    val partFile = spark.read.option("header", "true").option("delimiter", ";").csv("hdfs://quickstart:8020/user/cloudera/partfile")
    

    或者您可以使用 map 处理您的数据集,然后使用 toDF 转换为 DataFrame。例如,假设您希望一列是该行的第一个字符(作为 Int),另一列是第四个字符(也作为 Int):

    val partFile = spark.read.textFile("hdfs://quickstart:8020/user/cloudera/partfile")
    
    val processedDataset: Dataset[(Int, Int)] = partFile.map {
      line: String => (line(0).toInt, line(3).toInt)
    }
    
    import sqlContext.implicits._
    val df = processedDataset.toDF("value0", "value3")
    

    此外,您可以定义一个案例类,它将代表您的 DataFrame 的最终架构:

    case class MyRow(value0: Int, value3: Int)
    
    val partFile = spark.read.textFile("hdfs://quickstart:8020/user/cloudera/partfile")
    
    val processedDataset: Dataset[MyRow] = partFile.map {
      line: String => MyRow(line(0).toInt, line(3).toInt)
    }
    
    import sqlContext.implicits._
    val df = processedDataset.toDF
    

    在上述两种情况下,调用df.printSchema 会显示:

    root
     |-- value0: integer (nullable = true)
     |-- value3: integer (nullable = true)
    

    【讨论】:

    • 根据你的回答,我不得不稍微调整一下。根据分隔符拆分数据集: val partdata = partFile.map(p => p.split(",")) 我还必须更改此语句: val prdt = partdata.map{line => rows(line(0 ).toInt, line(1).toString, line(2).toInt, line(3).toString, line(4).toString)} 因为非数字数据是 'char' 格式,我必须转换他们到'字符串'。它现在正在工作。
    • @Sidhartha,很高兴知道它有效。如果是逗号分隔的文件,可以考虑我的第一个建议,使用spark.read.csv,可能更简单。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-12-10
    • 2017-07-29
    • 2011-11-07
    • 1970-01-01
    • 2011-03-10
    • 2017-02-10
    • 1970-01-01
    相关资源
    最近更新 更多