☰
基于Spark+Hive+LLM的农产品价格预测与推荐系统实战
2026/10/1 12:31:42 网站建设 项目流程

1. 为什么最终选择了这套技术组合

1.1 项目需求到底长什么样

先说结论:这个题目的本质不是“做一个价格预测算法”,而是“做一个以农产品数据为核心的复合型Web系统”。很多同学一看“Spark+Hadoop+Hive+LLM+Django”就懵了,觉得什么都要会,其实拆开来看,它要解决的是四个层面的问题:数据怎么存、数据怎么算、预测和推荐怎么做、结果怎么展示给用户。

毕业设计的评分一般看三块:工作量、完整度、技术难度。纯用Python写个线性回归做价格预测,页面用Bootstrap随便搭一下,也能过,但拿不到高分。而把分布式存储、数据仓库、计算引擎、大模型应用、Web框架全串起来,这就是一个典型的“工业级玩具系统”,既能体现对大数据工具链的掌握,又能蹭上LLM的热点。我当时选这个题目,就是因为它覆盖足够宽,可以讲的故事足够多。

1.2 技术选型背后的逻辑链

选Spark、Hadoop、Hive,不是炫技,而是为了处理“农产品价格和销量”这类带时间戳、带地域、带品类维度的数据。比如全国各地的白菜价格,每天一个记录,一年下来也就十几万行,这点数据量单机数据库完全扛得住,但毕设的意义在于展示“你能处理更大规模数据的能力”。所以我们需要一个合理的理由把这些技术串起来。

Hadoop提供HDFS做底层存储,Hive把结构化的价格、销量数据管理成表,Spark负责做ETL清洗和特征计算,Django作为Web后端把结果暴露给前端,LLM则承担“智能问答”和“推荐理由生成”这类需要自然语言能力的任务。这个链路在真实工业界也是成立的,只是规模不同。从学习角度讲,你通过这个项目至少能把大数据生态的“存储层-计算层-应用层”跑通一遍。

1.3 这套组合适合谁来参考

如果你正在准备计算机毕业设计,或者想从零搭建一个包含大数据处理和Web应用的项目,这篇文章应该能帮你少走几个月的弯路。我会把架构设计、集群搭建、数仓建模、Spark特征工程、预测模型、推荐算法、Django集成,以及我踩过的各种坑都掰开揉碎讲一遍。哪怕你之前没接触过Hadoop,看完整篇也能照着复现一个可运行的系统,因为每个环节我都会给出具体的操作路径和参数的推荐值。

2. 系统架构与数据流设计

2.1 分层思想:别把所有代码塞进一个Django项目

这个系统我强烈建议分层,不要试图在Django里直接调用Hadoop API。我当时把系统分成五层:数据采集层、数据存储层、离线计算层、算法服务层、Web应用层。

采集层用Python脚本爬取公开的农产品批发市场数据(比如价格行情、成交量),或者用CSV模拟数据;存储层就是HDFS和Hive表;计算层用Spark做ETL、特征工程,以及调用训练好的预测模型进行批量预测;算法服务层单独跑一个Flask/FastAPI服务,封装价格预测、销量预测、推荐和LLM问答的接口;Django只负责从这个服务取数据,再渲染到前端页面。这样做的好处是,每一块都能独立测试,答辩论的时候也更好讲清模块边界。

2.2 数据流:从原始数据到用户看到的图表

我实际跑通的数据流是这样的:爬虫脚本每天把各省农产品的“品类、价格、销量、日期、产地”写入HDFS的原始目录;然后用Hive建立外部表指向这个目录,再通过一条INSERT OVERWRITE语句把数据清洗后落入数仓的ODS层和DWD层;Spark从Hive读取DWD层数据,计算滑动平均、同比环比、节假日标记等特征,保存为特征表;模型服务在每天凌晨定时调用Spark训练的模型,输出未来三天的价格预测结果,写回Hive或MySQL;Django后端再按日期从MySQL读取预测结果,通过ECharts绘制趋势图。

LLM的角色出现在两个地方:一是在推荐模块中根据用户跟农产品的交互记录生成推荐理由,比如“您最近购买了苹果,而且山东烟台苹果当前价格处于近三个月低位,适合囤货”;二是在全局搜索框中做一个智能助手,回答“西红柿最近为什么涨价”、“西瓜和荔枝哪个更适合夏天卖”这类自然语言问题。看明白这个流向后,你就知道每个组件其实只负责一个很小的环节,并没有想象中那么难。

2.3 为什么需要LLM而不是纯规则

农产品推荐和电商推荐最大的区别是“时效性”和“解释性”。用户买农产品通常受季节、价格波动影响很大,规则引擎可以做到“近七日价格上涨的品类不推荐”,但无法给出一句人性化的解释。大模型在这里的价值不是替代预测算法,而是把预测结果和用户行为组织成一段可读性强的文案。我在项目里接入的是国内可用的开源大模型API,通过LangChain封装了一个简单的问答链,把数据库查询出来的最新行情数据作为上下文,再让大模型生成结论。这种方式实现起来比微调模型省事得多,而且对于毕设来说已经足够亮眼。

3. 大数据环境搭建与数据仓库实现

3.1 Hadoop环境选择:伪分布式还是集群

很多教程开局就让你准备三台虚拟机搭集群,其实对毕设来说完全是浪费。我第一次搭的是三节点集群,结果光调试SSH免密登录就花了两天,后来个人做实验直接改用Hadoop伪分布式模式。所谓伪分布式,就是在一个节点上同时跑NameNode、DataNode、ResourceManager、NodeManager,让所有进程都跑在同一台机器上。只要你的物理机内存有16G,伪分布式完全够用,而且跑Spark任务比集群还方便,因为不需要考虑网络间传输。

如果你非要搭集群,我建议用Docker Compose一次性拉起三个容器,每个容器配置2G内存,挂载一个数据卷。我在网上看到很多“从零开始安装Hadoop”的教程,大多会让你改一堆配置文件,其实核心就四个文件:core-site.xml、hdfs-site.xml、yarn-site.xml、mapred-site.xml。伪分布式模式下,只需要把core-site.xml里的fs.defaultFS设为hdfs://localhost:9000,hdfs-site.xml里把副本数从3改成1,YARN的资源配置调低一点就行。

3.2 Hive数仓建模:分区表和小文件治理

Hive在这一整套链路里是作为Spark的数据源和落地端存在的。我第一版直接建了普通表,结果Spark读取时扫描全表,查询速度慢得离谱。后来改成按日期分区,每天一个分区,查询时用WHERE date_partition='2025-03-01'就能只扫一个目录。分区字段类型建议用STRING而不是DATE,避免Hive和Spark之间的类型转换麻烦。

建表语句给大家参考:

CREATE EXTERNAL TABLE ods_market_daily ( product_name STRING, market_name STRING, price DECIMAL(10, 2), sales_volume INT, origin STRING, unit STRING, record_date STRING ) PARTITIONED BY (date_partition STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' LOCATION '/warehouse/ods/market_daily';

这里有个重要经验:写入分区数据时,如果当天文件特别多,HDFS上会产生大量小文件,导致Spark读取时频繁启动task,效率极低。我在项目里用Spark的coalesce(1)把每天的数据合并成一个文件再落库,或者在Hive侧设置“hive.merge.smallfiles.avgsize”参数控制合并阈值。这个坑不亲自踩一次,很难理解为什么明明只有几百M的数据,跑起来比几G还慢。

3.3 Spark核心任务:数据清洗与特征工程

Spark在整个系统里最实际的作用是“算特征”。光有原始价格数据,模型是没法学习的,因为模型需要知道今天的价格相对于昨天涨了没有、近一周的波动率是多少、该品类往年在同一天的平均销量是多少。

我用Spark SQL写了一段特征生成逻辑,核心是窗口函数的应用。比如要算每个品类近七天的价格均值,可以这样写:

SELECT product_name, record_date, price, AVG(price) OVER (PARTITION BY product_name ORDER BY record_date ROWS BETWEEN 6 PRECEDING AND CURRENT ROW) AS price_ma7 FROM dwd_market_fact

这个price_ma7就是价格移动平均,后面可以作为预测模型的一个输入特征。窗口函数在Hive和Spark里语法基本一致,但要注意如果你用的是Hive老版本,可能需要开启Hive窗口函数支持。实际写任务时,我还会计算销量七日差分、价格环比变化率、是否周末、是否节假日等十几个维度,最后统一写入特征表hive_feature。

这里给新手一个建议:Spark任务不要一次性编成一大段脚本,最好拆成三到四个阶段,每个阶段单独写一个object,用checkpoint隔开。因为一旦某一步报错,直接重跑那个阶段就行,不用从头再来。

3.4 直接能抄的Spark读取Hive代码

在Spark子项目里,我用的是SparkSQL读取Hive表,需要提前把hive-site.xml放到Spark的conf目录下,然后启用Hive支持:

SparkSession spark = SparkSession.builder() .appName("FeatureEngineering") .enableHiveSupport() .config("spark.sql.warehouse.dir", "hdfs://localhost:9000/user/hive/warehouse") .getOrCreate(); Dataset<Row> df = spark.sql("SELECT product_name, record_date, price, sales_volume, origin FROM dwd_market_fact WHERE date_partition >= '2024-01-01'"); df.createOrReplaceTempView("market_data"); Dataset<Row> featureDf = spark.sql( "SELECT product_name, record_date, price, " + "AVG(price) OVER (PARTITION BY product_name ORDER BY record_date ROWS BETWEEN 6 PRECEDING AND CURRENT ROW) AS price_ma7, " + "sales_volume - LAG(sales_volume, 1) OVER (PARTITION BY product_name ORDER BY record_date) AS sales_diff " + "FROM market_data" ); featureDf.write().mode("overwrite").saveAsTable("dwd_feature_table");

写完这个任务后,用spark-submit提交时记得带上Hive的依赖包:--packages org.apache.spark:spark-hive_2.12:3.2.0。Graphically,我刚开始没加这个包,直接报ClassNotFound,都是一点点踩出来的。

4. 价格与销量预测模型详解

4.1 预测什么、预测多远,比用什么模型更重要

首先要明确目标:价格预测和销量预测是两个不同的问题。价格预测是回归问题,目标是输出未来一天的精确价格;销量预测也类似,但它的数值波动更大,受促销、节假日影响明显。我建议毕设只做“未来三天的预测”,不要做长期预测,因为农产品市场价格规律性不强,预测十天以上基本就是玄学,答辩时也没法解释清楚。

我对价格预测采用了两阶段方案:先用Spark训练一个随机森林回归模型,然后把预测结果和真实历史数据的偏差作为特征,再用一个简单的线性回归做修正。这样看起来比单一模型更有层次感,而且可以在论文里写“集成学习策略”。销量预测则用了梯度提升树,原因很简单:销量的特征中很多是类别型变量(品类、产地、是否节假日),树模型对这种混合特征处理得比线性模型好。

4.2 特征工程的实用细节

特征工程直接决定预测效果的上限。我最终用的特征分为四类:统计类(近7日均价、近30日均价、环比、同比)、日期类(星期几、是否月初、是否节假日、节气)、品类类(品类ID编码、产地编码)、外部类(当天气温、降雨量,通过爬天气接口获得)。注意,这些特征必须在预测时也能拿到,不能用未来数据,否则就是数据泄露。

比如你要预测3月5日白菜价格,那么3月5日当天的天气可以提前从天气预报接口拿到,所以它可以作为特征。但3月5日当天的实际销量就不能用于价格预测,因为到预测时点还不知道。我在第一版模型里不小心把当天的销量滞后特征也塞进去了,训练时表现极好,实盘时完全失灵,后来排查数据泄露花了整整一个下午。这算是一个血泪教训。

4.3 模型训练与超参调节

由于特征表在Hive里,训练代码用Spark MLlib会比较顺。但Spark MLlib的随机森林调参并不是很直观,我用的是网格搜索加交叉验证,核心参数是maxDepth、numTrees、maxBins。对于这种数据量不大的场景,maxDepth设为10、numTrees设为200基本就够,再大不仅训练慢,还容易过拟合,验证集上效果反而变差。

在实际效果上,我的价格预测模型RMSE大约在0.4元左右,对于均价3元多的蔬菜来说,误差率已经能控制在15%以内。销量预测因为波动大,误差率在25%左右,勉强能用。如果你想进一步提升,可以尝试用Prophet做时间序列基线,然后和树模型做加权平均。我试过之后发现,短期预测树模型优势明显,长期趋势Prophet更稳,但毕设不需要追求极致效果,稳定的流程比精度更重要。

5. 农产品推荐系统的设计与实现

5.1 推荐场景跟电商有什么不同

农产品推荐最大的特点是“价格敏感”和“季节性敏感”。我卖浏览器的推荐只会根据你的点击记录推相似商品,但农产品不一样,上个月推荐草莓可能合理,这个月草莓下市了,再推就是坑人。所以我的推荐系统必须把价格趋势和时令因素融合进去。具体做法是:先做一个基于物品的协同过滤,得到候选列表,然后用一套规则和LLM对候选列表做重排和解释。

5.2 协同过滤实现细节

由于是毕设,我不建议用复杂的图神经网络,就用传统的物品协同过滤(ItemCF)就够了。用户-物品交互数据可以从订单表里抽取:用户ID、商品ID、购买时间、购买数量。ItemCF的核心是计算物品之间的相似度矩阵,通常用余弦相似度。我用Spark SQL写出“共同购买次数”矩阵,再转成相似度,整个过程大概一百行代码。

实现完ItemCF后,每个用户都能拿到一个“看过A也爱买B”的候选列表。但我发现一个问题:如果用户只买过一次红薯,推荐结果会偏向红薯相关品类,完全忽略他现在最可能需要的秋月梨。于是我又加了时令权重:当候选商品的当月销量在所有同类商品中排名前20%时,给它的推荐分乘上1.5倍权重。这样推荐列表就同时有了“个性化”和“应季性”。

5.3 用LLM生成推荐理由

推荐理由是这个项目相对有亮点的地方。我拒绝用固定模板“根据您的购买记录,为您推荐XXX”,因为那很呆。我的做法是:先用规则引擎从数据库里找出推荐商品的三个卖点,比如“近七天价格下降了8%”、“目前山东产区走货速度快”、“您上个月买过同品类商品”,然后把这三个卖点作为提示词塞给大模型,让它组织成一两句自然的话。

实际效果类似这样:“哥哥,你上次买的苹果是山西红富士,最近山东烟台的红富士批量上市,均价每斤便宜了5毛钱,而且走货量连续三天上升,这时候囤一点正好赶上清明节前的价格低点。”这种文案比冷冰冰的模板有感染力得多。实现方式也不难,用Django后端调用大模型API,把规则拼好的Prompt传过去,再设置temperature=0.7,避免回答太死板。

6. Django系统集成与可视化

6.1 Django项目的模块划分

Django端我建了四个app:market、analysis、recommend、assistant。market负责农产品行情展示,就是每天价格和销量的表格、趋势图;analysis负责预测结果页面;recommend负责“为你推荐”页面;assistant负责LLM智能问答页面。这种划分非常清晰,答辩时问到哪里你都不会乱。

数据访问方面,我对接的是MySQL。很多新手以为Django必须直接连Hive,其实根本没必要。Hive和Spark算好的结果,比如预测值、推荐列表,都用定时脚本同步到MySQL的几张表里。Django只对MySQL做增删改查,这样性能上毫无压力,而且Django的ORM用起来非常舒服,不用自己拼SQL。

6.2 异步任务与WebSocket推送

系统有一个“生成预测报告”的功能,按下按钮后需要跑Spark任务,可能要等几十秒,这时候不能让前端一直转圈。我用了Celery加Redis处理异步任务,Django提交一个任务后立刻返回“正在计算中”,等Celery worker完成后再通过WebSocket把结果推给前端。这里涉及Python Django WebSocket实现后台有数据前端推送的典型场景,不用成熟的Django Channels也能做,但Channels是标准方案。

我碰到一个经典问题:Celery worker里调用Spark的Java类库,会碰到路径和依赖问题。建议不要直接在worker里用pyspark,更稳妥的做法是worker通过subprocess调用已经打包好的spark-submit命令,然后把输出文件路径或状态码写回Redis。这样解耦了Python环境和Spark环境的差异。

6.3 页面可视化与前端框架选择

前端我没有用复杂的Vue或React,就用Django模板加Bootstrap,图表用ECharts。因为毕设的关键是功能闭环,不是前端炫技。ECharts的折线图可以直接接收后端传过来的JSON,像这样:

$.get('/analysis/price_trend/', {product: 'apple'}, function(res) { var chart = echarts.init(document.getElementById('chart')); chart.setOption({ xAxis: { type: 'category', data: res.dates }, yAxis: { type: 'value' }, series: [{ type: 'line', data: res.prices }] }); });

预测页面则是把“真实价格”和“预测价格”画成两条线,一眼就能看出模型拟合程度。推荐页面则是一个卡片流,每张卡片包括商品图、名称、现价、推荐理由,推荐理由由LLM动态生成,所以每次刷新可能句子都不一样,这种“非确定性”反而让答辩老师觉得你的系统真的有AI能力。

6.4 权限管理与发布部署

系统里有普通用户和管理员两种角色,管理员可以管理商品信息、触发Spark重跑任务;普通用户只能看行情、看预测、用推荐。Django自带的auth系统够用,我给两组用户建了不同的Group,在视图上用了装饰器@login_required和@permission_required控制访问。部署时我用Nginx加uWSGI跑在服务器上,Celery worker和Redis后台运行,Spark任务是手动或定时触发,没有像生产环境那样做完整的调度,但架构是完整的。

7. 踩坑记录与问题排查实录

7.1 Hadoop集群启动失败的排查思路

伪分布式模式最常见的坑就是启动后DataNode起不来。我的排查路径比较有代表性:先执行hdfs dfs -ls /看能不能访问,如果报错说NameNode在安全模式,就等一会儿或运行hdfs dfsadmin -safemode leave;如果DataNode进程不在,大概率是hdfs-site.xml里dfs.data.dir指向的目录没有写权限,或者残留的VERSION文件跟clusterID不匹配。解决办法是把tmp下的hdfs目录清空重新格式化。

这里提醒一个细节:格式化NameNode之前,一定要把dfs.namenode.name.dir和dfs.datanode.data.dir指向的目录也清干净,否则格式化之后datanode的clusterID跟namenode不一致,就是启动失败。

7.2 Spark内存问题与监测工具

我在跑特征任务时经常遇到Executor lost或者OOM。毕设阶段,我直接在spark-submit里设置了--driver-memory 2g --executor-memory 2g,如果还要更大,就得看物理机内存是否够。后来我发现一个很小的技巧:在Spark UI的Executors页面看每个Executor的Shuffle Read大小和GC时间,如果GC时间占比高,说明内存确实紧张,要调大executor内存或者减少并行度,而不是一味加资源。

还有一点,如果同时开着Hadoop的NameNode、DataNode、ResourceManager以及Spark,电脑会卡到怀疑人生。我建议在一个独立的Shell里跑一个内存监测工具,随时看Java进程占了多少RSS,及时关掉不需要的节点。这个项目本身就是要多工具配合,所以资源管理也是必修课。

7.3 Hive查询慢和小文件优化

前面提到小文件问题,再补充一下我在Hive端的优化实践。Hive默认的fetch task有时候会因为小文件过多退化成MapReduce,速度慢出天际。我在hive-site.xml里设置了hive.fetch.task.conversion=more,可以在简单查询时直接读文件不走MR。同时设置hive.merge.mapredfiles=true和hive.merge.size.per.task=256000000,让MapReduce跑完后自动合并输出文件。这一套做完,查询速度至少提升三倍。

7.4 Django数据库连接与耗时操作处理

Django默认是多线程服务,如果直接在视图里跑大规模ORM查询,数据库连接池会被占完。我的做法是把复杂的统计查询放到Redis缓存里,比如“某品类近七日价格走势”这个接口,第一次从MySQL查完之后缓存五分钟,后面直接读Redis。触发Spark重算的接口一定不要用同步请求,否则前端等几十秒直接超时,用户体验极差。用Celery异步任务后,前端用一个轮询获取任务状态,这样最稳妥。

8. 从项目本身延展出的几个升级空间

8.1 实时流计算:让价格预警成为可能

目前这个项目是离线预测,每天跑一次。如果接入Kafka,把各个市场的价格数据实时发送到Spark Streaming,就能实现“价格突变预警”。比如某市场白菜价格30分钟内上涨20%,系统立刻推送通知给订阅用户。这个功能并不难实现,在已有的Hadoop和Spark基础上,只需要加一个Kafka生产者脚本和一个Structured Streaming任务。如果你想把毕设方向再拔高,这是一个非常自然的扩展点。

8.2 模型自动重训与回归检测

离线模型用的是月度重训,但农产品价格随着季节变化,每个月市场模式可能都不同。可以写一个定时任务,每天对比预测值和真实值,当连续三天的误差超过阈值时,自动触发Spark重训任务。我在本地试过用Crontab加Shell脚本检测Spark任务退出码,失败就发送企业微信通知,效果很好。这一套东西放到论文里,就是“持续学习”章节,评委一定会追问细节,所以你要把重训的数据流程搞清楚。

8.3 数据大屏与移动端适配

毕设结尾如果能做一个投屏用的数据大屏,整体效果会非常加分。大屏上展示全国农产品价格热力图、当日涨价榜、主力推荐商品Top10,数据接口直接从Django的API读取,完全不需要单独的框架。我自己的体会是,大屏做出来以后,连自己都觉得这个项目“有没有用”这件事不重要了,“好不好看”才是实实在在的完成度。如果再打包一个PWA应用,让用户在手机上也能看到预测结果,那工作量也不会增加太多。

最后再分享一个小技巧:在答辩之前,一定要把Hadoop、Spark的冷启动时间算进去,提前把服务拉起,免得现场演示时一直转圈。你的系统越复杂,越容易在关键时刻掉链子,准备一个录屏的备用视频,永远比现场手忙脚乱要强。我自己就是在答辩前录了一份完整演示视频,结果现场服务器内存不足,幸好有备份才稳住局面。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询