【发布时间】:2017-05-17 06:55:24
【问题描述】:
我想从文本文件创建数据框。
案例类别限制为 22 个字符;我有 100 多个字段。
因此我在创建案例类时遇到问题。
我的实际目标是创建 Dataframe;
有没有其他方法可以创建Dataframe,而不是使用Case Class?
【问题讨论】:
标签: scala apache-spark dataframe
我想从文本文件创建数据框。
案例类别限制为 22 个字符;我有 100 多个字段。
因此我在创建案例类时遇到问题。
我的实际目标是创建 Dataframe;
有没有其他方法可以创建Dataframe,而不是使用Case Class?
【问题讨论】:
标签: scala apache-spark dataframe
一种方法是使用 spark csv 包直接读取文件并创建数据框。如果您的文件有标头,包将直接从标头推断架构,或者您可以使用结构类型创建自定义架构。
在下面的示例中,我创建了一个自定义架构。
val sqlContext = new SQLContext(sc)
val customSchema = StructType(Array(
StructField("year", IntegerType, true),
StructField("make", StringType, true),
StructField("model", StringType, true),
StructField("comment", StringType, true),
StructField("blank", StringType, true)))
val df = sqlContext.read
.format("com.databricks.spark.csv")
.option("header", "true") // Use first line of all files as header
.schema(customSchema)
.load("cars.csv")
val df = sqlContext.read
.format("com.databricks.spark.csv")
.option("header", "true") // Use first line of all files as header
.option("inferSchema", "true") // Automatically infer data types
.load("cars.csv")
您可以在databricks spark csv documentation page 上查看其他各种选项。
其他选项:
您可以使用如上所示的struct type创建一个schema,然后使用sqlContext的createDataframe创建dataframe。
val vRdd = sc.textFile(..filelocation..)
val df = sqlContext.createDataframe(vRdd,schema)
【讨论】:
当无法提前定义案例类时(例如,记录的结构被编码为字符串,或者将解析文本数据集并针对不同的用户以不同的方式投影字段),可以通过编程方式创建 DataFrame三个步骤。
StructType 表示的架构,以匹配在步骤 1 中创建的 RDD 中的行结构。SQLContext 提供的createDataFrame 方法将模式应用于Rows 的RDD。另一种方法是在StructType 中使用datatyoe 定义StructField。它将允许您定义多种数据类型。请参阅下面的示例以了解这两种实施方式。请考虑注释代码以了解这两种实现。
package com.spark.examples
import org.apache.spark._
import org.apache.spark.sql.SQLContext
import org.apache.spark.sql._
import org.apache.spark._
import org.apache.spark.sql.DataFrame
import org.apache.spark.rdd.RDD
import org.apache.spark.sql._
import org.apache.spark.sql.types._
// Import Row.
import org.apache.spark.sql.Row;
// Import Spark SQL data types
import org.apache.spark.sql.types.{ StructType, StructField, StringType }
object MultipleDataTypeSchema extends Serializable {
val conf = new SparkConf().setAppName("schema definition")
conf.set("spark.executor.memory", "100M")
conf.setMaster("local")
val sc = new SparkContext(conf);
// sc is an existing SparkContext.
val sqlContext = new org.apache.spark.sql.SQLContext(sc)
def main(args: Array[String]): Unit = {
// Create an RDD
val people = sc.textFile("C:/Users/User1/Documents/test")
/* First Implementation:The schema is encoded in a string, split schema then map it.
* All column dataype will be string type.
//Generate the schema based on the string of schema
val schemaString = "name address age" //Here you can read column from a preoperties file too.
val schema =
StructType(
schemaString.split(" ").map(fieldName => StructField(fieldName, StringType, true)));*/
// Second implementation: Define multiple datatype
val schema =
StructType(
StructField("name", StringType, true) ::
StructField("address", StringType, true) ::
StructField("age", StringType, false) :: Nil)
// Convert records of the RDD (people) to Rows.
val rowRDD = people.map(_.split(",")).map(p => Row(p(0), p(1).trim, p(2).trim))
// Apply the schema to the RDD.
val peopleDataFrame = sqlContext.createDataFrame(rowRDD, schema)
peopleDataFrame.printSchema()
sc.stop
}
}
它的输出:
17/01/03 14:24:13 INFO SparkContext: Created broadcast 0 from textFile at MultipleDataTypeSchema.scala:30
root
|-- name: string (nullable = true)
|-- address: string (nullable = true)
|-- age: string (nullable = false)
【讨论】:
通过 sqlContext 的 sqlContext.read.csv() 方法读取文件效果很好。因为它有许多可用的内置方法,您可以在其中传递参数并控制执行。但是在 1.6 之前的 spark 版本上工作可能没有这个可用。所以你也可以通过 spark-context 的 textFile 方法来实现。
Val a = sc.textFile("file:///file-path/fileName")
这会给你一个 RDD[String]。因此,您现在已经创建了 RDD,并且希望将其转换为数据框。
现在继续使用 StructTypes 为您的 RDD 定义架构。这允许您拥有尽可能多的 StructFields。
val schema = StructType(Array(StructField("fieldName1", fieldType, ifNullablle),
StructField("fieldName2", fieldType, ifNullablle),
StructField("fieldName3", fieldType, ifNullablle),
................
))
您现在有两件事:1) RDD,我们使用 textFile 方法创建的。 2) Schema,具有所需数量的属性。
下一步肯定是将此模式与您的 RDD 正确映射! 您可能会观察到您拥有的 RDD 是单个字符串,即 RDD[String]。但是您实际上想要做的是将其转换为您为其创建模式的许多变量。那么为什么不根据逗号分割你的RDD。下面的表达式应该使用映射操作来做到这一点。
val b = a.map(x => x.split(","))
你在评估中得到一个 RDD[Array[String]]。
但是你可能会说这个 Array[String] 仍然不是那么直观,我可以应用任何操作。 因此,您可以使用 Row API。使用 import org.apache.spark.sql.Row 将其导入 我们实际上会将您拆分的 RDD 与 Row 对象映射为一个元组。看到这个:
import org.apache.spark.sql.Row
val c = b.map(x => Row(x(0), x(1),....x(n)))
上面的表达式为您提供了一个 RDD,其中每个元素都是一个 Row。你现在只需要给它一个模式。再次,sqlContext 的 createDataFrame 方法为您完成了这项工作。
val myDataFrame = sqlContext.createDataFrame(c, schema)
这个方法有两个参数:1)你需要处理的RDD。 2)您要在其之上应用的架构。 结果评估是 DataFrame 对象。 所以最后我们现在创建了我们的 DataFrame 对象 myDataFrame。如果你在 myDataFrame 上使用 show 方法,你可以看到表格格式的数据。 您现在可以对其执行任何 spark-sql 操作了。
【讨论】: