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

Spark 3.3.0读取Go生成Parquet文件时查询执行计划报错

Spark 3.3.0读取Go生成Parquet过滤报错的解决思路

问题现象

使用Spark 3.3.0读取Go服务生成的Parquet文件时,执行包含特定字段过滤的Scala查询(如isClone条件过滤)会触发以下错误:

java.io.IOException: can not read class org.apache.parquet.format.PageHeader: Socket is closed by peer.

移除该过滤条件后查询可正常执行;部分Parquet文件中对parent_vol_id字段的过滤也会引发同类错误,且肉眼检查Parquet数据无明显异常。

可能原因

  • Parquet格式兼容性差异:Go语言的Parquet库(如xitongsys/parquet-go)与Spark依赖的Apache Parquet 1.12.x版本在PageHeader的编码/解码逻辑上存在不一致,当查询触发列裁剪(Column Pruning)时,Spark读取特定列的PageHeader时因格式不匹配导致连接被远端关闭。
  • 字段编码或统计信息异常:isClone、parent_vol_id这类字段在Go端可能使用了非标准编码(比如布尔值存储格式、字符串编码),或者字段的统计元数据(min/max值、空值计数)写入错误,Spark执行过滤时读取这些元数据触发报错。
  • 隐性文件损坏:肉眼可见的数据正常,但文件的Page结构或元数据存在隐性损坏,仅在读取特定列时暴露问题。

解决办法

1. 验证Parquet文件结构

使用parquet-tools工具检查文件元数据和目标列的结构,对比标准Parquet文件的差异:

# 查看文件元数据
parquet-tools meta /path/to/your-file.parquet
# 导出目标列数据和结构
parquet-tools dump --column isClone /path/to/your-file.parquet

2. 调整Spark读取参数(临时验证)

  • 禁用列裁剪,强制读取所有列(仅用于验证,不推荐生产环境长期使用):
spark.read.option("parquet.enable.column.pruning", "false").parquet("/path/to/parquet-files")
  • 禁用谓词下推,避免Spark提前读取过滤字段的元数据:
spark.read.option("parquet.pushdown.filter.enabled", "false").parquet("/path/to/parquet-files")

3. 修复Go端Parquet生成逻辑

  • 更换兼容性更好的Go Parquet库,比如segmentio/parquet-go,这类库严格遵循Apache Parquet规范。
  • 确保字段编码符合标准:布尔值使用1bit存储,字符串采用UTF-8编码,字段统计信息(min、max、空值数)正确生成。

4. 重新生成Parquet文件

用Spark读取无过滤条件的原始Parquet数据,重新写入一份标准格式的Parquet文件,再执行过滤查询:

val df = spark.read.parquet("/path/to/original-files")
df.write.mode("overwrite").parquet("/path/to/new-standard-files")

内容的提问来源于stack exchange,提问作者sagar nyamagouda

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 22:52:50