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

使用PySpark读取S3上MongoDB的bson.gz文件报pymongo_spark缺失错误

问题场景

基于Jupyter Notebook使用PySpark读取S3路径下扩展名为.bson.gz的MongoDB快照文件时抛出模块不存在报错,使用的测试代码如下:

from pyspark import SparkContext, SparkConf
sc.install_pypi_package("pymongo==3.2.2")

import pymongo_spark
pymongo_spark.activate()
conf = SparkConf().setAppName("pyspark-bson")
file_path = "s3://location/users.bson.gz"
bsonFileRdd = sc.BSONFileRDD(file_path)
bsonFileRdd.take(5)

运行后抛出的核心错误信息:

ModuleNotFoundError: No module named 'pymongo_spark'
报错原因
  • 仅安装了pymongo依赖,未安装对应Spark连接器的Python包pymongo-spark,导致导入失败。
  • 指定安装的pymongo==3.2.2版本过旧,与当前新版Spark、pymongo-spark连接器存在API兼容性问题。
  • 代码中调用的sc.BSONFileRDD是旧版连接器废弃的API,新版环境下不存在该方法,即使解决依赖问题也会触发二次报错。
  • 旧版连接器默认不支持.bson.gz格式的自动解压识别,直接读取会出现解析失败。
可行解决步骤

1. 安装匹配版本的依赖

先终止已初始化的SparkContext,避免依赖安装后不生效,再安装兼容版本的依赖包:

# 终止已存在的Spark上下文
try:
    sc.stop()
except:
    pass

from pyspark import SparkContext, SparkConf
# 安装兼容版本依赖
sc.install_pypi_package("pymongo>=4.0.0")
sc.install_pypi_package("pymongo-spark")

2. 修正读取代码,配置压缩支持

使用新版连接器官方API,添加gzip压缩编解码配置,确保压缩格式的BSON文件可被正常解析:

import pymongo_spark
pymongo_spark.activate()

conf = SparkConf()\
    .setAppName("pyspark-bson")\
    .set("spark.hadoop.io.compression.codecs", "org.apache.hadoop.io.compress.GzipCodec")
sc = SparkContext(conf=conf)

file_path = "s3://location/users.bson.gz"
# 调用新版API读取BSON文件
bson_rdd = sc.mongoBsonFile(file_path)
# 测试读取前5条数据
bson_rdd.take(5)

备选方案(连接器安装失败时使用)

如果集群环境无法正常安装pymongo-spark包,可直接通过Python原生bson库手动解析压缩文件,无需依赖Mongo Spark连接器:

# 安装所需依赖
sc.install_pypi_package("bson")

import gzip
import bson

# 以二进制格式读取S3上的压缩文件
raw_rdd = sc.binaryFiles("s3://location/users.bson.gz")

def parse_bson_compressed(file_part):
    _, content = file_part
    decompressed_content = gzip.decompress(content)
    return bson.decode_all(decompressed_content)

parsed_rdd = raw_rdd.flatMap(parse_bson_compressed)
# 测试读取前5条
parsed_rdd.take(5)
注意事项
  • 若S3路径存在访问权限控制,需提前在Spark配置中添加S3鉴权参数,或给运行Notebook的节点绑定对应S3读权限的IAM角色。
  • 读取GB级以上大体积快照文件时,建议先做重分区操作,避免单节点内存占用过高触发OOM。
  • 不要固定安装3.x版本的pymongo,会与新版pymongo-spark出现签名不兼容问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 14:39:08