PySpark DataFrame保存为Parquet重读后RECORD_ID去重计数异常
问题根因
不是Parquet无法解析特殊字符导致行覆盖,你遇到的是Spark 3.2.1版本Parquet矢量化读取器的已知bug,触发条件刚好匹配你的场景:
- 数据集包含存储超大文本的
raw_text列,写入Parquet时默认对该列启用字典编码 - 读取时默认开启的矢量化Parquet读取器在解析字典编码的大字符串列时,会出现字典页偏移计算错误,导致列值跨行错误复用
- 这个bug只会导致单列值串位、重复,不会改变文件内存储的总行数,和你观察到的「总行数和写入前一致、但RECORD_ID去重后数量骤降」的现象完全吻合
你自定义的ASCII处理UDF不是核心诱因,但Python UDF的序列化开销会放大写入时的页对齐异常概率。
快速验证
执行以下代码关闭矢量化读取器后重读数据,若去重计数恢复为213051即可实锤根因:
# 临时关闭矢量化Parquet读取 spark.conf.set("spark.sql.parquet.enableVectorizedReader", "false") verify_df = spark.read.parquet('path') verify_df.select(count_distinct("RECORD_ID")).show()
验证时可以抽几个重复的RECORD_ID查看对应的raw_text值,会明显看到文本内容混杂了其他行的片段,是典型的列值错位特征。
解决方案
按落地成本从低到高选择:
- 集群升级(推荐):将Azure Databricks集群升级到11.3 LTS及以上版本(对应Spark 3.3+),该版本已经正式修复了大字符串列的Parquet矢量化读取错位bug,不需要修改任何业务代码。
- 写入侧规避(不升级集群可选):写入Parquet时关闭字典编码,从根源避免触发bug:
该方案对写入性能影响小于10%,写入后的数据即使开启矢量化读取也不会出现值错位。# 写入时关闭Parquet字典编码 result_df.write.mode("overwrite").option("parquet.enable.dictionary", "false").parquet('path') - 读取侧规避(历史数据修复可选):如果已经生成的历史Parquet文件无法重写,就在读取作业中全局关闭矢量化Parquet读取配置即可拿到正确结果,代价是Parquet读取性能会下降20%左右。
- UDF优化:把你自定义的Python ASCII处理UDF替换为Spark内置函数,避免Python序列化带来的额外异常风险,性能还能提升5~10倍:
from pyspark.sql import functions as F # 替换原有ascii_udf逻辑,移除所有非ASCII字符 result_df = result_df.withColumn("raw_text", F.regexp_replace(F.col("raw_text"), r"[^\x00-\x7F]", ""))
额外排查项
- 写入时必须显式指定
mode("overwrite"),提前清空目标路径下的残留文件,避免历史任务生成的Parquet文件被一并读取导致数据异常。
内容的提问来源于stack exchange,提问作者Watershipdown
相关产品推荐
相关产品推荐

