Spark 2.2.1:如何保存含已知Schema的空DataFrame并写入Schema信息?
解决Spark保存空DataFrame时保留Schema的问题
当然可以实现保存带已知Schema的空DataFrame,让后续读取时能正确识别Schema!你的代码问题在于:当DataFrame完全没有数据行时,默认情况下Spark的DataFrameWriter不会写入任何Parquet元数据文件或空数据文件,导致读取时无法推断Schema。
下面给你两种实用的解决方案:
方案1:使用CTAS(Create Table As Select)语句(兼容所有Spark版本)
通过SQL的CTAS语句创建表,即使是空表,Spark也会强制写入Schema元数据到指定路径。示例代码如下:
import org.apache.spark.sql.{SparkSession, Row} import org.apache.spark.sql.types.StructType import org.apache.spark.sql.SaveMode def example(spark: SparkSession, path: String, schema: StructType) = { // 创建空DataFrame val emptyDF = spark.createDataFrame(spark.sparkContext.emptyRDD[Row], schema) // 注册临时视图 emptyDF.createOrReplaceTempView("temp_empty_view") // 用CTAS语句创建Parquet表,写入Schema元数据 spark.sql(s""" CREATE TABLE IF NOT EXISTS empty_parquet_table USING PARQUET LOCATION '$path' AS SELECT * FROM temp_empty_view """) // 现在读取路径下的文件,Schema会被正确识别 val readDF = spark.read.parquet(path) readDF.printSchema() // 可以正常输出预设的Schema }
方案2:开启Spark空数据集写入配置(Spark 3.1+适用)
Spark 3.1及以上版本提供了一个配置项spark.sql.parquet.writeEmptyFileWhenNoData,开启后会在保存空DataFrame时自动写入一个空的Parquet文件,这样读取时就能通过这个空文件获取Schema:
import org.apache.spark.sql.{SparkSession, Row} import org.apache.spark.sql.types.StructType import org.apache.spark.sql.SaveMode def example(spark: SparkSession, path: String, schema: StructType) = { // 开启写入空文件的配置 spark.conf.set("spark.sql.parquet.writeEmptyFileWhenNoData", "true") val dataframe = spark.createDataFrame(spark.sparkContext.emptyRDD[Row], schema) dataframe.write.mode(SaveMode.Overwrite).format("parquet").save(path) // 读取操作现在可以正常获取Schema val readDF = spark.read.load(path) readDF.printSchema() // 正确显示预设的Schema }
为什么原代码会失败?
当DataFrame没有任何数据行时,默认配置下Spark不会生成任何Parquet数据文件,也不会写入_metadata或_common_metadata这类保存Schema信息的元数据文件。后续读取时,Spark找不到任何能推断Schema的文件,就会抛出无文件可读、Schema未知的错误。而上面的两种方案,要么通过CTAS强制写入元数据,要么通过配置生成空文件来携带Schema信息,完美解决了这个问题。
内容的提问来源于stack exchange,提问作者Erik
相关产品推荐
相关产品推荐

