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
相关产品推荐
相关产品推荐

