Spark DF写入Parquet报错:列word_c类型不匹配(String/INT64)
问题:指定Schema读取Parquet后写入时出现类型不匹配报错
在Databricks中读取分散在不同文件夹的Parquet文件,指定了全为StringType的Schema,读取代码如下:
df = spark.read.option("mergeSchema", "true").schema(parquet_schema).parquet('/mnt/my_blobstorage/snap/*/*.parquet')
通过display(df)和df.printSchema()确认所有列均为StringType,但执行写入命令:
df.write.parquet('/mnt/my_blobstorage/saved/merged_df.parquet')
时出现报错:
Parquet column cannot be converted. Column: [word_c], Expected: StringType, Found: INT64
解决思路
- 强制重铸所有列为StringType:
遍历DataFrame所有列,显式转换类型,彻底消除隐式类型残留:from pyspark.sql.functions import col from pyspark.sql.types import StringType df_cast = df.select([col(c).cast(StringType()).alias(c) for c in df.columns]) df_cast.printSchema() # 再次确认类型 df_cast.write.parquet('/mnt/my_blobstorage/saved/merged_df.parquet') - 关闭
mergeSchema选项:
指定自定义Schema时开启mergeSchema会导致Spark尝试合并底层文件的原始Schema,与你指定的Schema冲突。关闭后Spark会严格按指定Schema读取:df = spark.read.schema(parquet_schema).parquet('/mnt/my_blobstorage/snap/*/*.parquet') - 排查单个文件夹的Parquet文件:
单独读取包含word_c列的可疑文件夹文件,查看其原始类型,确认是否存在转换异常:test_df = spark.read.parquet('/mnt/my_blobstorage/snap/目标文件夹/*.parquet') test_df.select("word_c").printSchema() - 写入时显式指定Schema:
写入时强制指定目标Schema,避免Spark自动推断类型:df.write.option("schema", parquet_schema.json()).parquet('/mnt/my_blobstorage/saved/merged_df.parquet') - 检查分区列影响:
如果word_c是分区列,其类型信息存储在文件夹路径元数据中,指定Schema可能无法覆盖,需要单独对分区列做类型转换。
内容的提问来源于stack exchange,提问作者Olgaraa
相关产品推荐
相关产品推荐

