对DataFrame调用map()方法是否总会将其转换为RDD?
对Spark DataFrame调用
map()方法是否总会将其转换为RDD? 没错,只要你对Spark DataFrame调用map()方法,最终都会得到一个RDD,下面结合你的示例来具体解释:
首先,Spark DataFrame本质上是带有结构化Schema的分布式数据集,属于Spark SQL的API体系,但底层仍基于RDD实现。而map()是RDD API中的经典转换操作,当你在DataFrame上调用这个方法时,相当于把DataFrame中的每一行(Row对象)传入自定义的转换逻辑,最终输出的结果会脱离DataFrame的Schema约束,成为一个普通的RDD。
看你给出的代码示例:
首先加载得到一个DataFrame:
scala> val custDF = sqlContext.read.format("com.databricks.spark.avro").load("/user/cloudera/practice1/problem7/customer/avro") custDF: org.apache.spark.sql.DataFrame = [customer_id: int, customer_fname: string, customer_lname: string]
这里custDF的类型明确是DataFrame,包含结构化的字段信息。
当你调用map()进行行转换后:
scala> val a = custDF.map(x=>x(0)+"\t"+x(1)+"\t"+x(2)) a: org.apache.spark.rdd.RDD[String] = MapPartitionsRDD[106] at map at <console>:36
结果a的类型变成了RDD[String],这直接验证了map()方法将DataFrame转换为了RDD。
补充一句:如果想在保留Dataset/DataFrame特性(比如Schema、优化执行计划)的前提下进行行转换,在Spark 2.x及以后的版本中,你可以使用Dataset API的map()方法(需要提供对应的Encoder),这样会返回一个Dataset[T]而不是RDD,但这和你示例中调用的RDD风格map()是不同的场景。
内容的提问来源于stack exchange,提问作者Sudhanshu
相关产品推荐
相关产品推荐

