简介此份基于MongoDB与Spark的大数据项目资料包完整包含项目文档、源码与优秀项目范例面向计算机相关专业人工智能、通信工程、自动化、电子信息等的在校学生和教师适用于毕业设计、课程设计、作业或项目初期演示。压缩包共81个文件以Java源码为主含60个Java文件、6个JSP页面及多个XML配置文件同时包含JS、CSS、Markdown说明与项目配置文件覆盖后端逻辑、前端展示与配置说明的完整链路整体仅375KB结构清晰便于查阅。已有59人学习下载。该项目为高分源码已获导师认可答辩评审分95分代码经测试运行成功。读者可借此理解MongoDB与Spark的整合实践在现成代码基础上修改扩展也可直接用于毕设、课设或作为项目模板。1. 一个资源包拆开的 MongoDBSpark 离线链路网上标注「基于 MongoDBSpark」的资源包解开后一般包含部署文档、样例工程和数据文件但真正值钱的是那条从 MongoDB 到 Spark 的离线分析链路。MongoDB 负责接住订单、行为日志这类半结构化业务数据Spark 负责把这些数据变成指标和宽表。为什么是这两个组件组合MongoDB 写入快、文档模型改造成本低但统计聚合能力弱Spark 擅长分布式批处理和复杂计算却不适合当业务主库。两个一拼正好补齐。下面按工程顺序拆开讲集合结构怎么设计、PySpark 读写参数怎么设、集群和内存怎么调、复购率怎么算最后落在核对与排错。适合正在做大数据方向毕设和入门项目的人也适合想快速搞清这套组合工程边界的开发者。2. MongoDB 文档建模与索引先设计结构再交给 Spark 扫描2.1 集合怎么设计才能少一个 join 环节做离线分析时MongoDB 集合承担的是事实表角色存储订单或点击日志。设计上优先考虑 Spark 读取时怎么少 join而不是怎么省存储。订单集合最稳的做法是把下单时的商品快照直接嵌进 items 数组用户相关只留 uid用户的其他维度信息单独放 users 集合。这样分析一条订单时不需要回查商品表商品改名改价不影响历史订单的统计值。{ order_id: no_20250101_0001, uid: usr_1024, shop_id: 12, status: paid, channel: android, items: [ {sku: a01, name: 蓝牙耳机, price: 99.0, qty: 1} ], paid_at: 2025-01-01T10:11:00Z }这个文档的要点是价格字段。商品表中 name 和 price 都会变动订单分析要的是成交瞬间的值所以 price 快照必须冗余在 order 文档内。如果引用商品表出报表时要么 join 商品主表要么把商品表历史版本也保存一份工作量马上大起来。MongoDB 的文档模型天然接受这种冗余和关系型数据库三范式思路不一样。反面典型是把关系库的关联表设计搬进来订单拆成主表加明细表Spark 读取时每次都要多路 join还会在业务服务里出现循环查子集合的 N1 查询数据访问层一旦写出findAll()一次性捞全集合再逐条关联量上去后服务和数据库都扛不住。2.2 索引怎么建才能配合 Spark 扫描对 Spark 读取来说全表扫描会被 MongoDB 端的光标分块处理索引的价值在于两点让下推的$match少扫文档让排序不要回表。考虑订单集合最常用的过滤是状态和时间窗口对应复合索引。db.orders.createIndex({status: 1, paid_at: -1}, {background: true}) db.orders.createIndex({uid: 1}, {background: true})两个命令先建复合索引再建 uid 单键索引。paid_at用 -1 是为了配合“最近 N 天”这类倒序过滤请求索引方向选错时 MongoDB 会走内存排序数据量大时分析任务开销明显上升。uid索引主要服务 Spark 端 join 用户维度时的点查。资源包里经常补一堆单列索引分析链路上用不到反而拖慢写入。索引类型典型场景注意点复合索引状态时间过滤、排序字段顺序要匹配查询条件单字段索引用户关联、分组基数太低可以选择性放弃TTL 索引行为日志自动过期字段必须是时间类型不能建在数组上哈希索引分片键随机分布只支持等值查询不支持范围过滤TTL 索引在日志类集合上很实用比如只保留 30 天点击流过期字段设置为 expireAtMongoDB 后台线程会自动清理。哈希索引这一行说清楚分片集合里如果没有合适的高基数字段可以拿_id做哈希分片代价是该字段上的范围查询会退化成广播式扫描所以分析任务里要避免在哈希字段上做范围过滤。2.3 把聚合下推给 MongoDB而不是全量拉到 SparkPySpark 读 MongoDB 前先在 MongoDB 侧做一层筛选和裁剪是最容易见效的优化。连接器允许在读取时指定 aggregation pipelineMongoDB 执行完 pipeline 再交给 Spark 分区网络传输量能小一个数量级。pipeline [ {$match: {status: paid, paid_at: {$gte: 2025-01-01}}}, {$project: {uid: 1, amount: 1, channel: 1}} ] df MongoSpark.read(spark) \ .option(database, bigdata) \ .option(collection, orders) \ .option(aggregation.pipeline, pipeline) \ .load()pipeline 里的字段名是 MongoDB 物理字段名和 Spark 侧列名不是一个体系$project里别写 DataFrame 的列名。还可以往 pipeline 里放$group把求和、计数在 MongoDB 端先做掉Spark 收到的就是汇总结果。什么时候不该下推呢如果分组字段基数很高、结果集接近原始数据规模下推意义就不大反而多一次计算。更简单的场景用 MongoDB 原生$merge直接落结果集合连 Spark 都不用启动这条边界经常被忽略。3. PySpark 读 MongoDB 的 ETL 链路参数、Schema 与写回3.1 最小读配置和分区参数怎么给读取集合前先把 uri、database、collection 三个参数配齐。用 SparkSession 级配置写一次全任务都能用也可以只在 DataFrameReader 上指定。推荐 session 级配置因为写回时还要复用同一套连接信息。from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(mongo-etl) \ .config(spark.mongodb.read.connection.uri, mongodb://127.0.0.1:27017) \ .config(spark.mongodb.read.database, bigdata) \ .config(spark.mongodb.read.collection, orders) \ .config(spark.mongodb.read.partitioner, MongoSamplePartitioner) \ .config(spark.mongodb.read.partitionerOptions.partitionSizeMB, 64) \ .config(spark.mongodb.read.readPreference, secondaryPreferred) \ .getOrCreate() df MongoSpark.read(spark).load() df.printSchema()partitionSizeMB 决定单个 Spark 分区处理的数据量默认 64MB。如果发现任务最后几个 stage 长尾严重、部分分区数据倾斜把 64 降到 32 或 16 会得到更细的切片反过来小集合没必要调小分区粒度太细时任务调度开销超过计算本身。readPreference 显式设成 secondaryPreferred让 Spark 扫描走从节点避免把主节点 IO 打满影响线上写入。参数作用建议spark.mongodb.read.connection.uri连接串副本集写全节点列表spark.mongodb.read.database数据库与业务环境隔离spark.mongodb.read.collection集合名只读目标集合spark.mongodb.read.partitioner分区器分布不均时换采样分区器partitionSizeMB分区大小64 起步倾斜调小aggregation.pipeline下推过滤与聚合能下推就下推readPreference读偏好分析任务用 secondaryPreferreduri 里如果有用户名密码注意密码含特殊字符时要先做 URL 编码mongodb://user:p%40sshost:27017用原文会解析失败。连接串写多个节点时用逗号分隔MongoDB 驱动会自己探测主从。3.2 Schema 推断为什么会翻车连接器默认按采样文档推断 DataFrame 的 schema采样量小的时候推断结果和真实数据结构经常对不上。常见情况集合里 99% 的 amount 是 Double采样恰好抽到几条整数字段推断成 LongTypepaid_at 有的记录是字符串、有的是日期推断出的类型不稳定某个字段只有采样未覆盖的文档才有值被推断成全部可空。显式声明 schema 是绕过这些坑最稳的办法。from pyspark.sql.types import DoubleType, StringType, StructField, StructType schema StructType([ StructField(order_id, StringType(), False), StructField(uid, StringType(), True), StructField(amount, DoubleType(), True), StructField(channel, StringType(), True), StructField(paid_at, StringType(), True) ]) df MongoSpark.read(spark) \ .option(database, bigdata) \ .option(collection, orders) \ .schema(schema) \ .load()这里的 paid_at 故意声明成 StringType。MongoDB 的日期经连接器转成 Spark timestamp 时时区处理容易出偏差反正后面 ETL 里还要按业务时区转换干脆在读取阶段保持字符串原样清洗时统一用to_timestamp转换。声明 schema 还带来一个额外好处字段顺序固定下来下游代码不会因为抽样文档的字段排序变化被带偏。提示连接器的 schema 推断基于采样显式声明 schema 的任务在数据分布变化后行为更稳定。3.3 写回 MongoDBSaveMode 和 _id 的语义ETL 结果写回 MongoDB 时最容易踩的坑是mode(overwrite)会把整个集合并掉重来如果集合还在被线上使用会造成瞬间不可读。按主键更新的需求要用 append 加 DataFrame 自带_id字段实现。result_df df.groupBy(uid).count() result_df.write \ .format(mongodb) \ .mode(append) \ .option(database, bigdata) \ .option(collection, repurchase_result) \ .save()连接器写入不自动生成业务主键_id不存在时 MongoDB 会自动分配 ObjectId重跑任务就变成重复插入。要幂等写前把业务主键映射到_id字段upsert_df result_df.withColumn(_id, F.col(uid)) upsert_df.write.format(mongodb).mode(append).save()_id用字符串是合法的MongoDB 要求_id唯一且非数组整数、字符串、UUID 都行。映射之后同一任务里相同 uid 会互相覆盖吗不会同一集合里重复_id会报主键冲突所以任务本身要先按 uid 聚合去重再写。重跑想整集合替换先在 Spark 侧计算目标结果再对对应集合做删除后重写。写模式实际语义适用场景append追加文档无主键冲突的增量结果overwrite删除集合再写入一次性全量重建ignore集合存在则跳过初始化临时表errorifexists集合存在报错目录或表存在检查4. 集群部署与内存调优MongoDB 分片与 Spark Executor 参数4.1 先副本集后分片分片键别偷懒自定义MongoDB 7.x 的安装部署无论单机还是集群官方安装手册都建议先初始化副本集再考虑分片。副本集解决高可用分片解决容量和写入吞吐。先搭副本集再开分片顺序反了会经历一次全量数据搬迁集群越大越折腾。部署配置用 YAML 更直观。systemLog: destination: file path: /data/mongodb/log/mongod.log logAppend: true storage: dbPath: /data/mongodb/db net: port: 27017 replication: replSetName: rs0 sharding: clusterRole: shardsvr这个配置同时是分片节点所以带sharding.clusterRole: shardsvr只做复制的话把这一段删掉。部署过程中 Windows 上的安装失败大多报The installer has encountered an unexpected error常见处理方式是管理员身份运行安装包彻底卸载旧版本清理安装目录和数据目录残留再看系统服务里有没有遗留的 MongoDB 服务占用端口。装完先用 MongoDB Compass 连接一次确认服务正常再进集群初始化。副本集初始化后在 mongos 上执行分片sh.enableSharding(bigdata) sh.shardCollection(bigdata.orders, {_id: hashed}){_id: hashed}是哈希分片订单_id单调递增范围分片会持续写最后一个分片哈希分片把写请求散到各个分片。分片键选择记住三条基数尽量高、写入不能只集中在单个分片、常见查询要能带这个字段。拿同一个字段做排序或范围过滤的场景就要重新评估是否适合哈希分片。4.2 Spark 安装与集群搭建的最小三步Spark 集群搭建在测试环境里用 Standalone 模式就够了。三台机器一台 master 两台 worker解压 tar 包后配置好 Java 环境三步启动$SPARK_HOME/sbin/start-master.sh $SPARK_HOME/sbin/start-worker.sh spark://master-host:7077 spark-submit \ --master spark://master-host:7077 \ --executor-memory 4g \ --executor-cores 2 \ job.pystart-master.sh负责启动调度节点worker 启动参数必须传 master 地址和端口 7077。spark-submit上的--executor-memory 4g是每个 executor 的堆内存--executor-cores 2控制每个 executor 占用 CPU 核数两个参数一起决定单个节点能起多少个 executor。不建议一味堆 executor-memory内存超过物理内存一半后 GC 停顿反而变大。4.3 先调 Spark 内存的四个参数Spark 任务常见的故障是在读完全量数据后做 groupByshuffle 阶段 OOM。报错形式通常是 executor 被 kill。真遇到内存溢出先看是执行内存不足还是存储内存不足再动参数。参数默认值调优方向spark.executor.memory1g4g 到 8g配合物理内存spark.memory.fraction0.6执行段压力大时可到 0.75spark.memory.storageFraction0.5频繁缓存 DataFrame 时加大spark.sql.shuffle.partitions200按 128MB 每分区估算spark.memory.fraction决定执行和存储共享堆的比例剩下的留给用户代码内部的临时对象。spark.memory.storageFraction是共享内存里存储段的比例如果计算链路依赖缓存 DataFrame 复用storageFraction 调大如果任务是纯 shuffle 重计算调小能给执行段更多空间。spark.sql.shuffle.partitions默认 200 对几 GB 的测试数据偏大数据量不大时改成 64 或 128 更稳。5. 用 Spark DataFrame 实现用户复购率与订单宽表5.1 复购率口径先定清楚再去重复购率的定义网上有很多版本有的按支付次数、有的按购买天数写 Spark 脚本之前先明确口径周期内购买次数大于等于 2 的用户数除以周期内有过购买行为的用户总数订单状态取 paid。同一天下多单是否算复购结果差异很大这里按天去重处理。from pyspark.sql.functions import col, count, date_format paid df.filter(col(status) paid) user_buy_days paid \ .select(uid, date_format(col(paid_at), yyyy-MM-dd).alias(buy_date)) \ .dropDuplicates([uid, buy_date]) user_buy_cnt user_buy_days.groupBy(uid).count() total_users user_buy_cnt.count() repurchase_users user_buy_cnt.filter(col(count) 2).count() rate repurchase_users / total_users print(f复购率: {rate:.4f})dropDuplicates([uid, buy_date])是关键同一用户同一天多笔订单只计一次口径是“购买天数大于等于 2”。如果业务口径按支付笔数算去掉去重直接groupBy(uid).count()即可。最后算比率时注意整数除法PySpark 的count()返回的是 int直接相除没问题但封装成指标输出时保留四到五位小数方便后续环比。5.2 从订单事实表生成用户宽表复购率只是其中一个指标数据分析项目里常把用户相关指标合成一张宽表。订单表按 uid 聚合出频次、金额、最近付款时间再和用户维度表做左连接得到的宽表每行一个用户后续做用户标签查询都不再碰订单集合。from pyspark.sql.functions import avg, max, round, sum user_agg paid.groupBy(uid).agg( count(order_id).alias(order_cnt), round(sum(amount), 2).alias(total_amount), round(avg(amount), 2).alias(avg_amount), max(paid_at).alias(last_paid_at) ) user_dim spark.read.parquet(/data/user_dim.parquet) wide_table user_dim.join(user_agg, uid, left) wide_table.write \ .mode(overwrite) \ .partitionBy(channel) \ .parquet(/data/user_wide)left join 保证没有消费的用户也保留在宽表里后面做人群筛选不会丢人。写入时按 channel 字段分区这个字段来自 user_dim 表如果渠道类别很少写出的文件数很有限查询时能直接做分区裁剪。宽表字段设计上order_cnt、total_amount、avg_amount 这类聚合指标是稳定列last_paid_at 时间列是后续做 RFM 模型的输入。user_dim 如果是几十万行几 MB 的维度表默认的 sort-merge join 会在 shuffle 上浪费资源给 user_dim 加一个 broadcast hint整个 join 就会变成 map 端 join少一个 stage。5.3 数据规模与 shuffle 分区怎么对应Spark 任务在集群上提交前先按集合大小估算分区数。一个订单集合 2GB连接器默认 64MB 一个分区读出来约 32 个分区数据是 20GB分区数约 320。shuffle 阶段的分区数要跟数据量匹配不是越大越好。数据规模读取分区数参考shuffle partitions2GB326420GB320400200GB32004000表格里的分区数按 64MB 每分区估算实际执行时还要考虑集群并行度executor 数乘每个 executor cores得到同时运行的任务数量。shuffle partitions 设成这个数的 2 到 4 倍是常见做法太高会有一堆空任务太低会看到一个 task 吃几 GB 数据。6. 验证链路连接器版本、幂等写入与数据核对6.1 连接器和 Spark 版本不匹配的最常见报错资源包源码在本地跑通、到集群上报NoSuchMethodError绝大多数和业务代码无关是 MongoDB Spark Connector 与 Spark 主版本不匹配。官方连接器的构建会标明支持哪个 Spark 版本线下载时选与集群 Spark 主版本一致的构建包用 Maven 管理时把连接器依赖版本固定下来临时换版本前先看依赖树。mvn dependency:tree | grep mongodbdependency:tree会列出所有传递依赖。连接器版本对不上时先确认项目中是否同时引入两套连接器坐标冲突时排除旧坐标保留对应 Spark 主版本的那一份。修完依赖后重新提交任务先看 driver 日志开头的 Spark 和连接器版本行。6.2 幂等写入检验数两边行数还不够任务重跑后验证结果只比对源表和结果集合的行数远远不够要核对业务主键是否重复以及关键指标的汇总值是否一致。用 pymongo 做一次快速核对脚本是常见做法。from pymongo import MongoClient client MongoClient(mongodb://127.0.0.1:27017) target client[bigdata][repurchase_result] src_count source_count_from_spark dst_count target.count_documents({}) assert src_count dst_count, fcount mismatch src{src_count} dst{dst_count}这段只验证了行数接下来还要统计_id冲突率也就是目标集合里重复_id的比例。金额类指标用 MongoDB 聚合把合计算出来再和 Spark 输出对比两边独立实现相同口径数值对不上时优先检查金额过滤条件是否一致。6.3 三条命令完成收尾核对落到报表和线上使用前用 mongosh 跑三条核对命令mongosh bigdata --eval db.repurchase_result.find({$or:[{uid:null},{amount:{$lt:0}}]}).count() mongosh bigdata --eval db.orders.aggregate([{$match:{status:paid}},{$group:{_id:null,total:{$sum:$amount}}}]).toArray() mongosh bigdata --eval db.repurchase_result.aggregate([{$group:{_id:null,total:{$sum:$amount}}}]).toArray()第一条看空值和负数异常第二条在源集合算金额合计第三条在结果集合算合计。两条合计对不上时回到 5.2 的 left join 语义检查。跑通这三条命令后把任务挂到调度器定时重跑每天结果集合先清空再写入_id冲突率归零MNGO 到 Spark 这条链路才算真正稳定。本文还有配套的精品资源点击获取