Spark读取S3分区Parquet文件列类型不一致报错求助
问题描述
我的S3上有如下格式的分区Parquet目录:
database/year=2022/mth=1/id=1/file1.parquet database/year=2022/mth=12/id=2/file2.parquet database/year=2022/mth=3/id=3/file3.parquet database/year=2022/mth=8/id=4/file4.parquet
这些文件中部分列的数据类型不一致,例如file1的impression列是double类型,file2的该列是decimal类型,其他部分列也存在类似问题。
我尝试用以下Scala代码读取数据:
sparkSession.read() .option("header", true) .option("basePath", basePath) .schema(schema) .parquet(basePath + input) .filter(col("year").equalTo(key).and(col("mth").isin(value.toArray())) .and(col("id").isin(idList.toArray())));
但报错:Column: [impressions], Expected: bigint, Found: DOUBLE
我已经试过mergeSchema参数和读取时指定schema,但都没解决问题,求可行的解决办法。
解决方案
1. 分文件读取并手动转换类型
跳过全局schema校验,逐个读取文件后统一转换到目标类型,再合并数据:
// 枚举所有目标文件路径 val filePaths = List( s"${basePath}/database/year=2022/mth=1/id=1/file1.parquet", s"${basePath}/database/year=2022/mth=12/id=2/file2.parquet", s"${basePath}/database/year=2022/mth=3/id=3/file3.parquet", s"${basePath}/database/year=2022/mth=8/id=4/file4.parquet" ) // 定义最终需要的schema val targetSchema = StructType(Seq( StructField("impression", LongType, nullable = true), // 补充其他列的类型定义 StructField("year", StringType, nullable = true), StructField("mth", StringType, nullable = true), StructField("id", StringType, nullable = true) )) // 逐个读取文件、转换类型后合并 val finalDF = filePaths.foldLeft(spark.emptyDataFrame) { (acc, path) => val rawDF = spark.read.parquet(path) // 按目标schema转换每一列的类型 val convertedDF = rawDF.select( targetSchema.fields.map { field => col(field.name).cast(field.dataType).alias(field.name) }: _* ) acc.unionByName(convertedDF) } // 应用过滤条件 val filteredDF = finalDF.filter( col("year") === key && col("mth").isin(value.toArray: _*) && col("id").isin(idList.toArray: _*) )
2. 调整Spark配置放宽校验
修改Spark配置,强制开启schema合并并禁用严格的类型校验:
// 配置Spark参数 sparkSession.conf.set("spark.sql.parquet.mergeSchema", "true") sparkSession.conf.set("spark.sql.parquet.failOnCorruptFile", "false") sparkSession.conf.set("spark.sql.parquet.enableVectorizedReader", "false") // 读取数据(不指定schema,先让Spark自动合并) val mergedDF = sparkSession.read .option("basePath", basePath) .parquet(basePath + input) // 手动转换到目标类型 val targetDF = mergedDF.select( col("impression").cast(LongType).alias("impression"), // 其他列同理转换 col("year"), col("mth"), col("id") ) // 应用过滤 targetDF.filter( col("year") === key && col("mth").isin(value.toArray: _*) && col("id").isin(idList.toArray: _*) )
3. 预处理文件统一类型
如果长期存在类型不一致问题,建议一次性预处理所有文件,将类型统一后再读取:
// 获取所有分区的唯一组合 val partitions = spark.read.parquet(basePath + input) .select("year", "mth", "id") .distinct() .collect() // 遍历每个分区,读取后转换类型并覆盖写回 partitions.foreach { row => val year = row.getAs[String]("year") val mth = row.getAs[String]("mth") val id = row.getAs[String]("id") val partitionPath = s"${basePath}/database/year=${year}/mth=${mth}/id=${id}" val rawDF = spark.read.parquet(partitionPath) // 按目标schema转换所有列 val convertedDF = rawDF.select( targetSchema.fields.map { field => col(field.name).cast(field.dataType).alias(field.name) }: _* ) // 覆盖写回原分区路径 convertedDF.write.mode("overwrite").parquet(partitionPath) } // 之后就可以正常读取了 sparkSession.read .option("basePath", basePath) .schema(targetSchema) .parquet(basePath + input) .filter(...)
内容的提问来源于stack exchange,提问作者Neethu Lalitha
相关产品推荐
相关产品推荐

