Databricks读取Parquet写入Delta表遇Schema转换不支持错误求助
解决方案
方法1:禁用Parquet向量化读取
向量化读取对部分逻辑类型与物理类型不匹配的字段处理存在问题,禁用后Spark会以传统方式解析字段:
# 临时禁用向量化读取 spark.conf.set("spark.sql.parquet.enableVectorizedReader", "false") # 重新读取Parquet数据 df = spark.read.format("parquet").load("/mnt/path") # 按需将SubID转换为目标类型(比如StringType) from pyspark.sql.functions import col from pyspark.sql.types import StringType df = df.withColumn("SubID", col("SubID").cast(StringType())) # 写入Delta表 df.write.format("delta").mode("overwrite").saveAsTable(path)
方法2:自定义匹配物理类型的Schema读取
既然Parquet中SubID的物理类型是INT64,先以LongType读取,再转换为需要的类型:
from pyspark.sql.types import StructType, StructField, StringType, LongType # 定义Schema,SubID设为LongType匹配物理存储 custom_schema = StructType([ StructField('GatewayID', StringType(), True), StructField('Version', StringType(), True), StructField('Generation', StringType(), True), StructField('AppVersion', StringType(), True), StructField('UnitEntity', StringType(), True), StructField('SubID', LongType(), True), StructField('SIMCardID', StringType(), True), StructField('UnitNumber', StringType(), True), StructField('UnitType', StringType(), True), StructField('ISOCountryCode', StringType(), True), StructField('ReportTime', LongType(), True), StructField('MessageFormat', StringType(), True), StructField('MessagesAsString', StringType(), True) ]) # 用自定义Schema读取数据 df = spark.read.format("parquet").schema(custom_schema).load("/mnt/path") # 转换为StringType(如果Delta表需要字符串类型) df = df.withColumn("SubID", col("SubID").cast(StringType())) # 写入Delta表 df.write.format("delta").mode("overwrite").saveAsTable(path)
方法3:检查并修复Parquet元数据
如果Parquet文件的元数据本身存在错误(物理类型与逻辑类型不匹配),可以用Parquet-tools工具查看元数据:
parquet-tools schema /path/to/your/parquet/file
确认字段的实际类型后,要么修复Parquet文件的元数据,要么调整读取逻辑适配。
内容的提问来源于stack exchange,提问作者Saswat Ray
相关产品推荐
相关产品推荐

