Spark读取Mongo聚合结果文档数异常,重复数据问题排查
问题原因分析
你遇到的核心问题是Spark MongoDB连接器没有在Mongo服务器端完整执行你的聚合管道,而是在Spark的每个数据分区上单独执行了聚合逻辑。
因为Mongo集合的数据会被Spark分成多个分区读取,当聚合管道被推送到每个分区执行时:
- 同一个用户的日志记录可能分散在不同分区中
- 每个分区单独执行
$group操作,会生成该分区内的用户-电影ID集合 - 最终合并所有分区结果时,就会出现重复的
user_id,且部分记录的movie_ids仅包含该分区内的电影(甚至可能因为分区数据问题出现错误的电影ID) - 这就导致Spark读取后的总记录数远大于Mongo端执行聚合的192条。
解决方案
针对这个问题,你可以通过以下几种方式修复:
1. 强制使用单分区读取聚合结果
设置Spark读取Mongo时使用单分区,确保聚合管道在Mongo服务器端全局执行,而不是分区级别执行。修改SparkSession的配置:
SparkSession spark = SparkSession .builder() .appName("LuckyBetsFP") .config("spark.mongodb.read.connection.uri", "mongodb://localhost:27017/exercises.movie_logs") .config("spark.mongodb.read.partitioner", "com.mongodb.spark.sql.connector.read.partitioner.MongoSinglePartitioner") .getOrCreate();
2. 禁用分区级聚合推送
如果使用的是Mongo Spark Connector 10.x及以上版本,检查是否开启了分区级聚合推送,将其关闭:
Dataset<Row> items = spark .read() .format("mongodb") .option("aggregation.pipeline", MOVIES_PIPELINE) .option("spark.mongodb.read.aggregate.partitionPushdown", "false") .load();
3. 验证连接器版本兼容性
确保你使用的Mongo Spark Connector版本与MongoDB服务器版本兼容。比如MongoDB 5.0+建议使用Connector 10.x,旧版本的连接器可能存在聚合管道推送的逻辑bug,导致管道未在服务器端完整执行。
验证方法
修改配置后,重新运行作业,检查结果的user_id唯一性和总记录数是否为192。也可以在Mongo端开启查询日志,确认聚合管道是否被完整发送到Mongo服务器执行。
内容的提问来源于stack exchange,提问作者zaro
相关产品推荐
相关产品推荐

