MongoDB+Spark大数据项目实战:从环境搭建到读写闭环
简介这份资源是围绕 MongoDB 与 Spark 构建的大数据项目完整资料包面向计算机、人工智能、通信工程、自动化、电子信息、物联网等专业的在校学生、教师及企业员工可用于毕业设计、课程设计、作业提交或项目初期立项演示也适合具备一定基础的小白进阶学习。压缩包共 81 个文件约 375KB以 60 个 Java 源码文件为核心配合 6 个 JSP 页面、4 个 XML 配置、2 个 JavaScript 脚本及 CSS、classpath、project 等工程文件另附 README 说明文档与项目截图整体结构清晰便于按模块阅读与二次开发。该资源为个人高分项目已通过导师指导与答辩评审评分达 95 分代码均经测试运行成功后才上传功能可用性有保障。目前已有 60 人学习关注读者可借此掌握 MongoDB 与 Spark 的整合思路、项目分层组织方式及关键配置要点并在此基础上修改扩展实现更多业务功能。1. 从一份 MongoDBSpark 资料包说起大数据项目到底该怎么落地很多人拿到「基于 MongoDBSpark 的大数据项目文档源码优秀项目全部资料」这类压缩包时第一反应是解压、找 README、跑 main 函数然后卡在环境上——MongoDB 安装失败、Spark 集群搭建报错、依赖版本对不上最后项目躺在硬盘里吃灰。这个标题背后其实是一条完整的技术链路用 MongoDB 存原始与结果数据用 Spark 做分布式计算中间靠 Connector 打通读写外面再套一层数据清洗、分析与可视化。它适合正在做大数据毕业设计、想补一个端到端项目经验或者需要把「MongoDB 数据库基本操作」和「Spark 数据分析案例」串起来的开发者。这篇笔记不讲空泛的大数据介绍而是按「数据怎么进 MongoDB → Spark 怎么读 → 怎么算 → 结果怎么回写 → 怎么验证」的顺序把每一步的命令、参数和翻车点讲清楚让你拿到任何一份同类资料包都能自己跑通。2. 环境先立住MongoDB 与 Spark 的安装边界2.1 为什么先锁版本再谈安装大数据项目翻车十次有八次死在版本上。MongoDB 的 Spark Connector 对两端版本有明确矩阵Connector 10.x 通常要求 MongoDB 4.4 和 Spark 3.xConnector 3.x 则对应 Spark 2.x。如果你拿到的资料包里pom.xml写的是mongo-spark-connector_2.11:3.0.1却装了 Spark 3.5编译能过运行必炸报NoSuchMethodError或ClassNotFoundException。我一般会先做一件事打开资料包的pom.xml或build.sbt把spark.version、scala.binary.version、mongo-spark-connector三个坐标抄下来再去装对应版本。Scala 2.11 配 Spark 2.4Scala 2.12 配 Spark 3.x这条线不能乱。MongoDB 这边社区版足够用重点确认两件事一是bindIp默认127.0.0.1只允许本机连Spark 集群里其他节点连不上二是认证如果开了authorization: enabled连接串必须带用户名密码和authSource。很多「MongoDB 安装失败」的帖子其实是服务起不来根源在/var/lib/mongodb权限或dbPath目录不存在跟安装包本身没关系。2.2 Linux 上装 MongoDB 的最小命令集以 Ubuntu 22.04 为例用官方 apt 源安装避免手动解压带来的库依赖问题# 导入 MongoDB 公钥并添加源以 7.0 为例按需换版本 curl -fsSL https://www.mongodb.org/static/pgp/server-7.0.asc | \ sudo gpg -o /usr/share/keyrings/mongodb-server-7.0.gpg --dearmor echo deb [ archamd64,arm64 signed-by/usr/share/keyrings/mongodb-server-7.0.gpg ] \ https://repo.mongodb.org/apt/ubuntu jammy/mongodb-org/7.0 multiverse | \ sudo tee /etc/apt/sources.list.d/mongodb-org-7.0.list sudo apt-get update sudo apt-get install -y mongodb-org # 启动并设为开机自启 sudo systemctl enable mongod sudo systemctl start mongod sudo systemctl status mongod逻辑说明第一段导入官方签名密钥保证包来源可信第二段写入 apt 源jammy对应 Ubuntu 22.04其他系统要换代号第三段安装的是mongodb-org元包会带上 server、shell、mongos 等组件。参数上-y免交互适合脚本化。装完必须status看一眼active (running)才算过。如果失败先看journalctl -u mongod -n 50九成是dbPath权限或端口 27017 被占。验证连接用mongoshmongosh --host 127.0.0.1 --port 27017 # 进入后执行 use bigdata_project db.stats()use在库不存在时不会立刻创建第一次插入数据才落盘这是 MongoDB 的懒创建特性别以为没生效。2.3 Spark 本地模式与 Standalone 集群的取舍做项目验证我强烈建议先用local[*]模式跑通逻辑再上集群。本地模式不需要改spark-env.sh一条命令就能起# 解压后进入目录直接提交任务 tar -zxvf spark-3.5.1-bin-hadoop3.tgz cd spark-3.5.1-bin-hadoop3 ./bin/spark-submit \ --master local[*] \ --class com.example.MongoSparkDemo \ --jars mongo-spark-connector_2.12-10.3.0.jar \ target/mongo-spark-demo-1.0.jarlocal[*]表示用本机所有 CPU 核--jars把 Connector 塞进 classpath避免ClassNotFoundException: com.mongodb.spark。如果资料包里用的是 Maven shade 打包Connector 已经打进去了--jars可省。要搭 Standalone 集群最少改两个文件conf/spark-env.sh里设SPARK_MASTER_HOST为本机 IPconf/workers里写 worker 节点主机名。然后sbin/start-all.sh。注意SPARK_MASTER_HOST别写localhost否则 worker 注册不上UI 里看不到节点。集群起来后提交任务把--master换成spark://master:7077并确保每个节点都能访问 MongoDB 的 27017 端口。提示MongoDB 和 Spark 混部在同一台机器做实验时内存要留够。Spark executor 默认吃 1gMongoDB 的 WiredTiger 缓存默认是内存的一半两者抢内存会导致 OOM 或频繁换页。3. 数据链路打通Spark 读写 MongoDB 的三种姿势3.1 Connector 的读写原理与配置项MongoDB Spark Connector 的本质是把 MongoDB 的集合映射成 Spark DataFrame读的时候按分区并行拉取写的时候按分区并行写入。它不走 JDBC而是用 MongoDB Java Driver 的聚合游标分批取数所以性能比 JDBC 好但配置项也更多。核心配置分两类连接类spark.mongodb.read.connection.uri/spark.mongodb.write.connection.uri和分区类spark.mongodb.read.partitioner、partition.size等。连接串格式是mongodb://user:passhost:port/database.collection?authSourceadmin。注意database.collection是点号分隔不是斜杠。如果开了认证authSource必须写对用户建在admin库就写admin建在业务库就写业务库名写错会报Authentication failed。读分区策略有三种DEFAULT按_id范围切、SHARDED分片集群用、SAMPLE采样切分。单机或副本集用DEFAULT就行它会根据_id的 min/max 自动算分区边界。如果集合里_id分布极不均匀比如全是递增插入DEFAULT会切出很多空分区这时候用SAMPLE更稳。3.2 用 Python 跑通第一个读写闭环资料包里如果是 Scala 项目逻辑一样这里用 PySpark 写方便新手直接抄from pyspark.sql import SparkSession # 构建 SparkSession同时注入 MongoDB 读写配置 spark SparkSession.builder \ .appName(MongoSparkDemo) \ .config(spark.mongodb.read.connection.uri, mongodb://127.0.0.1:27017/bigdata_project.raw_orders) \ .config(spark.mongodb.write.connection.uri, mongodb://127.0.0.1:27017/bigdata_project.result_orders) \ .config(spark.mongodb.read.partitioner, DEFAULT) \ .config(spark.mongodb.read.partition.size, 64) \ .getOrCreate() # 读把集合加载成 DataFrame df spark.read.format(mongodb).load() df.printSchema() df.show(5, truncateFalse) # 算按状态分组统计金额过滤掉取消单 from pyspark.sql.functions import col, sum as _sum result df.filter(col(status) ! cancelled) \ .groupBy(status) \ .agg(_sum(amount).alias(total_amount)) # 写结果回写 MongoDB覆盖模式 result.write.format(mongodb) \ .mode(overwrite) \ .option(replaceDocument, false) \ .save() spark.stop()逻辑说明SparkSession构建时把连接串写进 config读和写用不同的 URI避免读结果写回原集合造成脏数据。partition.size单位是 MB64 表示每个分区大约 64MB 数据太小分区多、调度开销大太大单分区内存压力大一般 32128 之间调。replaceDocumentfalse表示写入时只更新匹配文档的字段不整文档替换保留原有_id和其他字段。如果设true会按_id整文档覆盖容易丢字段。mode(overwrite)在 Connector 里的行为是删除目标集合再重建不是逐条覆盖生产环境慎用。想增量更新用append或update但update需要指定idFieldList。3.3 参数怎么调分区、批大小与内存读写性能调优先看三个参数参数作用建议值调错后果spark.mongodb.read.partition.size读分区大小(MB)32128太小任务碎片化太大 executor OOMspark.mongodb.write.batch.size写批条数5122048太小网络往返多太大单批内存高spark.mongodb.read.partitioner分区策略单机 DEFAULT分片 SHARDED选错导致数据倾斜或空分区写批大小默认是 512如果文档很大比如带二进制字段要往下调否则单批内存爆。读分区大小默认 64MB如果集合只有几万条小文档切出来可能就一个分区并行度上不去这时候手动设小一点比如 8让 Spark 多切几个任务。Spark 内存这边spark.executor.memory和spark.mongodb.read.partition.size要联动。假设 executor 给 4g分区 128MB一个 executor 同时跑 4 个 task就是 512MB 数据在内存里加上 shuffle 和对象开销勉强够。如果 executor 只有 2g分区还设 128必 OOM。我一般按「executor 内存 / 分区大小 ≥ 8」来估留足余量。注意MongoDB 侧也有内存压力。Spark 大批量读会占用 MongoDB 的连接和游标资源mongod的maxIncomingConnections默认 65536一般够但游标超时cursorTimeoutMillis默认 10 分钟大分区读太久会被服务端杀游标报CursorNotFound。遇到这个把partition.size调小或调大服务端超时。4. 从原始集合到分析结果一个可复现的清洗与聚合流程4.1 数据清洗的四个必做动作资料包里的原始数据往往很脏直接聚合会得到离谱结果。我一般固定做四件事去重、补空、类型转换、异常值过滤。以订单数据为例from pyspark.sql.functions import col, when, to_timestamp, trim, lower from pyspark.sql.types import DoubleType # 1. 去重按 order_id 保留最新一条 df df.dropDuplicates([order_id]) # 2. 补空amount 为空填 0status 为空填 unknown df df.fillna({amount: 0.0, status: unknown}) # 3. 类型转换字符串时间转 timestamp金额转 double df df.withColumn(created_at, to_timestamp(col(created_at), yyyy-MM-dd HH:mm:ss)) \ .withColumn(amount, col(amount).cast(DoubleType())) # 4. 异常值过滤金额为负或超过 100000 的剔除 df df.filter((col(amount) 0) (col(amount) 100000)) # 5. 文本规范化状态统一小写去空格 df df.withColumn(status, lower(trim(col(status))))逻辑说明dropDuplicates按指定列去重比distinct()精准后者是全字段比对脏数据字段稍有差异就判为不同。fillna传字典可以按列填不同默认值。to_timestamp的格式串必须和源数据完全匹配yyyy是年MM是月dd是日大小写敏感写成YYYY会得到错误年份。cast转 double 时非数字字符串会变 null所以要在过滤前做或者用when(...).otherwise(...)兜底。清洗完建议先cache()再聚合避免重复计算df.cache() print(清洗后条数:, df.count())cache()把 DataFrame 物化到内存后续多次 action 不用重算。但数据量大时别乱 cache会挤占 executor 内存反而拖慢。4.2 聚合与回写把结果落回 MongoDB清洗后的数据按业务维度聚合结果写回 MongoDB 供前端或可视化用from pyspark.sql.functions import date_format, count, avg, max as _max # 按天和状态聚合 daily df.groupBy( date_format(col(created_at), yyyy-MM-dd).alias(day), col(status) ).agg( count(*).alias(order_count), avg(amount).alias(avg_amount), _max(amount).alias(max_amount) ) # 回写用 append 保留历史按 daystatus 做幂等需在业务层处理 daily.write.format(mongodb) \ .mode(append) \ .option(spark.mongodb.write.connection.uri, mongodb://127.0.0.1:27017/bigdata_project.daily_stats) \ .save()date_format把 timestamp 转成日期字符串方便按天分组。count(*)统计行数avg和max分别算均值和最大值。回写用append每次跑批追加一批如果重复跑会产生重复数据生产上要么用overwrite按天覆盖要么在 MongoDB 侧建唯一索引{day:1, status:1}做 upsert。Connector 的update模式支持按idFieldList做 upsert但配置稍复杂新手先用overwrite验证逻辑。写完后去mongosh里验证use bigdata_project db.daily_stats.find().sort({day: -1}).limit(10) db.daily_stats.aggregate([ { $group: { _id: $status, total: { $sum: $order_count } } } ])这一步别省很多「Spark 跑成功但没数据」的情况是写到了错误的库或集合或者mode选错把数据清了。4.3 用 Spark SQL 做日期维度分析大数据 SQL 面试题里高频出现日期加减、月份提取Spark SQL 的写法要记牢-- 注册临时视图 CREATE OR REPLACE TEMP VIEW orders AS SELECT * FROM df; -- 日期加年、加减天数、转月份 SELECT created_at, add_months(created_at, 12) AS plus_1_year, date_add(created_at, 7) AS plus_7_days, date_sub(created_at, 1) AS minus_1_day, date_format(created_at, yyyy-MM) AS month_str FROM orders WHERE created_at IS NOT NULL;add_months加月份date_add/date_sub加减天数date_format提取月份字符串。注意date_add只接受整数天不能直接加月加月用add_months。如果created_at是字符串先to_timestamp转否则这些函数返回 null。Spark SQL 里日期函数对 null 输入一律返回 null不会报错所以清洗阶段补空很重要。5. 避坑与排查那些让项目跑不起来的细节5.1 连接超时与认证失败现象Spark 任务报MongoTimeoutException: Timed out after 30000 ms while waiting for a server。 原因MongoDB 的bindIp只绑了127.0.0.1Spark executor 在别的节点连不上或者防火墙没放 27017。 解决改/etc/mongod.conf的net.bindIp为0.0.0.0或具体内网 IP重启mongod再ufw allow 27017或对应防火墙规则。生产环境别开0.0.0.0绑内网网段。现象报Authentication failed。 原因连接串没带authSource或用户没建在对应库。 解决确认用户建在哪个库连接串加?authSourceadmin或对应库名。用mongosh手动连一次验证凭据排除 Spark 侧问题。5.2 分区不均导致任务长尾现象Spark UI 里大部分 task 秒完少数几个跑十几分钟。 原因_id是 ObjectId前 4 字节是时间戳递增插入的数据按_id范围切分后新数据全挤在最后一个分区。 解决换SAMPLE分区策略或按业务字段如user_id做partitionKey自定义分区。也可以在写入 MongoDB 时打散_id但改动大优先调读分区。5.3 写回时_id冲突现象写回报DuplicateKeyException。 原因目标集合已有相同_id的文档append模式不做 upsert直接插入冲突。 解决改用overwrite清空重写或配置update模式加idFieldList做 upsert。如果业务上_id由 Spark 生成确保生成逻辑唯一别用monotonically_increasing_id()当业务主键它只在分区内唯一。5.4 内存溢出与 GC 频繁现象executor 日志刷GC overhead limit exceeded或直接OutOfMemoryError。 原因partition.size太大单分区数据超 executor 内存或cache()了超大 DataFrame。 解决调小partition.size到 32 或 16增大spark.executor.memory去掉不必要的cache()。如果用了collect()把大结果拉回 driver改成write落盘或show(n)限制条数。5.5 版本不匹配的隐蔽报错现象NoSuchMethodError: com.mongodb.client.MongoCollection.find(...)。 原因Connector 依赖的 MongoDB Java Driver 版本和 Spark 环境里已有的 Driver 冲突常见于 Spark 自带 jars 里有旧版 Driver。 解决用--jars显式指定 Connector 及其依赖或用 shade 打包把 Driver 重定位。检查spark-submit的 classpath 顺序确保 Connector 的 jar 在前。6. 进阶把项目从能跑变成可验证、可复现跑通一次不代表项目稳。我一般会加三层验证数据量对账、抽样比对、幂等重跑。数据量对账是读 MongoDB 的count和 Spark 的count比差太多说明分区读漏了或过滤条件写错。抽样比对是随机抽几条_id在mongosh和 Spark 里分别查看字段值是否一致尤其注意日期和数值类型MongoDB 的ISODate和 Spark 的timestamp时区处理不同差 8 小时是常见坑。幂等重跑是把同一批任务跑两次看结果集合条数是否翻倍翻倍说明没做去重或 upsert。再往上走可以把清洗和聚合逻辑参数化用配置文件控制输入集合、输出集合、分区大小、过滤条件这样换一份数据不用改代码。资料包里的优秀项目通常有application.conf或config.yaml照着改就行。如果要做可视化把daily_stats集合接到 ECharts 或 Superset按day做时间轴status做系列就是一个能演示的校园大数据或电商分析大屏。我自己的习惯是每接一个新的大数据项目先花半天把 MongoDB 和 Spark 的版本矩阵确认清楚再用最小数据集跑通读写闭环最后才灌全量数据。这个顺序能省掉后面大量的排查时间。希望帮到你。本文还有配套的精品资源点击获取