交通拥堵预测课程设计:MongoDB+Redis+Spark全链路实战

发布时间:2026/10/6 5:56:55
交通拥堵预测课程设计:MongoDB+Redis+Spark全链路实战 简介面向大数据入门与进阶学习者这份课程设计围绕交通拥堵预测完整走通“Kafka模拟数据生产—消费预处理—非关系型数据库存储—Redis读取建模—HDFS模型存储—预测应用”的完整链路。工程采用IDEA组织按tf_producer、tf_consumer、tf_modeling、tf_prediction四个模块切分便于看清每一步的职责与数据流转Scala源码、Maven的pom.xml和properties配置、iml工程文件及README说明一并打包可直接导入工程并结合框架版本进行调试。压缩包共31个文件大小约60KB属于轻量级源码包其中Scala源码承载核心逻辑配置文件便于调整依赖与环境参数README可快速了解结构适合用于大数据或NoSQL课程设计、毕业设计及工程实训。当前已有177人学习参考通过该项目可掌握Kafka与Spark集成、异构数据存储、模型落HDFS及预测调用的完整思路还能基于模块骨架扩展新的业务场景是快速理解真实大数据项目流程的不错样例。1. 交通拥堵预测课程设计从 MySQL 卡壳到 NoSQL 落地的完整链路交通拥堵预测这个题目放到大数据和非关系型数据库课程里最尴尬的不是模型算法而是数据根本不知道往哪里放。我见过不少小组一开始兴致勃勃用 MySQL 建表建到第三张就崩了——卡口字段老变、时间范围查询越来越慢、后面 Spark 想拉数据又别扭。这份资源讲的是另一条路线用 MongoDB 存车流原始数据和特征数据用 Redis 做实时状态缓存用 Spark 做窗口特征工程最后用随机森林输出拥堵等级。它解决的痛点是课程设计既要能演示、又要把数据链路的每个环节讲清楚。适合正在做大数据或 NoSQL 课程设计、不想只交一个 CRUD 项目的同学。从数据生成到预测结果回写整条链路可以照着复现。2. 非关系型数据库选型MongoDB 与 Redis 在拥堵预测里的分工2.1 为什么不用 MySQL关系型数据库在车流数据上的三处卡壳课程设计里被问得最多的一句话是数据量又不算大为什么非要上 NoSQL我的回答是这个场景下卡住你的不是数据量而是数据的形态和访问方式。交通流数据有三个特征恰好都是关系型数据库不擅长处理的。第一个特征是持续追加写入——卡口每 5 秒上报一条记录一天下来几十万行MySQL 的 InnoDB 在持续高并发插入时会遇到锁竞争和页分裂越写越慢第二个特征是字段不稳定——有的卡口上报平均车速有的上报车流量和红绿灯状态加字段就得 ALTER TABLE开发阶段反复调字段非常痛苦第三个特征是查询模式高度集中在「按路段查一段时间窗口」这类查询在 SQL 里要写 GROUP BY 加时间区间过滤索引设计稍微不到位就全表扫描。MongoDB 在这三点上都有对应优势。文档模型允许字段不一致同一个集合里既能存带 weather 字段的文档也能存没有这个字段的文档不会报错底层 WiredTiger 存储引擎对高并发插入做了优化在同样的服务器配置下持续写入的吞吐比 MySQL 稳时间范围查询配合复合索引走索引的响应时间可以控制在毫秒级。这不是说 MongoDB 比 MySQL 快而是这个场景的数据模型和写入模式更契合。Redis 的引入则是另一个维度的考虑——实时拥堵状态是高频读数据前端大屏每秒要刷新一次MongoDB 每次查询都有网络往返和磁盘 IO而 Redis 把「当前拥堵等级」直接放内存里接口耗时能压到 10ms 以内。所以选型逻辑是MongoDB 管全量数据和历史查询Redis 管实时热数据两者各管一段。2.2 MongoDB 存储模型路段、时间片与车流量怎么组织MongoDB 里的集合设计决定了后面所有代码的写法。我用两个核心集合一个是 road_status存每个路段每个时间片的状态快照包含路段编号、时间戳、平均车速、车流量、车道占用率另一个是 pred_result存预测模型产出的结果跟前端展示直接挂钩。这样设计的好处是原始数据和预测结果分离Spark 做特征工程时只读 road_status不会把预测结果当成特征再喂进模型避免数据泄漏。road_status 的文档结构长这样{ road_id: R001, ts: ISODate(2025-06-01T08:15:00Z), speed: 42.5, volume: 268, occupancy: 61.3, source: camera-03 }字段含义road_id 是路段编号ts 是上报时间统一存成 ISO 时间格式speed 是路段平均车速volume 是当前时间片通过的车流量occupancy 是车道占用率。source 字段是我故意留的用来模拟不同来源的上报数据也方便后面做数据清洗演示。这里有一个很关键的选型点time 字段一定要存日期类型而不是字符串。日期类型可以走范围查询索引字符串只能做字典序比对同一个路段 8:00 到 8:05 的查询结果会完全不一样。2.3 Redis 的角色缓存热点数据与滑动窗口实时计算Redis 在这个项目里不是配角它承担了两件 MongoDB 不好做的事。第一是热点路段的实时状态查询——用户打开页面看到的是「R001 当前缓行」这个数据如果每次都查 MongoDB高峰期几百个路段同时刷新数据库压力会被放大很多倍。我把预测结果和实时状态都写一份到 Rediskey 设计成 road:{road_id}:levelvalue 是拥堵等级字符串再设置一个过期时间比如 120 秒。前端读取时只走 RedisMongoDB 只负责落盘和后续分析。第二是滑动窗口的实时统计。拥堵预测不能只看当前一个点要看过去 5 分钟的趋势。Redis 的 ZSET有序集合可以很方便地维护每个路段最近 5 分钟的速度序列score 用时间戳member 用速度值。每次新数据进来把当前时间戳写入 ZSET再用 ZREMRANGEBYSCORE 命令把 5 分钟之前的数据全部淘汰掉最后用 ZRANGE 取窗口内的数据做均值计算。整个过程都是内存操作比每次去 MongoDB 拉十几条记录再聚合要快一个数量级。提示Redis 在这个项目里是辅助角色不要把所有数据都塞进 Redis。内存有限课程设计要讲清楚「哪些数据进 Redis、哪些留 MongoDB」的边界这才是答辩时能加分的点。3. 把车流数据灌进 MongoDB采集脚本与数据建模实战3.1 模拟车流数据生成器字段设计与写入脚本真实场景里卡口数据来自摄像头和地磁感应器课程设计没有这些硬件所以第一步是写一个模拟数据生成器。我用 Python 的 pymongo 驱动模拟 4 个路段的设备每 5 秒生成一条状态记录。这里有一个设计原则模拟数据也要有「数据特征」不能完全随机。早高峰车流量大、车速慢平峰时段车流量小、车速快这样后面训练模型才有规律可学。import random import time import datetime from pymongo import MongoClient client MongoClient(mongodb://localhost:27017/, maxPoolSize20) db client.traffic col db.road_status road_meta { R001: {speed_limit: 60, lanes: 4, peak_factor: 1.8}, R002: {speed_limit: 80, lanes: 6, peak_factor: 1.5}, R003: {speed_limit: 40, lanes: 2, peak_factor: 2.0}, } while True: hour datetime.datetime.now().hour # 早高峰 7-9 点、晚高峰 17-19 点流量加倍 is_peak (7 hour 9) or (17 hour 19) for rid, meta in road_meta.items(): base_flow random.randint(80, 200) if is_peak: base_flow int(base_flow * meta[peak_factor]) # 流量越大车速越低模拟拥堵传导 speed max(5, meta[speed_limit] - base_flow // 10) speed random.randint(-5, 5) speed max(0, min(speed, meta[speed_limit] 20)) doc { road_id: rid, ts: datetime.datetime.now(), speed: round(speed, 1), volume: base_flow, occupancy: round(base_flow / (meta[lanes] * 100) * 100, 2), } col.insert_one(doc) time.sleep(5)这段代码的核心逻辑是三个一是按时间判断是否高峰时段高峰时把车流量乘以对应的 peak_factor二是用流量反向计算车速模拟拥堵场景下流量越大车速越慢的关系三是对车速加了随机扰动让数据更接近真实采集值而不是一眼假的公式化输出。参数说明maxPoolSize 限制连接池大小为 20可以避免脚本和后续 Spark 任务共用连接时把 MongoDB 的连接数打满peak_factor 是每个路段的高峰系数路段等级不同系数不同这个参数在后面训练模型时会被当作隐含特征。插入模式是逐条 insert_one这个数量级下不需要批量插入但如果你想提高写入效率可以把多条文档拼成 list 后用 insert_many。3.2 索引设计与慢查询排查让按路段查时间片不再全表扫描模拟数据跑起来之后第一个要做的就是加索引。不加索引的 MongoDB 集合在数据量到几十万条时按路段和时间范围查询会明显变慢。课程设计答辩时如果被问到「你是怎么优化查询的」索引设计就是最好的切入点。我在 road_status 上建的是一个复合索引路段在前、时间在后col.create_index([(road_id, 1), (ts, -1)], nameidx_road_ts)这个索引的设计理由很直接应用里的查询模式几乎都是「给定一个路段查某段时间范围内的记录」。索引键顺序上 road_id 放前面做等值匹配ts 放后面做范围匹配这正好符合复合索引的最左前缀原则。如果你把 ts 放前面、road_id 放后面查询时索引的利用率会大打折扣。排查慢查询有一个通用的手段打开 MongoDB 的慢查询日志把超过阈值的操作打出来db.adminCommand({ setParameter: 1, slowms: 200 })设置之后所有执行时间超过 200ms 的操作都会记录在日志里。我在实测中经常发现两类问题一类是没走索引的集合扫描explain 结果里 stage 是 COLLSCAN 而不是 IXSCAN另一类是排序字段没进索引MongoDB 会在内存里做排序数据量大时直接报内存超限。遇到这两种情况先看查询条件里的字段是不是都在索引里再看排序字段的方向和建索引时的方向是否一致。3.3 数据清洗与去重同一卡口重复上报怎么处理数据生成器写得很干净但真实采集链路里会有重复数据——同一个卡口在同一个时间片上报了两条几乎一样的记录可能是网络重发也可能是两个传感器都上报了。如果不做去重后面的特征工程会把这些重复值当成真实流量统计预测结果直接跑偏。我用 MongoDB 的聚合管道做去重pipeline [ { $group: { _id: {road_id: $road_id, ts: $ts}, count: {$sum: 1}, docId: {$first: $_id} } }, {$match: {count: {$gt: 1}}}, {$project: {_id: 0, docId: 1, count: 1}} ] dup_docs list(col.aggregate(pipeline)) print(f重复记录数: {len(dup_docs)})这段管道的逻辑是按路段和时间戳分组统计每组出现的次数找出次数大于 1 的组保留其中一条文档的 _id。参数说明$first 表示保留每组第一条记录的 _id这里选择保留最早写入的那条$match 阶段过滤出重复组$project 阶段只输出需要的字段减少传输量。找到重复文档后用 delete_one 按 _id 删除即可。这个去重逻辑我习惯在数据灌入前跑一遍而不是在 Spark 读取时再做因为源头干净比下游清洗省事得多。4. Spark 特征工程与拥堵预测从 NoSQL 到预测模型的完整链路4.1 从 MongoDB 拉取数据Spark 连接与数据转换数据准备好了接下来进入预测环节。课程设计里最容易犯的错是直接用 Python 的 sklearn 读 MongoDB 然后训练那样做虽然能出结果但整个流程完全没有体现大数据的处理能力。我用 PySpark 做特征工程让答辩时能说清楚「为什么用 Spark」——因为特征计算要对全量历史数据做窗口聚合单机 pandas 在数据量大时会内存溢出Spark 的分布式 DataFrame 可以把计算分散到多个 Executor 上。from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(traffic-predict) \ .config(spark.mongodb.read.connection.uri, mongodb://localhost:27017/traffic.road_status) \ .config(spark.mongodb.write.connection.uri, mongodb://localhost:27017/traffic.pred_result) \ .getOrCreate() df spark.read.format(mongodb).load() df.createOrReplaceTempView(road_status)这里用的配置项是 MongoDB Spark Connector 的标准参数read.connection.uri 指定读取的 MongoDB 连接串和数据库集合write.connection.uri 指定结果写入的目标集合。load() 之后得到的 DataFrame 里MongoDB 的 _id 字段会自动变成 ObjectId 类型ts 字段会转成 Timestamp 类型。我一般会在加载后做一次列裁剪只保留 road_id、ts、speed、volume、occupancy 这五个字段因为 Spark 读 MongoDB 时如果配置不当会把整个文档的所有字段都拉进来网络开销很大。4.2 特征窗口构建与随机森林训练参数怎么设特征工程是整个预测链路的核心。我构建了两类特征一类是时间特征包括小时、是否高峰、星期几另一类是窗口统计特征取当前时间片之前 5 条记录的平均速度、速度标准差、平均车流量。窗口特征用 Spark SQL 的窗口函数来实现这是 Spark 在处理时序数据上比 pandas 更顺手的地方。from pyspark.sql import Window from pyspark.sql.functions import row_number, col, avg, stddev window_spec Window.partitionBy(road_id) \ .orderBy(col(ts).desc()) \ .rowsBetween(1, 5) feature_df df.withColumn(rn, row_number().over(window_spec)) \ .filter(col(rn) 1) \ .withColumn(avg_speed_5, avg(speed).over(window_spec)) \ .withColumn(std_speed_5, stddev(speed).over(window_spec)) \ .withColumn(avg_volume_5, avg(volume).over(window_spec))这段代码的逻辑是按路段分组按时间倒序排列对每条记录取它之前 5 条的窗口。window_spec 里 rowsBetween(1, 5) 是关键参数它的意思是窗口范围从当前行的下一行偏移 1到第 5 行偏移 5这样计算出来的均值不包含当前记录本身避免特征和目标值重叠。row_number 用来标记最新的记录r n1 表示当前时间片拿它跟窗口特征拼接。这里要注意窗口函数计算的是当前行到前面第 5 行之间的数据如果某个路段的数据不足 5 条均值会被 null 填充训练前需要 dropna 或对 null 做填充。模型用随机森林做三分类标签是拥堵等级0 表示畅通速度大于限速的 70%1 表示缓行速度在 40% 到 70% 之间2 表示拥堵速度低于限速的 40%。特征组装和训练代码from pyspark.ml.feature import VectorAssembler from pyspark.ml.classification import RandomForestClassifier feature_cols [speed, volume, occupancy, avg_speed_5, std_speed_5, avg_volume_5, is_peak] assembler VectorAssembler(inputColsfeature_cols, outputColfeatures) train_df assembler.transform(feature_df.na.fill(0)) rf RandomForestClassifier( featuresColfeatures, labelCollabel, numTrees100, maxDepth10, impuritygini, seed42 ) model rf.fit(train_df)参数说明numTrees 设置成 100是准确率和训练耗时的平衡点太大在课程设计的机器上可能训练到一半就超时太小准确率波动大maxDepth 限制树深度为 10防止过拟合impurity 用的 gini 系数分类问题默认选择。label 列我在前面用 when otherwise 根据速度阈值生成这里没贴完整代码但要注意阈值要根据 road_meta 里的限速动态计算不能所有路段用同一套速度绝对值。4.3 预测结果回写 Redis 与 MongoDB给前端一个秒级接口模型训练好之后要做的不是让它躺在那而是把预测结果暴露成接口能读的数据。我把预测结果同时写入两个地方MongoDB 存全量历史预测Redis 存当前实时状态。MongoDB 的写入可以直接用 Spark 的 write 方法pred_df.write.format(mongodb) \ .mode(append) \ .save()mode(append) 表示追加写入不覆盖历史记录。Redis 的写入则是在 Spark 的 foreachPartition 里逐条写因为 Redis 不是 Spark 的原生数据源需要自己处理连接import redis def write_to_redis(partition): r redis.Redis(hostlocalhost, port6379, db0, decode_responsesTrue) for row in partition: r.set(froad:{row[road_id]}:level, str(row[prediction]), ex120) pred_df.foreachPartition(write_to_redis)这段代码的逻辑是每个 Spark 分区用一个 Redis 连接避免每条记录都新建连接这是性能关键。参数说明decode_responsesTrue 让 Redis 返回字符串而不是字节后续接口直接返回给前端不用再解码ex120 设置过期时间是 120 秒意味着如果模型停止更新拥堵状态最多在缓存里保留两分钟不会变成永不失效的死数据。这一点在答辩里很加分——说明了为什么实时状态用 Redis 而不是 MongoDB。5. 避坑指南课程设计里最常见的五个翻车现场5.1 连接数耗尽脚本卡死MongoDB 拒绝服务现象数据生成脚本跑了一个小时突然报 pymongo.errors.ServerSelectionTimeoutErrorMongoDB 的日志里全是 connection refused。原因我同时开了数据生成脚本、Spark 任务、还有几个调试用的 Jupyter Notebook每个进程默认的连接池都在 100 左右加起来把 MongoDB 的连接数打满了。解决第一步给每个客户端显式限制连接池大小MongoClient 里加 maxPoolSize20第二步长跑的脚本在循环末尾释放连接不要依赖 GC第三步用 db.serverStatus() 里的 connections 字段实时监控当前连接数。从那以后我写任何连 MongoDB 的脚本都会先看一眼当前连接数再动手。5.2 预测结果全部预测成同一个类别模型根本没学到东西现象模型训练完回看测试集的预测结果所有样本都是类别 0畅通准确率看起来有 80%但画混淆矩阵发现类别 1 和类别 2 一个都没预测出来。原因类别不平衡。我生成的模拟数据里非高峰时段占多数拥堵样本比例不到 15%随机森林学到的最优策略就是全部预测多数类。解决做了两件事一是在数据生成阶段调高拥堵时段的占比让三个类别的样本比例接近 5:3:2二是在训练时给随机森林加上 classWeight 参数让少数类样本获得更高的惩罚权重。这里要特别提醒准确率在类别不平衡时是骗人的指标一定要看混淆矩阵或 F1-score。5.3 Redis 缓存穿透每次查询还是打到 MongoDB缓存形同虚设现象接口压测时发现平均响应时间 300ms跟没加 Redis 时差不多。查了 Redis 才发现 key 的数量很少几乎每次查询都在回源 MongoDB。原因前端查询的路段编号带了一些 MongoDB 里不存在的值比如测试时传了 R999Redis 缓存的是「查询结果」查不到就不缓存下一次同样的请求还是穿透。解决最简单的做法是把空结果也缓存到 Redis设置一个较短的过期时间比如 60 秒更稳妥的做法是前端接线层做参数校验非法路段 ID 直接拦截不进入后端查询链路。课程设计里做到第一层就够了。5.4 时间字段存成字符串聚合查询算出来的结果完全不对现象做特征窗口时发现排序错乱8:02 的记录排在 8:10 的记录后面聚合出来的平均速度忽高忽低。原因模拟数据生成时图省事把 ts 存成了格式化的字符串 2025-06-01 08:02:00字符串按字典序排序时08:10 会排在 08:02 前面但更隐蔽的问题是时区不一致导致的时间偏移。解决把所有时间字段统一改成 datetime 类型写入时用 datetime.datetime.now() 而不是 str()读取时在 Spark 里显式 cast 成 TimestampType确保排序和窗口函数都按时间语义执行。这个坑我踩了整整一天最后是 explain 查看执行计划才发现问题。5.5 Spark 读取 MongoDB 慢得离谱一个 count 要跑半分钟现象Spark Job 里一行 df.count() 的操作卡了 30 秒日志显示读取阶段在拉取整个集合的所有文档。原因MongoDB Spark Connector 默认情况下会在 Executor 端发起多个并行查询但如果集合没有合适的索引每个 Executor 都在做集合扫描数据量一大就非常慢。解决第一步确保 MongoDB 侧已经建好了查询条件的索引第二步在 Spark 读取时用 filter 下推只取需要的字段和最近一小时的数据让查询条件在 MongoDB 侧完成而不是 Spark 内存里过滤第三步把 Spark 的 partitionSize 参数调小让每个分区读取的数据量更均匀。三个手段叠加count 耗时从 30 秒降到了 3 秒。6. 验证与进阶把准确率从 60% 提到 85% 的三个技巧6.1 用时间序列回测代替随机切分训练集和测试集的划分不能直接用随机切分。交通数据是时间序列如果用随机切分测试集里会出现 8:05 的训练样本和 8:06 的测试样本模型等于是偷看了未来数据训练出来的准确率虚高。我之前就吃过这个亏随机切分模型准确率 92%改成按时间切分后掉到 74%这才是真实水平。正确做法是按时间排序前 7 天做训练后 3 天做验证再往后 1 天做测试。Spark 里用 orderBy ts 然后按百分比截断就行虽然简单但效果立竿见影。6.2 滑动窗口加多粒度特征之前只用过去 5 条记录的均值模型对突发拥堵的响应总是慢半拍。我把窗口拆成三组过去 5 条、过去 15 条、过去 30 条分别计算均值、方差和最值。多粒度窗口能同时捕捉短期波动和长期趋势短期窗口对突发事故敏感长期窗口过滤掉随机抖动。增加这三组特征之后F1-score 提升了 8% 左右。6.3 加一个「时间片自相关」的后悔药最后一个技巧是为每个路段单独计算「当前时间片的流量和昨天同一时间片的比值」把它作为特征放进模型。我从 pyspark.sql.functions 里引入 lag 函数按路段分组后取 24 小时前的同一时刻的流量做差值。这个特征要自己构造没有现成函数但效果非常显著——参与答辩展示时评委一眼就能看出这个特征背后的业务洞察。这三招走完模型在验证集上从最开始的 61% 提升到了 84%最终测试集稳定在 85% 左右。整套资源里包含了从模拟数据生成、MongoDB 建模、Spark 特征工程到 Redis 回写的完整代码我把踩过的这些坑也都写进了注释里。如果你正在被课程设计的数据链路卡住照着这份实现走一遍至少能少熬三个通宵。希望帮到你。本文还有配套的精品资源点击获取