多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Hadoop+Spark景区客流量预测与景点推荐系统:从毕设到工程实战全复盘

Hadoop+Spark景区客流量预测与景点推荐系统:从毕设到工程实战全复盘 每年到了毕业季计算机大类的同学都在和毕设题目搏斗。我收到最多的私信就是想做大数据方向但不知道从哪里下手。如果你打开过课题库大概率见过智慧旅游大数据景区客流量预测景点推荐系统这类名字——几个词叠在一起看起来热闹实际做起来要踩的坑远比想象中多。这篇博客就把HadoopSpark景区客流量预测与景点推荐系统这个题目从选题拆解、数据采集、数仓建设、算法建模到部署调优完整复盘一遍。不管你是正在为毕设发愁的学生还是想了解大数据项目完整链路怎么搭的开发者这篇都能当成一份带避坑经验的项目笔记来看。整套项目最终交付源码、文档、PPT和讲解但我更想分享的是那些不会写进文档里的真实工程判断。1. 把智慧旅游大数据拆开看五个模块的职责与数据流转1.1 一个毕设题目拆出来的五个交付物很多同学看到这个题目第一反应是要写一个App。这是最大的误解。这类题目本质不在前端界面而在数据链路。把它拆开其实是五个独立但串联的模块旅游爬虫模块采集景区评论、门票信息、天气历史、节假日数据解决的是数据从哪来。离线数仓模块基于Hadoop生态HDFSYARNHive对原始数据进行清洗、去重、聚合解决的是数据怎么变干净、怎么组织。客流量预测模块基于Spark MLlib构建回归模型预测景区未来N天的客流解决的是数据能算什么。景点推荐模块基于Spark的ALS协同过滤算法给用户生成个性化景点TopN推荐解决的是数据怎么服务业务。可视化展示模块把预测结果、推荐结果、统计数据做成一整套展示页面解决的是结果怎么讲清楚。这五个模块合在一起恰恰构成一条完整的采集-存储-计算-应用流水线。答辩时评委最看重的不是某一个算法多高级而是你能不能把这条链路讲圆、讲透。1.2 数据从爬虫到推荐结果的完整流转路径我在做项目梳理时习惯先画一条数据流把所有模块串起来第一步数据采集。爬虫把景区评论、天气、票务等数据抓下来以CSV或JSON格式分片落盘到本地临时目录。注意这里不是直接写MySQL而是先落到文件再批量上传HDFS原因是爬虫的网络状态不稳定单个小文件直接写分布式文件系统容易造成大量零碎小文件后面跑Hive会很痛苦。第二步HDFS存储原始数据。原始数据进入HDFS的ODS层原始数据层按数据源分目录存放。这是数仓的地基后续所有计算都从这里开始。第三步Hive离线清洗。通过Hive SQL做ETL去除重复评论、补齐缺失天气、过滤异常值把数据从ODS层清洗到DWD层明细层再聚合出按景区日期为粒度的指标形成DWS层汇总层。第四步Spark建模。Spark读取DWS层的特征宽表一份数据喂给随机森林做客流预测另一份喂给ALS算法做协同过滤推荐训练出的模型和输出结果写回MySQL或者HBase供查询使用。第五步Web与可视化。后端接口读取MySQL里的预测结果和推荐结果前端大屏展示。这一步是整个项目最直观的脸面也往往是答辩现场评委停留最久的地方。这条链路的逻辑是下一环只消费上一环的输出每一环的输入输出都清晰每一环也都对应一个可演示的交付物。2. 技术选型不是越新越好HadoopSpark组合的真实理由2.1 为什么不是纯Python单机也不是Flink我最初也考虑过是不是直接用PythonSklearn就能把所有事情做完毕竟数据量就几万到几十万条单机完全能跑。但最后答案是否定的原因不是性能而是毕业设计的答辩逻辑。先看一组对比方案学习曲线单机数据量下的表现答辩加分点踩坑成本PythonSklearn全栈平缓上手快完全够用几乎没有体现不了分布式低但分数天花板低Flink全链路陡峭窗口机制难啃杀鸡用牛刀加分但招架不住深挖追问高中文资料少HadoopSpark中等生态资料极多分布式下照样跑高能讲出完整生态故事中等大坑都有解决方案纯Python方案的问题是当评委问出你的数据量这么大为什么不用分布式的时候很难给出有底气的回答。Flink方案的问题则相反它是真正的流处理引擎但毕设选题通常是离线分析场景硬上Flink会带来不必要的复杂度比如状态管理、精确一次语义、窗口在水位线下的行为这些问题即便是工作了两三年的人也不一定都完全吃透。HadoopSpark是熟练度性价比最高的组合Hadoop负责存储和资源调度讲分布式故事Spark负责计算性能和API友好度都远胜MapReduce两者加起来既能讲出生态又不会让实现复杂到失控。2.2 Hadoop生态里每一层各干什么很多同学的认知停留在Hadoop就是一个东西实际上这里至少有四层分工HDFS存储层把大文件切块分布到多台机器上默认块大小128MB。它解决的是数据太大单机存不下的问题。YARN资源调度层管理整个集群的内存和CPU统一分配资源给计算任务。MapReduce、Spark、Flink都跑在它上面。Hive数据仓库工具本质是把SQL翻译成MapReduce或Spark任务让不想写Java/Scala的人能用SQL做分析。Spark计算引擎基于内存迭代计算比MapReduce快不少适合机器学习算法的反复迭代。这套组合的好处在于Hive做数仓管理比SparkSQL更方便Spark做算法比纯MapReduce更现实物尽其用。3. 旅游数据从哪来爬虫模块的采集目标与反爬工程3.1 先想清楚要采什么数据字典先行动写爬虫之前一定先列数据字典明确采什么字段、解决什么问题。否则会陷入无限抓数据的泥潭。我整理了一份最简采集字典数据源类型核心字段在项目里干什么用景区评论评论内容、评分、用户ID、评论日期、景区名称构建ALS评分矩阵做情感打分天气历史日期、城市、最高温、最低温、天气类型客流量预测的关键特征票务信息景区名称、门票价格、当月销量流量热度的辅助特征节假日历日期、节假日名称、是否调休客流预测的强特征注意评论是推荐系统的原料天气和节假日是预测模型的原料两组数据服务不同的算法在采集阶段就要按业务用途分目录存放不要混在一起。3.2 三种采集策略按页面复杂程度选择真正写爬虫时三种策略基本覆盖了大多数网站第一种requestsXPath处理静态HTML页。这是最轻量的方案直接请求URL拿HTML再用lxml的XPath提取目标字段。适合景区官网的票务列表、静态评论页。核心代码思路import requests from lxml import etree headers { User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 } resp requests.get(url, headersheaders, timeout10) html etree.HTML(resp.text) # XPath定位评论内容 comments html.xpath(//div[classcomment-item]//span[classcontent]/text())第二种直接解析接口JSON。很多信息其实被前端用JS渲染了但背后一定有接口。打开浏览器的开发者工具切到Network按XHR过滤直接分析Ajax返回的JSON比Selenium高效得多。比如某些网站的评论接口返回的是标准JSON把里面的comment、score、date字段取出即可。第三种Selenium兜底动态渲染页面。如果接口加了复杂加密参数实在解析不动就用Selenium起一个真实浏览器来做兜底。但尽量少用因为它慢且占资源爬几万条数据会等到怀疑人生。3.3 反爬规避与数据落地的工程细节反爬是爬虫模块最容易翻车的地方。我分享几个实测有效的做法控制请求节奏。不要用固定间隔去访问而是用随机间隔比如time.sleep(random.uniform(2, 5))。固定间隔是反爬系统最容易识别的特征之一。User-Agent轮换。准备一个UA池每次请求随机取一个模拟不同的浏览器和操作系统。简单场景下这招就够了比上代理更省心。代理IP池。如果目标网站对访问量敏感单个IP短时间请求过多会被限流。可以准备一个代理池按请求失败率自动切换代理。这一块水很深但毕设场景下只做基本切换即可用量不大。断点续爬。这是一个被大多数人忽略、但能救命的设计。爬虫跑了几万条突然被限流中断如果没有断点续爬之前的进度全部作废。做法很简单定时把已爬取数据的最后一条索引写入本地文件重启后从断点继续。数据落地分片写文件而不是直连数据库。我建议爬虫先把数据写入本地JSON/CSV片文件每1万条一个文件统一命名最后用脚本一次性hdfs dfs -put上传。直连数据库有两个问题数据库连接频繁、大事务耗时而且如果爬虫中途挂了数据库里会留下一半脏数据。4. 清洗与建模的中间层用Hive把原始数据喂成特征宽表4.1 毕设够用的三层数仓结构很多同学不理解为什么要分ODS、DWD、DWS觉得多此一举。真正常见的情况是爬虫数据直接拿来做算法结果发现数据里有大量重复和缺失模型效果一言难尽。数仓分层的本质是让每一层只干一件事把脏活累活隔离在ODS到DWD这一步。ODS层原始数据层原封不动地放爬虫上传的文件HDFS上的原始目录。DWD层明细数据层清洗后的明细数据字段规范、去重完成一行就是一条有效记录。DWS层汇总服务层按景区日期做聚合得到建模所需的统计指标日均评论数、美誉度、近7天流量等。4.2 数据清洗的四条硬规则我建议清洗规则只定四条覆盖90%的问题去重。用ROW_NUMBER()按用户ID评论内容评论日期分区取第一行。核心SQLINSERT OVERWRITE TABLE dwd_comment SELECT uid, scenic_id, content, score, comment_date FROM ( SELECT uid, scenic_id, content, score, comment_date, ROW_NUMBER() OVER (PARTITION BY uid, content, comment_date ORDER BY comment_date DESC) AS rn FROM ods_comment ) t WHERE rn 1;缺失值处理。评论缺失内容直接丢弃天气字段缺失则用该城市前后两天的均值补齐或用日期关联到历史天气维表直接join填入。异常值过滤。客流不可能为负评分不可能超过5。这类明显异常值直接过滤掉不要硬留着喂给模型。字段规范化。把时间字段统一成yyyy-MM-dd格式景区名称去除首尾空格评分转为浮点数。规范化的好处在join特征宽表时才能真正体会到。4.3 特征宽表里到底放了哪些字段预测模型不直接吃原始表吃的是一张特征宽表每行是一个景区日期的样本每列是一个特征。我设计的宽表字段如下特征类别字段名说明时间特征week_day一周中的第几天0-6时间特征is_weekend是否周末时间特征is_holiday是否节假日天气特征max_temp / min_temp当天最高/最低温天气特征weather_type晴、多云、雨、雪等历史客流last_1_day_traffic前1天实际客流历史客流last_7_day_avg前7天平均客流历史客流last_year_same_day去年同期客流热度特征comment_cnt_7d近7天评论数标签traffic当日实际客流预测目标这个宽表做完之后等于把预测问题简化成了给定这些特征回归出一个客流数值后面的模型工作都建立在这张表之上。5. 客流量预测从时间序列基线上手再到随机森林5.1 为什么先要一个能跑的基线我第一次做预测时直接端着随机森林就跑结果模型表现倒是能看但答辩被老师问了一句你对比过基线模型吗直接愣住了。后来才明白任何预测项目都要先有一个简单粗暴的基线用它衡量复杂模型的增量价值。基线模型我建议用时间序列里的Holt-Winters指数平滑法。它是一种对趋势和季节性做平滑的方法思路直观用Spark没实现没怎么做可以直接用pandas在本地跑from statsmodels.tsa.holtwinters import ExponentialSmoothing model ExponentialSmoothing( train_series, trendadd, seasonaladd, seasonal_periods7 # 一周为周期 ).fit()基线的意义在于它是无脑规律派的天花板。如果随机森林连基线都打不过说明特征工程出了问题而不是算法不够好。5.2 Spark MLlib随机森林回归的完整流程基线跑完之后正式模型可以用Spark MLlib实现。我最终选择随机森林回归RandomForestRegressor核心代码from pyspark.sql import SparkSession from pyspark.ml.feature import VectorAssembler, StringIndexer from pyspark.ml.regression import RandomForestRegressor from pyspark.ml.evaluation import RegressionEvaluator spark SparkSession.builder \ .appName(TrafficPrediction) \ .config(spark.executor.memory, 2g) \ .getOrCreate() df spark.sql(SELECT * FROM dws_traffic_features) # weather_type 是字符串先转成数值索引 indexer StringIndexer(inputColweather_type, outputColweather_idx) df indexer.fit(df).transform(df) feature_cols [week_day, is_weekend, is_holiday, max_temp, min_temp, weather_idx, last_1_day_traffic, last_7_day_avg, last_year_same_day, comment_cnt_7d] assembler VectorAssembler(inputColsfeature_cols, outputColfeatures) data assembler.transform(df).select(features, traffic) train, test data.randomSplit([0.8, 0.2], seed42) rf RandomForestRegressor( featuresColfeatures, labelColtraffic, numTrees50, maxDepth10, seed42 ) model rf.fit(train) predictions model.transform(test) evaluator_rmse RegressionEvaluator(labelColtraffic, predictionColprediction, metricNamermse) evaluator_mae RegressionEvaluator(labelColtraffic, predictionColprediction, metricNamemae) print(RMSE:, evaluator_rmse.evaluate(predictions)) print(MAE:, evaluator_mae.evaluate(predictions))随机森林在毕设场景下有几个不可替代的优势不需要做复杂的特征放缩抗过拟合能力强而且numTrees和maxDepth两个参数调起来有规律答辩时可以说清楚每个参数的影响。5.3 评估结果与节假日失效的识别我的项目跑下来平日的RMSE大概在客流均值的10%以内看着还行。但把测试集按日期拆开看节假日样本的预测偏差普遍到了30%以上。这个问题定位下来有两个原因一是训练集里节假日样本占比太低模型没见过足够多的假期模式二是刚开始把is_holiday只当做一个普通数值特征模型没有给到应有的权重。解决办法分两步第一对节假日样本做加权让它们在损失函数里占更高比重第二把邻近节假日天数距离上次节假日天数这类衍生特征加进去让模型能感知到节假日前后的客流连续变化。调整之后节假日预测偏差降到了20%以内。这个评估时要按数据切片看误差的习惯也是做预测项目最值得养成的。6. 景点推荐ALS协同过滤的矩阵构建与冷启动兜底6.1 用户-景区评分矩阵怎么来推荐系统的原料是用户对景区的评分但我们拿到的原始数据通常没有明确的评分只有评论、点击、收藏之类的隐式反馈。这里我用了一个简单可解释的策略显式评分如果用户给景区打了分那直接用这个分数。隐式评分根据评论的情感倾向打分。具体做法是维护一个正向词表好玩推荐震撼和一个负向词表失望坑无聊统计评论命中正向词和负向词的个数按比例映射到1-5分。行为加权点击记0.5分收藏记1分成为会员记1.5分累加到基础评分上但最终封顶5分。最终生成一张(userId, scenicId, score)三列的表这就是ALS的输入。注意评分矩阵会非常稀疏——大多数用户只去过一两个景区一张十万行的评论表生成的矩阵可能只有几千条有效三元组后面会遇到这个问题。6.2 ALS调参与TopN推荐的工程实现Spark MLlib自带了ALS交替最小二乘实现核心参数就四个rank隐含特征维度我用了15。太小拟合不够太大容易过拟合一般取10-20。iterations交替迭代次数默认10就够调太大了收益不明显。lambda正则化系数0.1起步。数据稀疏时调高到0.5防止过拟合。alpha如果用的是隐式反馈这个参数控制置信度权重我设为1.0。推荐代码from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator als ALS( maxIter10, regParam0.1, userColuserId, itemColscenicId, ratingColscore, coldStartStrategydrop ) recommender als.fit(ratings_train) # 为所有用户推荐Top5景点 recs recommender.recommendForAllUsers(5)coldStartStrategydrop是必须要写的——没有它遇到稀疏矩阵里没有出现过的用户ALS会直接返回NaN。这在排错时很坑细节会在后面展开。6.3 冷启动用户别慌规则推荐兜底ALS再训练也解决不了新注册用户没有任何行为数据的冷启动问题。这里我做了一个兜底策略如果用户历史行为数少于3条不走ALS直接走规则推荐。规则推荐只做两件事一是按景区评论数排序输出近期热门TopN二是按景区类型城市的关联规则做组合推荐比如用户搜了自然风光类景区就推同城市同类型的其他景区。兜底策略并不智能但它保证了系统在任何输入下都有合理输出。答辩时把冷启动问题主动讲出来反而比藏着掖着更显真功夫。7. 部署与资源规划单机也能模拟集群但内存是真的坑7.1 三种部署环境按预算选毕设的硬件条件千差万别我见过的三种典型情况部署方式所需资源优点缺点Hadoop伪分布式Spark local一台普通电脑部署简单入门最快不能体现真正的分布式调度虚拟机三节点集群一台16G内存电脑更贴近真实集群能演示NameNode挂掉切换电脑内存压力大部署周期长云服务器多节点2-4台服务器演示效果好不受本机性能限制花钱按需购买容易超预算我推荐大部分学生选择虚拟机三节点因为答辩演示时你能理直气壮地说我配了三个节点NameNode挂掉后Standby节点自动接管这是一个非常有说服力的现场演示点。但如果电脑只有8G内存伪分布式也能跑完整个项目只是注意Spark的local模式默认会用满本机内存。7.2 最容易内存崩溃的配置文件部署中最容易出问题的不是Hadoop本身而是YARN的内存参数和Spark的executor参数没匹配上。我自己的经验是如果机器总内存16G留给Hadoop和Spark的可用内存不要超过12G要给操作系统留足够余量。关键配置项配置文件配置项建议值yarn-site.xmlyarn.nodemanager.resource.memory-mb8192总内存一半yarn-site.xmlyarn.scheduler.maximum-allocation-mb8192spark-defaults.confspark.executor.memory2gspark-defaults.confspark.driver.memory1gspark-defaults.confspark.executor.cores27.3 Spark任务提交与内存计算心法提交任务的命令长这样spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 2g \ --num-executors 3 \ --executor-cores 2 \ traffic_prediction.py很多人不理解为什么executor内存设了2g进程实际占用却不止2g。这是因为每个executor还有额外的off-heap开销Spark默认会额外申请约max(384MB, 0.1*executorMemory)的内存。也就是说2g的executor实际占用约2.2g3个executor加driver集群总资源至少要8g才不挤。资源分配的金标准是所有executor内存 overhead driver内存 NameNode/DataNode内存加起来不能超过yarn.nodemanager.resource.memory-mb否则任务会一直卡在等待资源的状态这是很多初学者排查一整天才发现的坑。8. 四个卡住进度的问题与完整排查思路8.1 Hive聚合任务在shuffle阶段直接OOM现象执行GROUP BY聚合的Hive SQL时报Error: Java heap space任务反复失败。排查链路先看YARN日志确认是Map端OOM还是Reduce端OOM。定位后发现Map端输出数据量巨大shuffle传输时Reducer拉取数据内存爆了。进一步查数据分布发现某个热门景区的评论量远超其他景区数据严重倾斜。根因数据倾斜。单个key的值太多shuffle时全部压在一个Reducer上。修复设置SET hive.groupby.skewindatatrue;让Hive做两阶段聚合先局部聚合再全局聚合同时对大key加随机前缀把压力打散。这个参数对group by场景几乎是特效药。8.2 Spark任务提交后一直卡在Accepted状态现象spark-submit提交任务后YARN界面显示任务一直处于ACCEPTED既不运行也不报错。排查链路打开ResourceManager页面发现集群总可用内存显示正常但运行中的任务列表为空。随后查看NameNode日志发现没有调度到任何Container进一步查scheduler.maximum-allocation-mb发现只有1G而Spark executor申请了2G——任务请求的资源超过了YARN的分配上限。根因YARN调度器单容器最大内存限制小于Spark申请的executor内存。修复把yarn.scheduler.maximum-allocation-mb调到8192。这个坑的隐蔽之处在于不是报错而是无限等待很多教程不会主动提到YARN的分配上限。8.3 ALS损失函数出现NaN现象ALS训练日志里迭代过程的损失输出从正常数值逐渐变成NaN最终模型完全没有可用结果。排查链路先怀疑学习率过大但ALS并没有显式学习率再怀疑评分矩阵有缺失值检查评分表没有空值。最后用SQL统计了每个用户的评分数量发现大量用户只有1条记录——评分矩阵稀疏到几乎无法形成有效的用户相似度。根因ALS在极度稀疏的矩阵上用户因子和物品因子的交替更新会出现数值不稳定正则化系数不足时loss容易发散。修复先过滤掉评分少于5条的用户再做regParam调参从0.1调到1.0同时启用coldStartStrategydrop保证预测阶段不产出NaN。最终loss稳定收敛。8.4 长假预测全面失效现象10月1日前后的客流预测偏差极大模型输出数值远低于实际值。排查链路把误差按日期切片绘制折线发现误差峰值集中在国庆节假日区间。翻看训练集发现一年之中节假日样本不足20天模型几乎没学到节假日模式。根因训练样本中节假日占比过低代价函数被非节假日样本主导。修复先做样本加权给节假日样本更高权重再做特征衍生增加距最近节假日天数是否长假前1天等特征让模型能感知临近节假日的客流爬坡趋势。修复后长假预测偏差从30%以上降到20%以内。这个按切片分析误差的方法在预测类项目里可以通吃。9. 毕业答辩与项目交付让评委确信你真的做懂了9.1 现场演示的节奏设计毕设演示的时间通常不会超过10分钟节奏一旦失控评委还没看到核心就切换PPT了。我建议按数据-预测-推荐三段式走第一段先打开大屏页面展示爬虫采集到的数据规模、景区排行榜、天气统计等可视化图表。这段用来证明项目有真实数据支撑。第二段现场跑一条预测流程。选择某景区未来3天的预测调出历史实际流量作对比。建议提前准备几个预测效果好的案例放在数据库里遇到评委指定的随机景区也能从容应对。关键公式和模型结构放在这一段讲清楚。第三段展示推荐结果。找一个有历史行为的账号比较ALS推荐和热门兜底推荐的差异顺手解释冷启动策略。9.2 高频追问的应答框架评委最常问的问题我梳理成一套应答框架你的数据量多大真的需要分布式吗回答要点是单机确实能跑但项目重点是要完整走通分布式存储和计算的链路HDFS解决了多源数据的统一存储Spark完成了跨节点的并行训练这是单机方案无法体现的工程实战价值。为什么选随机森林而不是深度学习回答要点是样本量只有数千条深度学习优势发挥不出来随机森林可解释性强、训练成本低还能输出特征重要性在项目有限的算力条件下是更稳的选择。ALS的原理是什么准备一句话版本把用户和景区映射到同一个隐含特征空间交替固定一方去优化另一方最终用内积得到评分预测。9.3 文档、代码和PPT的整理套路交付物里源码文档PPT讲解四件套文档是很多人不重视但实际拉分最大的。我习惯的文档结构是三分法需求文档讲清楚项目解决什么问题设计文档讲清楚每个模块为什么这么选型、数据怎么流转部署文档则要细到照着敲命令就能复现环境。PPT控制在15页以内项目背景2页、系统架构2页、技术实现6页、成果3页、总结与反思2页。代码整理上建议按模块分目录每个模块附带README说明这个模块的入口是什么、输入输出是什么、运行命令是什么。讲解时可以拿着数据流图从头串到尾不照着PPT念。我后来做类似项目时最深的体会是毕设真正加分的地方不是用了多少新技术而是你能不能把每个选择背后的理由讲清楚。技术栈选得再高级说不清为什么就只是堆名词相反一个朴实但链路完整的项目配上扎实的工程判断反而更能打动评委。这套智慧旅游大数据题目的价值也恰好在这里——它把课程里散落的Hadoop、Spark、爬虫、推荐算法串成了一条完整的业务线做完一遍你才算真正把这些技术拼图装进了自己的知识体系里。
返回列表