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

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:
    # 写入时关闭Parquet字典编码
    result_df.write.mode("overwrite").option("parquet.enable.dictionary", "false").parquet('path')
    
    该方案对写入性能影响小于10%,写入后的数据即使开启矢量化读取也不会出现值错位。
  • 读取侧规避(历史数据修复可选):如果已经生成的历史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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 05:39:40