☰
Hadoop+Spark+Django电商评价系统全链路实战:从集群搭建到可视化大屏
2026/10/3 20:59:55 网站建设 项目流程

做大数据项目这些年,我见过太多毕设和工程demo,号称“大数据”其实就是在单机MySQL里跑了个聚合查询。但你拿到的这个标题——hadoop+Spark+django基于大数据的电商行业产品评价系统(源码+文档+调试+可视化大屏),是正经把分布式存储、分布式计算、Web服务层和可视化大屏串在一起的全栈大数据项目,不是拼凑demo。这篇文章把我自己从头到尾撸这个系统的思路、踩过的坑、以及可以直接抄作业的配置和代码全部分享出来,不管你是要做毕设、课设,还是想搞明白一个完整的大数据业务链路是怎么跑的,都有参考价值。

先说清楚这套系统到底是干什么的:电商产品评价数据量一旦上来,几千万条评论放在MySQL里,一个带LIKE的模糊查询能把你数据库拖死。Hadoop负责把海量评价数据分布式存下来,Spark负责把这些数据快速清洗和计算,得出评分分布、情感倾向、高频关键词等指标,Django则是把它包装成接口和后台页面,最后把结果扔到可视化大屏上展示。适合人群:正在做大数据方向毕设的学生、想熟悉“数据仓库+计算引擎+Web应用”全链路的后端工程师、以及准备用这个题目参加竞赛的队伍。


1. 整体架构设计与技术选型思路

1.1 为什么是Hadoop+Spark+Django这套组合

拿到需求之后,第一个要搞明白的问题不是“怎么写代码”,而是“每层技术到底在解决什么问题”。你如果只用一个爬虫加Pandas做词频统计,那叫脚本;用上HDFS和Spark集群,才叫大数据系统。这套组合的分工非常清晰:

  • Hadoop(HDFS + YARN):负责底层文件的分布式存储和资源调度。评价数据以日志或JSON文件的形式落入HDFS,后续Spark直接从HDFS读取,而不是从MySQL读。
  • Spark(Spark SQL / DataFrame):负责对原始评价数据做ETL清洗、聚合统计、情感分析打分等计算。它跑在YARN上,利用多节点并行计算能力,处理千万级评价数据时优势明显。
  • Django:负责业务Web层。它提供登录、后台管理、数据查询接口、大屏数据源接口等。Django只管读结果数据,不用碰原始大数据。
  • 可视化大屏:基于ECharts,从Django的API拉取Spark算好的结果,渲染出评分分布、销售趋势、用户画像词云、地域热力图等图表。

选这套组合而不是纯Flask+MySQL的原因有两个:一是数据量级的假设,评价数据达到百万甚至千万级别时,单机会成为瓶颈,分布式是刚需;二是技能树覆盖面的问题,Hadoop和Spark是招聘市场的高频关键词,这部分经验和毕设含金量远高于写SQL。

1.2 模块拆解与数据流转路径

整个系统的数据流向我建议按下面这条线来设计,这也是我实际项目里跑的成熟链路:

采集/生成原始评价数据 → 上传HDFS /data/ecommerce/raw → Spark ETL清洗(去重、过滤、分词、情感打分) → 结果宽表存回HDFS /data/ecommerce/result → Django通过Thrift/API或直读文件方式获取结果 → 缓存在MySQL/Redis → 可视化大屏请求Django接口渲染图表

在这里,Django不直接连HDFS读文件,更不直接连Spark。原因很现实:Django的同步阻塞模型不适合跟Spark Application交互,如果每个大屏请求都去触发一次Spark任务,系统绝对会卡死。正确做法是把Spark计算结果落到MySQL或者Redis,Django只做薄薄的API层。这一步隔离是很多新手架构上翻车重灾区,后面实操环节我会给具体的落库方案。


2. Hadoop与Spark环境搭建与集群部署细节

2.1 伪分布式到集群的跳跃节点

搜热词的人很多在搜“hadoop伪分布式搭建”和“hadoop集群搭建”,说明大家都是从零起步。如果你是单机学习,伪分布式足够跑通开发流程,但如果毕设答辩被问到“你这算大数据吗”,伪分布式配置会显得很单薄,所以我还是建议至少搞一个3节点集群或者用Docker模拟多个节点。

先说Hadoop核心配置,伪分布式切换集群,重点注意这几个文件:

core-site.xml

<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://node01:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/data/hadoop/tmp</value> </property> </configuration>

这里有个我一开始忽略的坑:hadoop.tmp.dir不设置,默认会指向系统/临时目录,重启机器后namenode元数据直接丢失,数据找不回来。必须手动指定一个持久化目录。

hdfs-site.xml

<configuration> <property> <name>dfs.replication</name> <value>3</value> </property> <property> <name>dfs.namenode.secondary.http-address</name> <value>node02:50090</value> </property> </configuration>

伪分布式的时候dfs.replication要改成1,因为只有一个DataNode,写成3会导致副本等待超时。集群环境才配置3。

yarn-site.xml里重点有一个参数yarn.nodemanager.resource.memory-mb,默认会读宿主机的物理内存,如果你服务器内存有限,一定要手动限流,否则集群节点分分钟内存溢出。

完成配置后启动顺序我建议这样:

# 格式化namenode(只在首次执行) hdfs namenode -format # 逐个启动 start-dfs.sh start-yarn.sh # 或者简单点 start-all.sh # 验证 jps

jps能看到NameNode、DataNode、ResourceManager、NodeManager就是正常的。我曾经犯过低级错误:格式化后没有删干净旧的tmp目录,重启后NameNode直接进入安全模式,整个集群无法写入文件。碰到这种情况,先看日志,别急着format,安全模式下执行hdfs dfsadmin -safemode leave应急处理。

2.2 Spark部署与内存调优实战

Spark装起来不难,难的是让它稳定跑在YARN上。首先版本要匹配,以我常用的CDH和Apache版本为例,我是Apache Hadoop 3.3.x配Spark 3.2.x,这个组合很稳。装完后你的spark-env.sh里至少要配置这几个参数:

export SPARK_HOME=/opt/spark export HADOOP_CONF_DIR=/opt/hadoop/etc/hadoop export SPARK_MASTER_HOST=node01 export SPARK_DRIVER_MEMORY=2G export SPARK_EXECUTOR_MEMORY=4G

内存这块特别多坑。比如你执行spark-submit --master yarn --executor-memory 10G,很可能直接报YarnScheduler: Initial job has not accepted any resources,因为每个容器除了executor内存,还要额外开销Overhead内存,而YARN的yarn.scheduler.maximum-allocation-mb没调大,资源不够就挂起。经验公式:executor-memory设的值只是堆内存,实际申请的容器内存大约是它的1.2倍左右。我常用的配置是:

spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 3 \ --executor-memory 4g \ --executor-cores 2 \ --driver-memory 2g \ --conf spark.default.parallelism=12 \ --conf spark.memory.fraction=0.6 \ --conf spark.memory.storageFraction=0.3 \ your_script.py

spark.memory.fraction默认0.6意思是在Executor内存里,最多60%用于执行和存储,剩下的留作预留安全区。如果频繁OOM,优先调大executor-memory,而不是盲目调高num-executors——后者只是增加并行任务数量,单位任务内存不足照样崩。

如果你的服务器本身不大,其实也可以先走本地模式跑通整个流程:

spark-submit --master local[4] spark_etl.py

本地模式适合写代码调试,但注意它不走HDFS也没关系吗?不是,本地模式一样可以读HDFS上的数据,前提是HADOOP_CONF_DIR设置正确,能够拿到core-site.xml里的fs.defaultFS配置。

2.3 Hadoop与Zookeeper集成要点

集群模式下,Hadoop HA(高可用)几乎是必修课,网上热词“hadoop和zookeeper整合实战”就是这块。默认的单NameNode一旦宕机,整个集群不可写。要接入Zookeeper做自动故障转移,JournalNode集群负责同步元数据,两个NameNode一个Active一个Standby。

我的整合步骤是这样的:

  1. 在Zookeeper中创建命名空间(通常自动创建,也可以手动加/hadoop-ha目录)。
  2. 修改hdfs-site.xml增加HA相关配置:
<property> <name>dfs.nameservices</name> <value>mycluster</value> </property> <property> <name>dfs.ha.namenodes.mycluster</name> <value>nn1,nn2</value> </property> <property> <name>dfs.namenode.rpc-address.mycluster.nn1</name> <value>node01:8020</value> </property> <property> <name>dfs.namenode.rpc-address.mycluster.nn2</name> <value>node02:8020</value> </property> <property> <name>dfs.client.failover.proxy.provider.mycluster</name> <value>org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider</value> </property>
  1. core-site.xml里把fs.defaultFS改成hdfs://mycluster。
  2. 配置dfs.ha.automatic-failover.enabled为true,启动zkfc(Zookeeper Failover Controller)进程。

这块最耗时间的坑是:两边NameNode元数据没有完全同步,导致Standby启动后一直处于Safemode或者Out of sync。最佳实践是先在主节点hdfs namenode -initializeSharedEdits,把fsimage推到共享存储的JournalNode上再启动备节点。

如果你只是毕设用,HA不一定要在生产级别跑得很完整,但答辩时能有条理地讲出这个机制,会是一个很大加分项。


3. 电商评价数据处理与Spark核心实现

3.1 数据模型设计与埋点字段规划

做评价分析,首要任务不是急着写Spark代码,而是先设计好你期望的评价数据长什么样。电商评价数据常见的原始字段包括:

  • order_id:订单号,用于防重
  • user_id:用户ID
  • product_id:商品ID
  • product_category:商品类目
  • rating:评分(1-5)
  • review_content:评论文本
  • review_time:评价时间
  • region:用户省份(可选)
  • tags:商家回复标签或用户打标

我用Python脚本生成模拟数据时,规范成JSON格式上传到HDFS。这里字段规划直接影响后续统计口径,比如你中途想加一个“价格区间”维度,结果原始数据里没有,那只能回头补数据,教训很深刻。

3.2 Spark ETL:清洗、去重、聚合与情感标签

说下完整Spark脚本的核心逻辑,用PySpark写比较好理解:

from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, IntegerType spark = SparkSession.builder \ .appName("EcommerceReviewETL") \ .getOrCreate() # 定义schema,读取性能比spark推断schema快2-3倍 schema = StructType([ StructField("order_id", StringType(), True), StructField("user_id", StringType(), True), StructField("product_id", StringType(), True), StructField("rating", IntegerType(), True), StructField("review_content", StringType(), True), StructField("review_time", StringType(), True), StructField("region", StringType(), True) ]) df = spark.read \ .option("multiline", "true") \ .schema(schema) \ .json("hdfs://mycluster/data/ecommerce/raw") # 清洗:去重,过滤空值 df_clean = df.dropDuplicates(["order_id", "user_id", "product_id"]) \ .filter(F.col("review_content").isNotNull()) \ .filter(F.col("rating").between(1, 5)) # 特征工程:增加评价月份、评分分组列 df_clean = df_clean.withColumn("review_month", F.substring("review_time", 1, 7)) \ .withColumn("rating_group", F.when(F.col("rating") >= 4, "正向") \ .when(F.col("rating") == 3, "中性") \ .otherwise("负向"))

这里说下dropDuplicates这个点。原始评价数据不一定干净,同一个用户对不同商品能评,同个订单也可能多次更新,你必须指定去重键。我们用它指定的订单号+用户+商品组合键,并且保留最新一条。如果用distinct()全列去重,可能因为评论时间不同而失效,那等于没去。

聚合统计这一层要产出大屏上能直接用的数据:

# 每日评分趋势 daily_score = df_clean.groupBy("review_month", "rating").count() # 商品评价TOP10 product_top = df_clean.groupBy("product_id").agg( F.count("*").alias("review_cnt"), F.round(F.avg("rating"), 2).alias("avg_rating") ).orderBy(F.desc("review_cnt")).limit(10) # 正负向占比 sentiment_stat = df_clean.groupBy("rating_group").count()

然后统一写回HDFS的parquet格式,parquet的压缩率和查询性能比JSON好太多:

daily_score.write.mode("overwrite").parquet("hdfs://mycluster/data/ecommerce/result/daily_score") product_top.write.mode("overwrite").parquet("hdfs://mycluster/data/ecommerce/result/product_top")

3.3 情感分析维度的附加实现方案

如果你想让系统更有亮点,可以给评论文本加一个简单的情感分析模块。不要被“情感分析”吓到,不用非得上深度学习模型。用基于词典的SnowNLP就够了,跑批的时候对评论文本打一个情感分:大于0.6是正向,小于0.4是负向,中间中性。这样就可以在Spark SQL里做窗口函数统计:

SELECT product_id, SUM(CASE WHEN sentiment = 'positive' THEN 1 ELSE 0 END) as pos_cnt, SUM(CASE WHEN sentiment = 'negative' THEN 1 ELSE 0 END) as neg_cnt FROM review_sentiment GROUP BY product_id

这里要注意的是SnowNLP这种词典模型对电商语境不一定准,“这个产品性价比绝了”这类口语化文本可能会被误判。我提供两个优化方向,你按需选:第一,清洗评论时做停用词过滤;第二,加入自定义的电商领域情感词表,比如“物流快”“客服好”直接标记正向。这两种都不复杂,但效果提升明显。


4. Django服务层与大屏接口实现

4.1 Django工程搭建与数据读取策略

过了Spark这一层,后面就是Web工程师的主场了。Django是Python社区生态最完整的Web框架,自带的Admin后台、ORM和DRF(Django REST Framework)能大幅减少重复开发。创建一个干净的工程:

django-admin startproject review_system cd review_system python manage.py startapp api

按前面说的,Django不直接连Spark,它要读的是Spark输出的结果。这里有两种方案:方案一是Spark把结果直接写入MySQL,Django用ORM查询;方案二是将HDFS结果导出到CSV/JSON后导入MySQL。我个人推荐方案一,也就是在Spark里能直接用jdbc连接MySQL写入:

result_df.write.format("jdbc") \ .option("url", "jdbc:mysql://localhost:3306/review_db?useSSL=false") \ .option("driver", "com.mysql.cj.jdbc.Driver") \ .option("dbtable", "daily_score_stat") \ .option("user", "root") \ .option("password", "your_password") \ .mode("overwrite") \ .save()

注意:Spark的lib目录要放MySQL驱动jar包,否则报ClassNotFound。版本对应MySQL 8.x用mysql-connector-java-8.0.x.jar。

Django这边你只需要定义好模型:

from django.db import models class DailyScoreStat(models.Model): review_month = models.CharField(max_length=7) rating = models.IntegerField() count = models.IntegerField() class Meta: db_table = 'daily_score_stat'

然后写API视图,用DRF的ModelViewSet或者直接JsonResponse都行。大屏接口追求快,所以我在Django里还加了一层Redis缓存,设置120秒过期,这样可以有效减少数据库压力:

import redis import json from django.http import JsonResponse r = redis.Redis(host='127.0.0.1', port=6379, db=0) def dashboard_data(request): cache_key = 'dashboard:overview' data = r.get(cache_key) if data: return JsonResponse(json.loads(data)) # 查库组装数据 data = build_dashboard_data() r.setex(cache_key, 120, json.dumps(data)) return JsonResponse(data)

4.2 可视化大屏前端实现

大屏这块没什么黑魔法,最核心的就是ECharts。你可以用纯HTML+JS,也可以用Vue+Django模板混合。我最常用的方案是Django模板+原生ECharts,部署简单,不需要额外构建前端项目。

页面布局上,我推荐一个经典的三栏式大屏:中间主区域放全品类评分散点图或者趋势折线图,左侧放商品TOP10排行和正负向占比图,右侧放词云和地域分布热力图。用echarts-wordcloud插件做词云之前,需要先由Spark对评论内容做jieba分词和TF-IDF关键词提取。这里有个小的效率点:分词这种CPU密集任务在Spark里做是爽的,但是词表量级很大时要控制输出数量,我一般抽Top 200,大屏展示足够了。

ECharts的数据从接口拉,我推荐图表配置与数据分离,即前端每个图对应一个JS初始化函数,通过fetch拉取Django返回的JSON,再设置到option里:

fetch('/api/dashboard/trend/') .then(res => res.json()) .then(data => { myChart.setOption({ xAxis: { data: data.months }, series: [{ name: '好评数', data: data.positive }] }); });

为了保证大屏不出错,前端一定要做一层数据兜底,比如接口返回空数组时,也要渲染一个带有默认值的图表,否则大屏开着开着某个图变成空白,那在演示的时候很尴尬。

4.3 Django执行查询和对象删除的踩坑记录

热词里有“django执行查询-删除对象”,这块我也给几个经验性的建议。Django ORM的查询和删除看起来简单,但大数据量下别乱用。

比如你要删除几个月前的旧评价记录:

# 不要这么写:遍历对象逐个删除,几万条能把你卡死 old_reviews = Review.objects.filter(created_at__lt='2024-01-01') for r in old_reviews: r.delete() # 要这么写:一次性批量删除 deleted_cnt, _ = Review.objects.filter(created_at__lt='2024-01-01').delete()

批量删除返回受影响行数,一步到位。另外一个常见坑是默认懒加载,查询到的QuerySet不会真正执行SQL,.delete()本身会直接执行数据库删除操作,但如果你在循环里再次查询或修改对象属性,会产生N+1的数据库请求,可以把普通查询改成select_related或prefetch_related来解决。

如果你是执行带聚合的复杂查询,我更建议直接写原生SQL或视图:

from django.db import connection with connection.cursor() as cursor: cursor.execute(""" SELECT DATE_FORMAT(c.created_at, '%Y-%m') as month, AVG(c.rating) as avg_score FROM api_comment c GROUP BY month """) rows = cursor.fetchall()

Django ORM对付复杂统计分析写起来很绕,Native SQL对比起来反而清晰直观。


5. 调试、部署与常见问题排查

5.1 调试技巧:从Spark到Django的链路追踪

做完一个完整链路之后,调试是噩梦环节。我总结出一套好用的排查顺序:自底向上,先数据,再计算,再存储,最后接口。

比如大屏显示“商品TOP10”没有数据,我会这样查:

  1. 检查HDFS原始数据是否存在:hdfs dfs -ls /data/ecommerce/raw
  2. 检查Spark ETL是否跑成功:看YARN日志yarn logs -applicationId xxx
  3. 检查结果parquet文件是否生成:hdfs dfs -ls /data/ecommerce/result
  4. 检查Spark写入MySQL的数据:select * from daily_score_stat limit 10
  5. 最后才看Django API返回:curl 'http://127.0.0.1:8000/api/dashboard/top/'

这个顺序能在5分钟内定位问题所在层级。我有一次怎么都查不出前端为什么没数据,最后发现Django接口返回的是NaN,而ECharts不认识NaN,渲染直接挂了。这提醒我:后端要做规格化处理,把None、NaN全部转成0。

5.2 高频报错与解决方案速查表

我把整个系统开发过程中最常见的问题整理成一张速查表,这会节省你大量查资料的精力:

报错场景关键错误提示解决方法
Hadoop启动后DataNode起不来Incompatible clusterIDs删除data目录下的current/VERSION文件后重新格式化
Spark连接HDFS超时java.net.ConnectException: Connection refused检查防火墙和core-site.xml的9000端口
Spark执行OOMJava heap space / Container killed调大executor-memory并降低并行任务数
Spark写MySQL失败ClassNotFoundException com.mysql.jdbc.Driver确认驱动jar是否放在$SPARK_HOME/jars目录
Django跨域大屏请求失败CORS policy安装django-cors-headers,配置CORS_ORIGIN_ALLOW_ALL
MySQL数据量稍微大一点ORM就很慢N+1 queries使用select_related/prefetch_related或原生SQL
ECharts词云图中文不显示字体加载异常确保开发者工具控制台没有字体404问题,使用自定义富文本字体样式

这里特别说下Hadoop的Incompatible clusterIDs问题。它的诱因是格式化NameNode后,DataNode里遗留了旧集群的ClusterID,再去启动就会拒绝。不要一上来就删除整个HDFS数据,风险很大,正确操作是找到DataNode数据目录中的current/VERSION,把clusterID改成和NameNode一致的,再重启DataNode。

5.3 上线部署的经验之谈

Django跑开发服务器只适合本地调试,正式大屏展示建议直接用uWSGI+Gunicorn+Nginx的组合。Nginx能扛住静态文件和并发连接,Django只处理动态请求。我的一个配置参考是:

[uwsgi] chdir = /opt/review_system module = review_system.wsgi:application master = true processes = 4 harakiri = 30 socket = /tmp/review.sock chmod-socket = 664 vacuum = true

Nginx这边把大屏页面的路由直接代理到uwsgi的socket上,并且给静态资源设置expires缓存,大屏刷新速度会快很多。

另外要提醒一句:整个系统部署时,一定要在环境变量和配置文件里把Hadoop/Spark相关的路径、内存参数独立拆到config.py里,别写死。不然换个服务器要重新改一段很长的代码,改到怀疑人生。


实战经验总结与后续扩展建议

我前后完整地做过不止一次类似系统,体会是很深的:大屏只是表象,分布式存储和计算链路才是系统的脊梁。给正在动手的你三点建议:第一,环境搭建阶段不要追求多个组件一步到位,先把Hadoop单节点跑通,再搭Spark,再写Django,层次推进;第二,数据一定不要只用几条测试数据糊弄,要用脚本生成至少几十万条以上的模拟数据,否则Spark的分布式优势根本体现不出来,答辩时数据规模也会被问住;第三,不要忽视文档的整理,按我上面这个链路去写系统设计文档,每一层解决什么问题、为什么选这个组件写清楚,这份文档的价值不亚于代码本身。

最后再分享一个小技巧,在调试Spark任务时,尽量把spark.sql.shuffle.partitions调成和集群CPU核心数接近的值,默认200,在数据量不大时会造成大量小文件碎片,而且很多task空转,白白浪费时间。我一开始没注意,直到看了Spark UI才发现大部分executor都在处理几乎为空的任务,调成12之后整个ETL时间缩短接近60%。这种细节在实战里往往比调大内存更有用。

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

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

立即咨询