最近我手里正好在折腾一个完整的课设项目:Hadoop+Spark+Django基于Python的重庆旅游景点数据分析系统,标题里把源码、文档、调试、可视化大屏都点出来了,一看就是冲着毕业设计/课程设计去的。这套技术栈在近几年的数据类课设里出现频率非常高,但说句实话,很多同学卡住的不是某个框架怎么调用,而是Hadoop、Spark、Django、可视化大屏这几层怎么串成一条完整链路。
这篇文章不打算做理论复述,我就按自己实际做完这个项目的顺序,把数据采集、HDFS存储、Spark分析、Django接口、ECharts大屏、以及调试过程里踩过的坑一次性写清楚。适合两类人看:一是准备做大数据方向课设/毕设,需要快速搭一个能演示、能答辩的完整系统的同学;二是已经能单点使用Spark或Django,但还不清楚多个组件如何协同工作的人。我会尽量把“为什么这样做”也讲明白,而不只是给你一段能跑就行的代码。
1. 为什么是"Hadoop + Spark + Django"这套组合,而不是别的
1.1 系统到底要解决什么问题
先回到需求本身。系统名称是“重庆旅游景点数据分析系统”,核心要做的是:采集重庆主要旅游景点的基础数据,对景点热度、评分、票价、区域分布等维度做分析,最后通过网页大屏把分析结果可视化展示出来。
数据量本身不大,重庆A级景区加热门打卡点撑死几百条,评论数据多也就几万条。这个量级用一台普通电脑、一个MySQL数据库就能算完。那为什么还要上Hadoop和Spark?因为这不是一个纯生产项目,而是一个学习型、展示型的大数据项目。它的价值不只是“算出结果”,而是完整地展示“大数据处理流程”。换句话说,Hadoop管存储,Spark管计算分析,Django管Web服务和数据展示,每一层各司其职,最终形成一个从数据到报表的闭环,这才是这套技术栈组合的真正意义。
如果换成纯爬虫+Flask+ECharts,两天就能做完,但体现不出分布式存储和分布式计算的内容;如果只有Hadoop+Spark,又没有一个Web端把结果展示出来,项目不够完整。Hadoop+Spark+Django这个组合,刚好把存储、计算、应用三层全串起来,覆盖面广,每个组件都有明确作用,这也是它在课设里被反复选择的原因。
1.2 每层技术选型的具体理由
| 层 | 技术 | 在这个项目里的职责 | 替代方案 | 为什么不选替代方案 |
|---|---|---|---|---|
| 数据采集 | Python + requests/Scrapy | 抓取景点列表、评分、票价、评论数 | 手工Excel整理 | 体现自动化数据获取能力,工作量足够 |
| 数据存储 | Hadoop HDFS | 分布式文件存储,存原始数据和清洗后的数据 | MySQL直接存 | 体现分布式存储概念,撑起“大数据”名头 |
| 离线计算 | Apache Spark | 做景点热度排行、区域统计、票价分布等 | 单纯Pandas | Spark算分布式计算,可演示提交任务流程 |
| 结果存储 | MySQL | 保存Spark分析结果,供Django快速读取 | HBase / Hive | Django访问MySQL最方便,ORM直接映射 |
| Web后端 | Django | 提供JSON接口,渲染大屏页面,管理数据 | Flask / FastAPI | 项目名指定Django,且自带Admin后台方便管理 |
| 可视化 | ECharts | 大屏图表绘制、地图、排行、词云 | Tableau / PowerBI | 网页内嵌展示,与Django天然配合 |
这个表基本就是答辩时老师问“你为什么要这样选型”的标准答案。还有一个隐藏的考量:每个技术点都对应课程大纲里的独立章节,比如HDFS对应分布式文件系统,Spark对应大数据处理,Django对应Python Web开发,一个项目覆盖多门课的知识点,对课程设计的评分是有利的。
1.3 整体架构和数据处理流向
我的项目实际分成四层:
- 数据采集层:Python爬虫抓取重庆景点数据,做初步清洗,输出CSV/JSON文件。
- 存储层:把CSV文件上传到HDFS的指定目录,保留原始数据。
- 分析层:Spark读取HDFS数据,用DataFrame做多维度统计分析,结果写回MySQL。
- 展示层:Django后端从MySQL读取统计结果,通过API输出JSON,前端用ECharts绘制可视化大屏。
这个流程里最容易让人糊涂的是“为什么Spark分析完不直接给前端,还要经MySQL转一手”。原因很简单:Django的ORM原生支持MySQL,但直接在Django里连HDFS并执行Spark任务非常别扭,而且每次打开大屏都现算一遍很浪费资源。把分析结果写入MySQL,相当于一次性算好,之后前端只是查表,响应速度快,代码也简单。我在实际项目里就是让Spark负责“算完”,Django负责“取数和展示”,分工明确。
2. 重庆旅游数据从哪来:定向爬虫与数据规整
2.1 数据源和字段设计
数据源我选了国内几家主流OTA平台的重庆景点列表页,包括景点名称、所在区域、景点类型、评分、票价、热度等公开信息。这里必须强调一句:爬虫只用于课程设计和个人学习,务必控制请求频率,不要影响目标网站正常访问,也要遵守网站的robots协议。我做的处理是单线程、每请求间隔2到5秒随机延时,总共采集几百条记录,量不大,不会对服务器造成压力。
字段设计上,我最终保留了这么几项:
- spot_name:景点名称
- region:所属区县,比如渝中区、沙坪坝区、南岸区、武隆区
- spot_type:景点类型,比如自然风光、历史人文、主题乐园、红色旅游
- score:评分,保留一位小数
- price:参考票价,纯数字,没有门票的填0
- comment_cnt:评论数量,用于代表热度
- hot_desc:热度描述,比如“5A景区”“网红打卡地”这类标签
这个字段设计是经过思考的,不是随手写的。score用于评分排行,price用于票价区间分析,comment_cnt用于热度排名,region用于地图分布展示,spot_type用于类型占比分析。每一个可视化图表都能对应上一列数据。
2.2 页面解析与爬虫骨架
现在的OTA平台大多走JSON接口,直接HTML解析反而麻烦。我先在浏览器开发者工具里找到返回景点列表的XHR请求,复制URL和必要的请求头,然后用requests直接请求JSON。这个项目里我用了最朴素的写法,没上Scrapy,因为数据量小,requests完全够用。
import requests import json import time import random import csv headers = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36" } def fetch_spot_list(page): url = "https://example.com/chongqing/spot/list" params = { "page": page, "pageSize": 20, "city": "chongqing" } resp = requests.get(url, params=params, headers=headers, timeout=10) data = resp.json() return data.get("data", {}).get("list", []) def parse_spot(item): return { "spot_name": item.get("name", "").strip(), "region": item.get("district", "").strip(), "spot_type": item.get("category", "").strip(), "score": float(item.get("score", 0)), "price": parse_price(item.get("price", "0")), "comment_cnt": int(item.get("commentCount", 0)), "hot_desc": item.get("hotTag", "") }需要提醒一点:有些平台对时间戳和Cookie做了强校验,直接requests拿不到数据。这时候我建议先用Selenium模拟浏览器操作,把页面渲染完成后从DOM里取数,或者直接查看页面里嵌入的window初始数据。我在这个项目里两种方式都试过,requests能拿到的优先用requests,拿不到的再用Selenium兜底,毕竟后者更耗资源。
2.3 清洗规则和文件落地
爬下来的数据不可能拿来直接用,清洗是不可避免的。我遇到的杂数据主要有三类:价格字段是“¥80起”“免费”“暂无报价”这种字符串;区域字段有的是“渝中”,有的是“重庆市渝中区”,格式不统一;部分景点的评论数为空。我写了统一的清洗函数处理:
def parse_price(text): if not text or text in ("暂无报价",): return 0 text = text.replace("¥", "").replace("元", "").strip() if text == "免费": return 0 if "起" in text: return float(text.replace("起", "")) return float(text)区域字段我用映射表统一成“渝中区”“沙坪坝区”“南岸区”“武隆区”这种标准格式,避免同一个区出现两种写法导致Spark聚合时被分成两条。
清洗完的数据我先落了一份CSV到本地,编码必须用UTF-8-sig,同时加header。这里有个实用经验:UTF-8-sig带BOM,用Excel打开不乱码,Spark读的时候用UTF-8指定编码也能正常解析,两边都照顾到。之后用hdfs dfs命令把CSV文件传到HDFS的/data/chongqing_spot/目录下,原始数据归档保存,Django端不再直接读这个目录。
3. 伪分布式环境搭建和Spark分析链路,顺手整理的最稳配置
3.1 机器环境与版本搭配
如果你的电脑配置一般,不建议用三台虚拟机搭真集群,伪分布式模式完全够用。伪分布式不等于“假的分布式”,它只是在一台机器上同时跑NameNode、DataNode、ResourceManager、NodeManager等进程,让每个组件的配置逻辑和集群完全一致。真正要换到多节点集群时,只需要把配置复制过去、改一下主机名和IP就行。
我这边的版本组合是:
- 操作系统:Ubuntu 22.04 虚拟机,分配4核8G内存
- JDK:OpenJDK 1.8
- Hadoop:3.3.6
- Spark:3.5.x,用的是spark-3.5.0-bin-hadoop3
- Python:3.8
- MySQL:8.0
这里最容易翻车的是Hadoop版本和Spark版本不匹配。Spark的预编译包会有bin-hadoop3或bin-hadoop2.7这种标记,一定要选择bin-hadoop3的版本,否则Spark读写HDFS时会报依赖冲突。JDK方面,Hadoop 3.x官方已经支持Java 8以上,但如果你要用JDK 11,记得检查Hadoop文档里的兼容矩阵,我图省事直接用的8,没遇到问题。
3.2 Hadoop伪分布式配置要点
Hadoop安装完成后,需要修改几个核心配置文件。第一个是core-site.xml,主要配置NameNode地址:
<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/home/user/hadoop_tmp</value> </property> </configuration>hadoop.tmp.dir务必改成你自定义的目录,不要用默认的/tmp,否则Linux系统重启后临时文件被清理,会导致NameNode元数据丢失,只能重新格式化,这是一个非常常见的坑。
第二个是hdfs-site.xml,设置副本数和NameNode元数据目录:
<configuration> <property> <name>dfs.replication</name> <value>1</value> </property> <property> <name>dfs.namenode.name.dir</name> <value>/home/user/hadoop_tmp/name</value> </property> <property> <name>dfs.datanode.data.dir</name> <value>/home/user/hadoop_tmp/data</value> </property> </configuration>由于是伪分布式,副本数必须设成1,否则DataNode会一直尝试复制第二份数据到其他节点,日志里持续出现块副本不足的告警。配置文件改完之后,需要做三件事:生成SSH免密登录、格式化NameNode、启动进程。
ssh-keygen -t rsa -P '' -f ~/.ssh/id_rsa cat ~/.ssh/id_rsa.pub >> ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys hdfs namenode -format start-dfs.sh start-yarn.sh验证方式很简单,浏览器打开http://localhost:9870,能看到NameNode界面就说明HDFS启动成功。很多教程会让你访问50070端口,那是Hadoop 2.x的端口,Hadoop 3.x已经改为9870,这个不同点容易让人误判服务没启动。
3.3 Hadoop和Zookeeper到底要不要整合
项目相关的热搜词里有“Hadoop和Zookeeper整合实战”,这里顺便说清楚。Hadoop的ZooKeeper主要用于High Availability,也就是HDFS NameNode的自动故障切换。如果你只做伪分布式单节点,不做NameNode主备,ZooKeeper并不是必须的。但很多课设要求体现高可用概念,所以我在这套系统里把ZooKeeper的配置也带上,但没有启用HDFS HA,只在文档里说明了HA的原理和整合步骤。
如果你确实想在单机上做HDFS HA演练,大概流程是:部署三个ZooKeeper节点模拟集群,再配置两个NameNode节点,通过journalnode同步元数据。这个配置过程比较繁琐,对于“重庆旅游景点数据分析”这个核心任务来说并不是必要的。我更建议在答辩时用语言解释清楚“如果数据量继续增长,如何扩展为HA集群”,而不是真的在演示环境里搭三套ZK。
3.4 Spark读取HDFS数据和任务提交
Spark这边我用的local模式,直接通过spark-submit提交Python脚本。分析脚本放在本地,用spark-submit运行时会自动打包上传到Driver,不需要手动指定依赖,这是最省事的方式。
spark-submit \ --master local[4] \ --driver-memory 2g \ --executor-memory 2g \ /home/user/spark_jobs/spot_analysis.py脚本里用SparkSession读取HDFS上的CSV文件,关键点是设置编码和header:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("ChongqingSpotAnalysis") \ .getOrCreate() df = spark.read \ .option("header", True) \ .option("encoding", "UTF-8") \ .csv("hdfs://localhost:9000/data/chongqing_spot/spot_raw.csv") df.printSchema() df.show(5)这时候如果打印中文列名不乱码,说明编码没问题。如果显示为乱码,大概率是CSV编码和option配置不一致,把CSV转成UTF-8-sig后重新上传一次即可。跟HDFS打交道的整个链路里,最容易出问题的永远是编码,后面我会专门讲。
4. Spark在景点数据上到底分析出了什么:指标设计和代码骨架
4.1 分析指标的确定思路
数据分析不能漫无目的地乱跑,每个指标都要对应一个可视化图表的“叙事需求”。我们这个系统的大屏要讲清楚重庆旅游的几件事:哪些区县景点最多、哪些景点最受欢迎、票价分布什么样、高分景点有哪些。于是指标自然拆成五类:
- 各区域景点数量分布,用于地图展示。
- 各类型景点的平均评分和评论总量,用于类型热度分析。
- 票价区间分布,把票价分为0元、0到50元、50到100元、100元以上四档。
- 评论数TOP10景点,作为热度排行榜。
- 高分景点TOP10,作为推荐榜单。
这里有一个设计上的小心思:评论总数代表“热度”,平均分代表“口碑”,热度不一定等于口碑,把两者分开看才有分析价值。比如某个景点评论数很高但评分只有3.8,说明它是流量型景点,体验可能存在争议;另一个景点评论数不高但评分4.9,说明是小众优质景点。这种对比在答辩中很容易结合图表展开,显得思路清晰。
4.2 Spark核心分析代码
整个分析过程在Spark里非常直白,就是groupBy、agg、orderBy的组合。先看区域统计:
from pyspark.sql import functions as F region_stats = df.groupBy("region").agg( F.count("spot_name").alias("spot_count"), F.round(F.avg("score"), 2).alias("avg_score"), F.sum("comment_cnt").alias("total_comment") ).orderBy(F.desc("spot_count")) region_stats.show()类型统计和热度排行同理,唯一需要注意的是sum和avg这类聚合函数在处理空值时会自动跳过,但如果你用Python的lambda自定义UDF,就得手动处理空值。所以我在清洗阶段就确保price、score、comment_cnt这些数值型字段都不为null,实在无法填写的填0,这样后面聚合就不会出现诡异数据。
票价区间我用了when+otherwise生成一个bucket列再分组:
price_bucket = df.withColumn( "price_level", F.when(F.col("price") == 0, "免费") .when(F.col("price") <= 50, "50元及以下") .when(F.col("price") <= 100, "50-100元") .otherwise("100元以上") ) price_stats = price_bucket.groupBy("price_level").count()这个代码很基础,但它是整个项目里最有“讲头”的部分。答辩时可以解释:为什么要用Spark而不是直接用Pandas?答案不是我为了炫技,而是因为当数据量到达百万级以上时,Pandas单机内存会撑不住,Spark可以通过分区并行处理。本项目数据量小,跑Spark和跑Pandas时间差不多,但处理思路是分布式的,这正是课程设计需要的知识点。
4.3 分析结果写回MySQL
Spark算完结果后,我选择写回MySQL。这里不直接用PySpark的jdbc连接也行,简单做法是把分析结果转成Pandas DataFrame,再用SQLAlchemy写库。虽然多了一步转换,但代码读起来直观,依赖也少。
region_stats_pd = region_stats.toPandas() region_stats_pd.to_sql( name="tb_region_stats", con=mysql_engine, if_exists="replace", index=False )tb_region_stats表结构在Django里对应一个模型,字段就是region、spot_count、avg_score、total_comment。其他几个指标也按同样方式写各自的表。这样Django端不需要连Spark环境,只查MySQL就能拿到所有大屏数据。
这里要特别提醒:toPandas()会把全部数据拉到Driver内存,如果数据量巨大,Driver内存可能爆炸。但本项目数据只有几百条,放心用没毛病。如果你真的要处理上亿数据,应该用DataFrame.write.jdbc直接写到MySQL,不要经过toPandas。
5. Django后端与大屏对接的关键细节
5.1 Django项目整体结构
创建好的Django项目结构如下,app命名为spot_data:
chongqing_spot_project/ ├── manage.py ├── config/ │ ├── settings.py │ └── urls.py └── spot_data/ ├── models.py ├── views.py ├── urls.py └── migrations/models.py里我建了三个核心模型,对应Spark写回的三张表。比如区域统计模型:
from django.db import models class RegionStats(models.Model): region = models.CharField(max_length=50, verbose_name="区域") spot_count = models.IntegerField(default=0, verbose_name="景点数量") avg_score = models.FloatField(default=0, verbose_name="平均评分") total_comment = models.IntegerField(default=0, verbose_name="评论总数") class Meta: db_table = "tb_region_stats" verbose_name = "区域统计"需要注意:Django默认会自动添加id主键,而Spark写库时没有id字段,如果Django模型不声明主键,查询时会报字段不存在。我的解决办法是在模型里自定义主键字段,或者把Spark写入SQL时drop掉id字段。实际用下来,直接在模型里声明primary_key=True最省事。
5.2 API接口设计与JSON返回
大屏需要的数据不是一个接口能搞定的,我设计了几个独立API,这样前端受损面小,单个接口出问题也不影响整屏:
- /api/overview/:总景点数、总评论数、平均评分、最高分景点
- /api/region/:各区域景点数量、平均评分,供地图和柱状图
- /api/type_rank/:各类型景点数和评论热度
- /api/price_dist/:票价分档统计
- /api/top10/:评论数TOP10榜单
view层非常简单,用Django的ORM查询后直接返回JsonResponse:
from django.http import JsonResponse from .models import RegionStats def region_api(request): rows = RegionStats.objects.all().values( "region", "spot_count", "avg_score", "total_comment" ) return JsonResponse({"data": list(rows)}, safe=False)这一步几乎没有任何技术难度,但有一个细节很容易踩:JsonResponse默认无法序列化Decimal字段,如果数据库字段用了Decimal类型,直接list(rows)会报Object of type Decimal is not JSON serializable。解决办法是设置json_dumps_params={"ensure_ascii": False},并且在模型里把数值字段定义为FloatField,或者用自定义的JOSNEncoder处理。我在项目里直接把avg_score定义成FloatField,一了百了。
5.3 ECharts大屏渲染和常见问题
大屏页面我放在Django的templates里,用原生HTML+JS+ECharts CDN实现,没上前端脚手架。页面布局采用常见的上下结构:顶部是全屏标题,中间左侧是重庆市地图,中间是热门景点TOP10,右侧是景点类型占比和票价分布,底部再放一个区域数量柱状图。
所有图表都是先调用Django接口获取数据,然后根据数据动态setOption。比如地图部分,我调用了echarts的重庆地图GeoJSON,这个地图数据可以从网络上获取,也可以手动注册到echarts.registerMap里。代码核心是:
fetch("/api/region/") .then(response => response.json()) .then(res => { let mapData = res.data.map(item => ({ name: item.region, value: item.spot_count })); myChart.setOption({ series: [{ type: "map", map: "chongqing", roam: true, data: mapData }] }); });这里有几个在演示前必须自测的点:
- ECharts的CDN地址在离线环境中不可用,如果要现场演示,最好把echarts.min.js下载到本地static目录。
- 地图GeoJSON要提前导成重庆各区县的名称,且要和Spark里的region字段完全一致,比如“渝中区”不能写成“渝中”,否则地图显示不出数据。
- 大屏尺寸要适配不同分辨率,我在外层容器用了百分比布局,图表初始化时监听window.resize,否则全屏和窗口缩放时图表会变形。
5.4 CORS跨域和后台美化
如果前后端分离开发,Django跑在8000端口,前端页面跑在另一个端口,跨域请求是绕不开的问题。虽然我把页面直接放进Django模板里,跨域问题不存在,但在调试阶段如果单独打开过前端静态页面,就会遇到CORS报错。解决办法是安装django-cors-headers:
pip install django-cors-headers然后在settings.py里添加:
INSTALLED_APPS = [ "corsheaders", ... ] MIDDLEWARE = [ "corsheaders.middleware.CorsMiddleware", ... ] CORS_ALLOW_ALL_ORIGINS = True只在开发调试阶段这么开,实际部署时应该限定具体允许的域名。另外给Django Admin换个样式,让管理后台看起来更专业,可以装一个django-unfold,但它对Django版本有要求,版本不匹配会直接报错,所以要不要加全看个人需求。
6. 我在这个项目里踩过的六个坑,每一个都值得录进调试文档
6.1 中文编码问题:几乎贯穿全流程
这是我在项目里遇到最多、也最隐蔽的问题。爬虫输出的CSV用Excel打开正常,但Spark读取后中文全部变成乱码。排查链路是:先用hdfs dfs -cat直接查看HDFS上的文件,发现文件内容本身正常,那问题一定出在Spark读取参数上。
最终确认是CSV保存时用了GBK编码,Spark读取时默认UTF-8。解决方式是把CSV统一用UTF-8-sig编码重新保存,同时在Spark读取时显式指定编码:
df = spark.read.option("encoding", "UTF-8").csv(...)这个坑一旦踩过,之后遇到任何中文乱码,我都会第一时间怀疑“文件编码和读取编码不一致”。还有一次是Django返回的JSON在浏览器里中文正常,但在ECharts图表里变成“\uXXXX”,那是JsonResponse的ensure_ascii参数未设置导致的,改成False就正常了。
6.2 NameNode格式化两次,HDFS状态不一致
我第一次启动Hadoop时,格式化完NameNode,start-dfs.sh后启动失败,日志提示namenode目录不可用。排查发现是格式化过一次又重启后没有清理临时目录,导致NameNode在启动时读取到损坏的元数据文件。
解决方法是直接停掉所有Hadoop进程,删除自定义的hadoop_tmp目录,重新创建,然后再次执行hdfs namenode -format。这个操作会清空HDFS上所有数据,所以格式化前必须确认没有重要数据,或者把原始数据备份在本地。做完项目我养成的习惯是:先把CSV文件放在本地一份,HDFS清掉也不慌,再传一次就行。
6.3 Spark写MySQL时驱动类加载失败
Spark通过JDBC直连MySQL时经常看到ClassNotFoundException: com.mysql.cj.jdbc.Driver,原因是没有把MySQL驱动jar包放到Spark的jars目录,或者用spark-submit提交时没有用--packages指定驱动。我在项目里的笨办法是下载mysql-connector-java.jar直接放进$SPARK_HOME/jars目录下,重启SparkSession就好。如果你不想污染Spark安装目录,也可以用:
spark-submit --packages mysql:mysql-connector-java:8.0.33 ...但这要求环境能联网下载maven依赖,如果演示现场没有外网,还是提前把jar放好最稳妥。
6.4 前端图表初始化时容器宽高为0
大屏页面加载时,图表总是显示不出来或者只显示一小块,非要手动resize才正常。原因是ECharts在DOM未完全渲染时就被调用,容器宽度还是0。我加了window.onload保证页面加载完成后再init,同时对地图和图表初始化做了debounce处理:
window.addEventListener("load", () => { initCharts(); });如果你用Vue或React,还要注意组件生命周期,图表init应放在mounted钩子里而不是created里。这个坑在纯Django模板里相对好避,因为模板渲染是同步的。
6.5 大屏数据不实时更新,演示时掉链子
第一次演示时,我发现Spark算完新数据,Django大屏还是显示旧数据。原因是我把Spark结果写进了MySQL的同一张表,if_exists="replace"会先删表再建表,但Django进程里ORM可能存在缓存,或者前端浏览器缓存了旧JSON。
解决方法是给API响应加上版本号参数,前端请求时带一个时间戳,例如/api/region/?v=20250601,后端不处理这个参数也能绕过浏览器缓存。同时每次Spark计算完,在Django端重启一次服务进程,或者干脆用Redis缓存并设置过期时间。我在文档里建议的是最朴素的方案:每次重新导数据后,重启Django进程,简单可靠,不引入额外依赖。
6.6 伪分布式下Spark任务提交时间过长
Spark任务本身几秒就完成,但spark-submit提交过程经常要等十几秒甚至更久,每次调试都让人以为卡死了。原因是Driver和Executor通信在local模式下也会有大量网络重试和日志输出,尤其是Hadoop和Spark都在同一台机器并启用了IPv6时,绑定localhost会失败。
解决方式是统一使用127.0.0.1而不是localhost,或者在/etc/hosts里把主机名映射到127.0.0.1,避免Spark尝试通过真实主机名连接外网。还有一个技巧:提交任务时加上--conf spark.driver.host=127.0.0.1,明显缩短启动时间。
7. 源码运行、文档编写和答辩演示的经验
很多课设项目的源码本身不难,难的是让评审老师快速看懂你的东西。我在交付项目时整理了一份简要文档,包含需求分析、系统设计、数据库设计、核心代码说明、测试运行步骤、问题与调试记录六个部分。实测下来,调试记录是最受欢迎的部分,因为它能体现你真实做过这个项目,而不是从网上随便下载的。
文档里我特意写清楚了“从零复现步骤”,包括:
- Python爬虫运行,生成spot_raw.csv;
- 上传CSV到HDFS,命令为:hdfs dfs -mkdir -p /data/chongqing_spot && hdfs dfs -put spot_raw.csv /data/chongqing_spot/;
- 运行Spark分析脚本,结果自动写MySQL;
- 启动Django服务,浏览器访问大屏地址。
这个顺序就是系统的完整演示链路,也是答辩时的口头讲述顺序。很多同学演示时东点一下西点一下,评委根本看不出逻辑;而按这个顺序走,从数据采集到存储、计算、展示,每一步都有明确的输入输出,评委自然能跟上你的思路。
答辩时我被问过几个高频问题,在这里把参考答案也整理出来:
- “为什么用Hadoop存储,而不是直接用MySQL?”回答要点:HDFS体现分布式存储理念,适合海量原始数据的持久化,分析结果才进MySQL,属于分层设计。
- “数据量这么小,用Spark有必要吗?”回答要点:Spark的价值在于分布式计算能力,本项目数据量小,但重点是跑通流程;换成更大数据量只需求改集群规模。
- “热点数据是如何定义出来的?”回答要点:没有单一指标,而是综合评论数、评分、价格档位和景区等级,用加权排序得到热度榜单。
- “可视化大屏的数据多久刷一次?”回答要点:基于离线计算,每次Spark任务完成后再刷新;如果做实时更新,可以引入Kafka和Spark Streaming。
如果你最后还有余力,建议在系统里加两个小功能提升完成度:一个是用Spark读取JSON格式的数据文件,另一个是在Django后台用django-unfold美化认领数据管理界面。这两个功能都能在答辩时多讲几句,而且对代码量贡献不大。
最后再分享一个我实操下来的体会:这类“Hadoop+Spark+Django+可视化大屏”的项目,放在一起最容易卡住的不是技术深度,而是组件之间的“暗链接”。环境版本、编码格式、数据格式、端口配置,每个环节错一点,整条链路就断掉。做项目的过程中别急着往下推进,先花半小时把环境验证扎实,后面反而最省时间。按照上面这套链路走一遍,把坑提前踩平,你的系统就能稳稳跑在演示现场。