Spark创建空DataFrame报错:不支持org.apache.spark.sql.types.DataType类型
问题:创建指定Schema的空DataFrame报错
代码尝试
val sparkConf = new SparkConf() .setAppName("app") .setMaster("local") val sparkSession = SparkSession .builder() .config(sparkConf) .getOrCreate() val sparkContext = sparkSession.sparkContext var tmpScheme = StructType( StructField("source_id", StringType, true) :: Nil) var df = conf.SparkConf.sparkSession.createDataFrame(tmpScheme)
报错信息
Schema for type org.apache.spark.sql.types.DataType is not supported ...
环境信息
- 导入语句:
import org.apache.spark.sql.types.{BooleanType, IntegerType, StringType, StructField, StructType}
- Spark版本:3.2.2
"org.apache.spark" %% "spark-sql" % "3.2.2", // spark "org.apache.spark" %% "spark-core" % "3.2.2", // spark
问题原因及解决方法
核心问题有两个:
SparkSession引用错误
你写的conf.SparkConf.sparkSession是完全错误的引用方式——代码里已经初始化好sparkSession变量,直接使用这个变量即可,不需要通过SparkConf去获取。createDataFrame方法调用错误
Spark的createDataFrame没有直接接收Schema作为唯一参数的重载方法。当你只传入StructType(它是DataType的子类)时,Spark无法匹配到合适的方法重载,因此抛出"不支持DataType类型"的错误。这和你有没有显式导入DataType无关,因为StructType本身就继承自DataType,Spark内部会识别这个类型。
正确的实现方式
需要传入一个空的数据源(比如空RDD或空序列),再结合指定的Schema来创建空DataFrame,以下是两种常用方法:
方法1:使用空RDD创建
import org.apache.spark.sql.Row val emptyRDD = sparkContext.emptyRDD[Row] val df = sparkSession.createDataFrame(emptyRDD, tmpScheme)
方法2:使用空序列创建
import org.apache.spark.sql.Row val df = sparkSession.createDataFrame(Seq.empty[Row], tmpScheme)
如果场景允许使用类型安全的方式,也可以定义样例类后创建空DataFrame:
case class Source(source_id: String) val df = sparkSession.createDataFrame(Seq.empty[Source])
内容的提问来源于stack exchange,提问作者palandlom
相关产品推荐
相关产品推荐

