You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.25 15:55:14