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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 11:28:24