Hadoop+Spark+Django电商评价系统全链路实战:从集群搭建到可视化大屏

发布时间:2026/10/3 21:00:17
Hadoop+Spark+Django电商评价系统全链路实战:从集群搭建到可视化大屏 做大数据项目这些年我见过太多毕设和工程demo号称“大数据”其实就是在单机MySQL里跑了个聚合查询。但你拿到的这个标题——hadoopSparkdjango基于大数据的电商行业产品评价系统(源码文档调试可视化大屏)是正经把分布式存储、分布式计算、Web服务层和可视化大屏串在一起的全栈大数据项目不是拼凑demo。这篇文章把我自己从头到尾撸这个系统的思路、踩过的坑、以及可以直接抄作业的配置和代码全部分享出来不管你是要做毕设、课设还是想搞明白一个完整的大数据业务链路是怎么跑的都有参考价值。先说清楚这套系统到底是干什么的电商产品评价数据量一旦上来几千万条评论放在MySQL里,一个带LIKE的模糊查询能把你数据库拖死。Hadoop负责把海量评价数据分布式存下来Spark负责把这些数据快速清洗和计算得出评分分布、情感倾向、高频关键词等指标Django则是把它包装成接口和后台页面最后把结果扔到可视化大屏上展示。适合人群正在做大数据方向毕设的学生、想熟悉“数据仓库计算引擎Web应用”全链路的后端工程师、以及准备用这个题目参加竞赛的队伍。1. 整体架构设计与技术选型思路1.1 为什么是HadoopSparkDjango这套组合拿到需求之后第一个要搞明白的问题不是“怎么写代码”而是“每层技术到底在解决什么问题”。你如果只用一个爬虫加Pandas做词频统计那叫脚本用上HDFS和Spark集群才叫大数据系统。这套组合的分工非常清晰HadoopHDFS YARN负责底层文件的分布式存储和资源调度。评价数据以日志或JSON文件的形式落入HDFS后续Spark直接从HDFS读取而不是从MySQL读。SparkSpark SQL / DataFrame负责对原始评价数据做ETL清洗、聚合统计、情感分析打分等计算。它跑在YARN上利用多节点并行计算能力处理千万级评价数据时优势明显。Django负责业务Web层。它提供登录、后台管理、数据查询接口、大屏数据源接口等。Django只管读结果数据不用碰原始大数据。可视化大屏基于ECharts从Django的API拉取Spark算好的结果渲染出评分分布、销售趋势、用户画像词云、地域热力图等图表。选这套组合而不是纯FlaskMySQL的原因有两个一是数据量级的假设评价数据达到百万甚至千万级别时单机会成为瓶颈分布式是刚需二是技能树覆盖面的问题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或者RedisDjango只做薄薄的API层。这一步隔离是很多新手架构上翻车重灾区后面实操环节我会给具体的落库方案。2. Hadoop与Spark环境搭建与集群部署细节2.1 伪分布式到集群的跳跃节点搜热词的人很多在搜“hadoop伪分布式搭建”和“hadoop集群搭建”说明大家都是从零起步。如果你是单机学习伪分布式足够跑通开发流程但如果毕设答辩被问到“你这算大数据吗”伪分布式配置会显得很单薄所以我还是建议至少搞一个3节点集群或者用Docker模拟多个节点。先说Hadoop核心配置伪分布式切换集群重点注意这几个文件core-site.xmlconfiguration property namefs.defaultFS/name valuehdfs://node01:9000/value /property property namehadoop.tmp.dir/name value/data/hadoop/tmp/value /property /configuration这里有个我一开始忽略的坑hadoop.tmp.dir不设置默认会指向系统/临时目录重启机器后namenode元数据直接丢失数据找不回来。必须手动指定一个持久化目录。hdfs-site.xmlconfiguration property namedfs.replication/name value3/value /property property namedfs.namenode.secondary.http-address/name valuenode02: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 # 验证 jpsjps能看到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_HOSTnode01 export SPARK_DRIVER_MEMORY2G export SPARK_EXECUTOR_MEMORY4G内存这块特别多坑。比如你执行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.parallelism12 \ --conf spark.memory.fraction0.6 \ --conf spark.memory.storageFraction0.3 \ your_script.pyspark.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。我的整合步骤是这样的在Zookeeper中创建命名空间通常自动创建也可以手动加/hadoop-ha目录。修改hdfs-site.xml增加HA相关配置property namedfs.nameservices/name valuemycluster/value /property property namedfs.ha.namenodes.mycluster/name valuenn1,nn2/value /property property namedfs.namenode.rpc-address.mycluster.nn1/name valuenode01:8020/value /property property namedfs.namenode.rpc-address.mycluster.nn2/name valuenode02:8020/value /property property namedfs.client.failover.proxy.provider.mycluster/name valueorg.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider/value /propertycore-site.xml里把fs.defaultFS改成hdfs://mycluster。配置dfs.ha.automatic-failover.enabled为true启动zkfcZookeeper Failover Controller进程。这块最耗时间的坑是两边NameNode元数据没有完全同步导致Standby启动后一直处于Safemode或者Out of sync。最佳实践是先在主节点hdfs namenode -initializeSharedEdits把fsimage推到共享存储的JournalNode上再启动备节点。如果你只是毕设用HA不一定要在生产级别跑得很完整但答辩时能有条理地讲出这个机制会是一个很大加分项。3. 电商评价数据处理与Spark核心实现3.1 数据模型设计与埋点字段规划做评价分析首要任务不是急着写Spark代码而是先设计好你期望的评价数据长什么样。电商评价数据常见的原始字段包括order_id订单号用于防重user_id用户IDproduct_id商品IDproduct_category商品类目rating评分1-5review_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和DRFDjango REST Framework能大幅减少重复开发。创建一个干净的工程django-admin startproject review_system cd review_system python manage.py startapp api按前面说的Django不直接连Spark它要读的是Spark输出的结果。这里有两种方案方案一是Spark把结果直接写入MySQLDjango用ORM查询方案二是将HDFS结果导出到CSV/JSON后导入MySQL。我个人推荐方案一也就是在Spark里能直接用jdbc连接MySQL写入result_df.write.format(jdbc) \ .option(url, jdbc:mysql://localhost:3306/review_db?useSSLfalse) \ .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_length7) 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(host127.0.0.1, port6379, db0) 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。你可以用纯HTMLJS也可以用VueDjango模板混合。我最常用的方案是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__lt2024-01-01) for r in old_reviews: r.delete() # 要这么写一次性批量删除 deleted_cnt, _ Review.objects.filter(created_at__lt2024-01-01).delete()批量删除返回受影响行数一步到位。另外一个常见坑是默认懒加载查询到的QuerySet不会真正执行SQL.delete()本身会直接执行数据库删除操作但如果你在循环里再次查询或修改对象属性会产生N1的数据库请求可以把普通查询改成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”没有数据我会这样查检查HDFS原始数据是否存在hdfs dfs -ls /data/ecommerce/raw检查Spark ETL是否跑成功看YARN日志yarn logs -applicationId xxx检查结果parquet文件是否生成hdfs dfs -ls /data/ecommerce/result检查Spark写入MySQL的数据select * from daily_score_stat limit 10最后才看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_ALLMySQL数据量稍微大一点ORM就很慢N1 queries使用select_related/prefetch_related或原生SQLECharts词云图中文不显示字体加载异常确保开发者工具控制台没有字体404问题使用自定义富文本字体样式这里特别说下Hadoop的Incompatible clusterIDs问题。它的诱因是格式化NameNode后DataNode里遗留了旧集群的ClusterID再去启动就会拒绝。不要一上来就删除整个HDFS数据风险很大正确操作是找到DataNode数据目录中的current/VERSION把clusterID改成和NameNode一致的再重启DataNode。5.3 上线部署的经验之谈Django跑开发服务器只适合本地调试正式大屏展示建议直接用uWSGIGunicornNginx的组合。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 trueNginx这边把大屏页面的路由直接代理到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%。这种细节在实战里往往比调大内存更有用。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询