Python+Spark电影数据分析系统:从爬虫到可视化大屏的完整实现
2026/9/8 6:17:03 网站建设 项目流程

你有没有在答辩现场被评委问倒过?我之前帮一个学弟看他的“基于Python的Spark电影数据分析系统”,功能铺得挺全:requests爬虫抓豆瓣电影、Spark做分析、Flask写接口、ECharts出可视化大屏,整个链路都有。结果评委只问了一句:“你的数据量多大?为什么非要用Spark?”他愣了几秒,才开始讲“因为是大数据项目”。这个场景其实特别典型——demo能跑通,但底层的架构逻辑没有想透。

这个标题看起来很“毕业设计模板”,但拆开看,它本质上是一套完整的数据工程闭环:数据采集 → 数据清洗 → 分布式分析计算 → Web接口 → 可视化展示。你把它做扎实了,简历上能写的东西远不止“会用框架”那么简单。这篇我打算按真实的开发顺序,把每一层的选型原因、核心代码、踩坑记录都过一遍,尤其是我个人觉得最容易被查得细的Spark资源问题和爬虫限流问题,尽量讲透,方便你直接照着复现。

1. 别被标题骗了:这其实是一条完整的数据工程流水线

1.1 这套系统解决的真实问题

很多人一看到“大数据”三个字,就以为非得搭一个几十台机器的集群。实际上在毕业设计范畴里,这套系统解决的是这样几个问题。

第一,数据从哪来。豆瓣电影是中文互联网里结构化程度比较高的电影数据库,一部电影有评分、评价人数、类型、地区、年代、导演、演员等多个维度字段,非常适合作数据分析的原料。用requests去抓取榜单接口,能拿到半结构化的JSON数据,这一步实际上练的是“把网页数据变成可计算数据”的能力。

第二,数据怎么存、怎么洗。抓下来的原始数据不能直接用,评分可能缺省、类型可能是“剧情/爱情/冒险”这种拼接字符串、上映日期里混着年份和地区,这些都必须做统一的清洗和标准化,才适合交给计算引擎。

第三,怎么体现“大数据计算”。如果数据量只有几千条,用pandas处理完全够,但毕业设计里你必须展示对Spark的理解——RDD/DataFrame、Spark SQL、行动算子/转化算子、提交到YARN集群上的资源调度。所以这个项目里Spark承担的定位可以设计成:对清洗后的千万级电影记录做离线聚合分析,产出各种指标结果。

第四,结果怎么展示。Spark算完推出一堆CSV或表,但非技术用户看不懂。所以后端用Flask把聚合结果封装成JSON接口,前端用ECharts渲染成折线图、柱状图、饼图,这才是“可视化”的实际含义。

1.2 技术栈组合逻辑:为什么是Python+Spark+Flask

这套组合在技术选型上是非常典型的“各管一段”:

  • requests爬虫:Python生态里最轻量的HTTP库,写爬虫学习和维护成本低,适合快速落地。
  • Spark:负责核心的分布式离线计算。虽然这个项目里你可以跑在Local模式,但代码层面和提交到YARN集群的写法完全一致,能力边界在。
  • Flask:轻量级Web框架,写几个API接口很快,不需要像Django那样引入整套ORM和Admin体系,对一个以数据分析为主的项目来说刚刚好。
  • ECharts:纯前端的图表库,数据驱动,接口给什么它画什么,和后端解耦。

这套链路还有一个隐性好处:每一层都是独立可替换的。爬虫抓的数据可以换成别的来源,Spark的计算任务可以被调度系统定时触发,Flask接口的响应结果可以被任意前端消费。这就是“工程化”和“写了一堆脚本”的区别。

1.3 项目目录结构与数据流向

我建议你按下面这种目录组织代码,层次清晰,答辩也好讲:

movie-analysis/ ├── crawler/ │ ├── douban_spider.py # requests爬虫 │ └── movies.csv # 原始抓取结果 ├── clean/ │ └── clean_movies.py # Spark DataFrame清洗 ├── analysis/ │ ├── stats_by_year.py # 年份维度聚合 │ ├── stats_by_genre.py # 类型维度聚合 │ └── stats_top_movies.py # 高评价/热门排行 ├── backend/ │ ├── app.py # Flask主入口 │ └── db_helper.py # MySQL查询封装 ├── frontend/ │ ├── index.html │ └── dashboard.js └── submit.sh # spark-submit脚本

数据流向:豆瓣JSON → requests抓取 → CSV → Spark清洗 → 聚合结果表 → MySQL → Flask JSON接口 → ECharts可视化。这个链路讲清楚,评委的“你这个项目的整体架构是什么”就稳了。

2. 第一道关卡:用requests把豆瓣电影数据变成“原料”

2.1 抓取目标与接口分析

爬虫部分不用去碰豆瓣的静态页面,那种HTML解析麻烦且容易遇到验证码。比较常规的做法是直接抓豆瓣电影的榜单JSON接口。这类接口返回结构化数据,字段干净,省去了大量解析工作,非常适合作为学习requests的实战案例。

抓取时有两个重点。一个是翻页参数,榜单接口通常靠startlimit控制偏移量,start=0&limit=20表示第一页,start=20表示第二页,以此类推;另一个是类型参数,比如你想抓“剧情”类电影,就传对应的类型编号。把这些参数封装成字典,通过params传参,代码会干净很多。

import requests import time import random HEADERS = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36", "Referer": "https://movie.douban.com/" } def fetch_movies(start=0, limit=20): url = "https://movie.douban.com/j/chart/top_list" params = { "type": "11", # 示例类型,实际以目标页面为准 "interval_id": "100:90", "action": "", "start": start, "limit": limit, } resp = requests.get(url, params=params, headers=HEADERS, timeout=10) if resp.status_code == 200: return resp.json() return []

2.2 请求头、限速与429限流的工程化处理

初学爬虫最容易犯的错是“一次性疯狂请求”。豆瓣这类站点对高频访问非常敏感,很容易直接返回HTTP 429 Too Many Requests。你在网上搜索“douban requests 429”会看到一堆人遇到这个错误,其实根因就是请求频率超过限制。

我的处理经验有四点:

  1. 请求头要够“像浏览器”User-AgentReferer是必须的,否则很容易拿到403。有条件可以准备几个UA轮换。
  2. 请求间隔必须随机。固定间隔0.5秒看起来礼貌,但模式太规律一样容易被识别。用time.sleep(random.uniform(2, 5))把时间间隔打散,模拟真人操作。
  3. 对429做指数退避重试。第一次失败后等2秒再试,第二次等4秒,第三次等8秒,最多重试3~4次。这个策略比断了就重跑要高效得多。
  4. 只抓学习需要的量。毕业设计抓几千条到几万条足够展示技术链路了,没必要顶着反爬去全量抓。抓取行为严格限制在个人学习研究用途,控制并发,遵守目标网站的协议约定。

下面这段是带重试的版本,直接可以参考:

def fetch_movies_with_retry(start=0, limit=20, max_retries=3): for attempt in range(max_retries): try: resp = requests.get(url, params=params, headers=HEADERS, timeout=10) if resp.status_code == 200: return resp.json() if resp.status_code in (403, 429): wait = (2 ** attempt) + random.random() print(f"请求频率受限,{wait:.2f}秒后重试...") time.sleep(wait) continue except requests.RequestException: time.sleep(random.uniform(3, 6)) return []

这套“随机间隔 + 指数退避”的思路不止适用于豆瓣,几乎所有HTTP接口的限流问题都是这么解。

2.3 字段设计与去重增量存储

抓回来的JSON字段很多,但不是每个都要存。按分析目标,建议至少保留以下字段:

字段示例说明
movie_id1292052电影唯一标识,去重主键
title肖申克的救赎片名
rating_score9.7豆瓣评分
rating_people2800000评价人数
pub_date1994-09-10(美国)上映日期,后续提取年份
genres剧情/犯罪类型,按斜杠拆分
duration142分钟片长,提取分钟数
regions美国制片国家/地区
directors弗兰克·德拉邦特导演
actors蒂姆·罗宾斯/摩根·弗里曼主演

存储层面,初学者最容易踩的坑是“重复抓取”。解决思路很简单:以movie_id为唯一键,每次抓完先和已有数据比对,存在就跳过。单机做可以先读CSV里的id集合放到内存,数据量大了再换成数据库唯一索引。这一招看着基础,但在答辩时讲出来,说明你考虑了实际工程里的幂等性问题。

3. 像做菜前挑菜一样:Spark里的数据清洗与特征加工

3.1 脏数据长什么样

爬虫抓完的数据一定是不干净的。我拿电影数据给你举几个真实例子。

评分字段看起来很规整,但偶发情况下是空字符串或者暂无评分,直接cast成float会变成null,如果后续聚合用AVG(score),null值不影响其他值,但如果你没过滤,页面展示上就会很难看。评价人数同理,有的接口在热度不足时给0,这个0到底是“0人看过”还是“缺省值”,需要你自己定义清洗策略。

类型字段是“剧情/犯罪/惊悚”这种拼接字符串,不拆开的话,你没法统计“犯罪片平均分高还是喜剧片平均分高”——因为一个电影会被算进所有类型的桶里。上映日期更麻烦,不同条目格式不一样,有的是1994-09-10(美国),有的只有1994,有的是2019-01-01(中国大陆),必须用正则抽年份。

还有一个常见坑是中文字段里的空格和特殊字符。比如导演:张三这种带冒号的,或者末尾带换行符的,清洗阶段要统一trim处理,否则分组统计时同一个导演会被算成两个。

3.2 Spark DataFrame清洗落地

清洗在Spark里做,比的不是“怎么把数据弄干净”,而是“怎么用DataFrame API把清洗逻辑写成可复用的计算流”。我一般按两步走。

第一步,统一类型和抽字段:

from pyspark.sql import SparkSession from pyspark.sql.functions import col, regexp_extract, split, when spark = SparkSession.builder.appName("movie_clean").getOrCreate() df = spark.read.csv("crawler/movies.csv", header=True, inferSchema=True) df_clean = df.select( col("movie_id").cast("long").alias("movie_id"), col("title"), col("rating_score").cast("double").alias("rating_score"), col("rating_people").cast("long").alias("rating_people"), regexp_extract(col("pub_date"), r"(\\d{4})", 1).cast("int").alias("year"), split(col("genres"), "/").alias("genre_arr"), regexp_extract(col("duration"), r"(\\d+)", 1).cast("int").alias("minutes"), col("regions"), col("directors"), col("actors") ).filter(col("rating_score").isNotNull() & col("year").isNotNull())

第二步,处理异常值。比如评价人数小于等于0的,可以做标记但不删除——因为你可能在分析“冷门电影”时还会用到它;比如片长缺失的,可以填充一个方法值-1,或者在画图时过滤掉。关键是清洗策略要能自圆其说,不能拍脑袋。

3.3 清洗之后的数据长什么样

清洗完成后一定要做“数据质量四查”:查总数、查缺失率、查字段类型、查抽样结果。我习惯把清洗后的DataFrame写回一份新的CSV或Parquet,后续Spark分析任务直接读这份干净数据,避免每次都在计算任务里重复清洗逻辑。

如果用的是Parquet格式,还有个额外好处:它是列式存储,带schema,Spark读起来比CSV快不少。在毕业设计阶段你体会不明显,但你在答辩时说出来“我用Parquet做中间存储,减少重复清洗的IO开销”,评委是认的。

清洗层最容易被忽略的是分区策略。如果你的原始数据达到了亿级,按year做分区后,分析任务可以直接做分区裁剪,比如只看2020年之后的数据,Spark不需要扫描全部文件。这个点稍微提一下,就能证明你不只是会调API,而是真的理解分布式存储。

4. Spark分析引擎:为什么要把它放进来,以及YARN资源那点事

4.1 单机Pandas vs 分布式Spark,毕业设计该怎么选

这是一个绕不开的问题。如果数据量只有几千条,坦诚讲,Pandas处理起来比Spark快得多,启动速度也快。但你把它作为“大数据毕业设计”,选Spark有几个不能忽略的理由。

第一,架构能力展示。Spark代码写的是分布式执行计划,数据量上来以后可以直接横向扩容,这一点Pandas做不到。你在项目里用Spark,说明你有“数据量增长时系统架构怎么响应”的意识。

第二,统一SQL与编程API。Spark SQL可以让你像写SQL一样做聚合,底层自动优化执行计划。对评委来说,Spark DataFrame API和SQL的写法都比手写一堆循环有说服力。

第三,生态完整。后面如果扩展推荐系统,可以用Spark MLlib的ALS;如果做实时分析,可以接Spark Streaming。这些扩展点是大数据分析项目最自然的落脚点。

所以我的建议是:思路按分布式写,数据量不用真的大到离谱。你抓个几千上万条做完整链路演示完全没问题,但代码里要体现“如果数据量上亿,这套代码依然能跑”的设计原则。

4.2 核心分析指标与Spark SQL实现

分析指标不要贪多,四到五个维度就能把可视化大屏填得很漂亮。我常做的指标如下。

年度产量与均分趋势:看电影数量和平均分随年份的变化,能直观反映内容市场的冷热。

trend_df = spark.sql(""" SELECT year, COUNT(*) AS movie_cnt, ROUND(AVG(rating_score), 2) AS avg_rating, SUM(rating_people) AS total_votes FROM movie WHERE year BETWEEN 1980 AND 2024 GROUP BY year ORDER BY year """) trend_df.write.mode("overwrite").option("header", True).csv("output/agg_year")

类型分布:把genre_arr用explode展开成一行一个类型,再统计每个类型的电影数量和平均分。这一步会揭示一个有意思的规律:数量多的类型往往是平均水平一般的类型,数量少的小众类型反而平均分偏高。

from pyspark.sql.functions import explode genre_df = df_clean.select(explode(df_clean.genre_arr).alias("genre"), "rating_score") genre_stats = genre_df.groupBy("genre").agg( count("*").alias("movie_cnt"), round(avg("rating_score"), 2).alias("avg_rating") ).orderBy(col("movie_cnt").desc())

评分区间分布:用when把分数切成[0,5)[5,6)[6,7)[7,8)[8,9)[9,10]六个区间,统计数量占比。这种数据最适合做柱状图或南丁格尔玫瑰图。

评分与评价人数的关系:直接看散点图。你会发现“叫好”和“叫座”不是一回事——很多9分以上的电影评价人数远少于8分左右的大众片,这个结论可以写进你的分析报告里,是真实的数据洞察。

4.3 Spark on YARN只用到1个CPU核:一个高频问题的排查链路

这个问题的出现频率高到离谱,随手一搜就是“spark on yarn cpu只能用1个”。现象是:你在spark-submit里配了--executor-cores 4,但YARN上每个Executor实际只分配到1个vCore,或者整个App只申请到1个核

我在帮别人排查时发现,原因集中在下面几个层面,按概率排序:

  1. spark.executor.cores在YARN模式下默认是1。很多人只在命令行里写了--num-executors--executor-memory,以为核数会像Standalone模式那样拿到机器全部CPU,实际上YARN模式下如果要多个核,必须显式指定--executor-cores
  2. YARN调度器单Container上限设置过小yarn.scheduler.maximum-allocation-vcores决定了每个Container最多能拿多少核,如果你的--executor-cores超过这个上限,YARN会直接拒绝或者按上限来。
  3. 节点可用资源不够。Executor要申请4核4G,而节点剩余资源只剩2核2G,调度器就会降级分配。
  4. 提交参数没有生效。比如通过spark-submit传了参数,但spark-defaults.conf里的同名配置优先级覆盖了命令行的设置,或者压根没把参数写对位置。

排查链路我建议这样走:先到YARN的ResourceManager页面里看App运行时的实际Executor资源,再用yarn logs -applicationId查启动日志里打印的Executor resource request,最后对比yarn-site.xml里调度器的上限。八成问题出在配置没真正生效。

4.4 内存模型与executor参数推荐

Spark的内存问题比CPU还隐蔽。简单说,一个Executor进程申请的内存不只有spark.executor.memory这部分,还有一块spark.executor.memoryOverhead,用来跑JVM以外的系统开销。默认情况下overhead大约是executor内存的10%且不小于384MB。也就是说,你申请--executor-memory 4g,实际YARN上占用的容器内存是4g + 至少384m。

对毕业设计的小集群,我用得比较顺手的参数是这样的:

配置项推荐值说明
--num-executors2~3executor数量,不贪多
--executor-cores2~4每个executor的核数
--executor-memory3g~5g留足overhead空间
--driver-memory1g~2gdriver端做collect时不至于爆
spark.default.parallelism4~8参考总核数设置

一个容易踩的坑是内存和核数不匹配。比如给一个Executor配了8核1G内存,每个任务实际能用的内存非常小,一跑groupBy就OOM。一般经验是单个Executor核数尽量不要超过5,核数和内存按“1核配1~2G”的比例比较稳。

5. Flask后端+可视化大屏:把分析结果翻译成人和图表都能懂的语言

5.1 Flask API设计:只做接口不参与计算

到了Flask这一层,核心原则是:后端千万不要写任何分析计算的代码。Spark算好的聚合结果已经落到MySQL或CSV里了,Flask只负责查询和序列化。

我建议用Blueprint把接口拆成模块,至少提供这几个:

from flask import Blueprint, jsonify api = Blueprint("api", __name__) @api.route("/api/movie/year_trend") def year_trend(): # 查MySQL表 agg_year,按年份排序 return jsonify(code=0, data=rows) @api.route("/api/movie/genre_dist") def genre_dist(): return jsonify(code=0, data=rows)

如果你用的是MySQL,注意PyMySQL连接时的字符集参数要显式指定charset="utf8mb4",否则中文和特殊符号容易乱码。每次请求现连数据库虽然慢一点但开发省事,答辩演示完全够用;追求性能可以加一个简单的@lru_cache,但接口如果有动态参数就不要乱用。

5.2 ECharts可视化:数据格式约定与图表案例

前后端联调最怕的是“前端要用A格式,后端返了B格式”。我踩过几次坑之后,现在一律统一成{ code: 0, data: [...] }data是数组,每个元素是一个扁平对象,字段名用驼峰。前端拿到直接map

以年度趋势折线图为例,前端代码大致是这样:

fetch('/api/movie/year_trend') .then(res => res.json()) .then(res => { const years = res.data.map(item => item.year); const avgRatings = res.data.map(item => item.avgRating); myChart.setOption({ tooltip: { trigger: 'axis' }, xAxis: { type: 'category', data: years }, yAxis: { type: 'value' }, series: [{ name: '平均评分', type: 'line', data: avgRatings, smooth: true, areaStyle: {} }] }); });

大屏上我建议至少放五个模块:顶部总览卡片(总影片数、总评价人数、评分最高影片)、年度趋势折线图、类型分布玫瑰图、评分区间柱状图、热门影片Top10条形图。多了信息过载,少了显得单薄,这五个刚好能把一整面屏铺满。

5.3 联调过程中最容易踩的3个坑

第一个坑是跨域问题。如果你的前端页面直接用文件打开,接口请求大概率会被CORS策略拦掉。开发期最简单的解法是Flask端加一个统一的after_request响应头,放行本地请求。

第二个坑是日期字段序列化。Spark写出来的年份是整数没问题,但如果有日期类型,Flask的jsonify不认datetime对象,必须自己转成字符串,否则直接报TypeError

第三个坑是图表数据为空的渲染。某些年份没有电影,接口返回空数组,ECharts的折线图会断掉。最好在后端就补默认值,或者前端对空数组做特殊处理,别让评委在演示时看到一条断掉的线。

6. 从“能跑”到“能答辩”:毕业设计项目的延伸与避坑

6.1 数据量不够大怎么办

这是很多人会慌的问题:“我只有几千条数据,能叫大数据吗?”

我的理解是,毕业设计的关键在于链路完整性和技术选型的合理性,而不是数据规模本身。你可以换这几个思路强化它:

  • 除了榜单页,再抓一批短评数据,做分词和词频统计,瞬间多出一个“文本分析”模块,数据量也能翻几倍。
  • 用Spark的随机扩样功能模拟更大数据量,比如把原始数据复制成10份并加上扰动,演示“数据量增长后Spark依然线性处理”的能力。
  • 把MySQL换成HDFS存储中间结果,展示你懂分布式文件系统的接入。

总之,别让“数据少”变成你的软肋,而是展示你“知道怎么在有限数据上做合理演示”的判断力。

6.2 答辩高频问题与回答思路

评委对你的项目一般会从这三个角度提问。

你的数据量多大?怎么保证抓取合规?”如实说量级(几千到几万条),强调个人学习用途、控制请求频率并遵守目标站点的公开协议。不要含糊,也不要吹自己抓了上亿条。

为什么用Spark,Pandas不行吗?”这是一个展示架构思考的好机会。你可以回答:Spark的DataFrame/SQL天然支持分布式计划,在数据量增长时可以横向扩展;而且Spark的聚合计算优于Pandas迭代处理,还能与后续的MLlib、Structured Streaming无线衔接。重点是让评委感受到你是“选型”而非“跟风”。

系统如果要做实时推荐,怎么改?”这个问题答好了会很加分。你可以说:把离线聚合任务换成Spark Structured Streaming消费Kafka里的用户行为数据,实时计算热门电影TopN;推荐部分用Spark MLlib的ALS协同过滤,用户评分矩阵作为训练数据,产出的TopN推荐结果写入Redis,Flask接口查Redis返回。不需要真的实现,但你能讲清楚链路,就已经胜过了大部分只做增删改查的同学。

6.3 几条实操建议再唠叨一遍

整个项目从零到能演示,我的经验是优先把主链路跑通,再去抠优化。主链路就是“爬数据→洗数据→算指标→出图表”这四个环节,先能用起来,再慢慢换大集群、调参数。

第二个建议是每层都留好日志。爬虫要打印抓了多少条、失败了多少条,Spark任务要打印输入输出路径,Flask接口要打印请求时间。这些日志不光是调试的工具,答辩时你随手展示一下,比嘴上说“我做了日志记录”有说服力得多。

第三个建议是提前准备好README或架构说明文档,包含环境搭建步骤、启动命令、每层输入输出说明。我帮过不少学弟排错,发现他们项目代码没问题,但换一台机器就不知道怎么启动——Spark、MySQL、Flask、前端每个组件都要单独启动,没有文档的话,你自己两周后再看都会蒙。

最后再说一个针对“建议收藏”这句话的体会。一个毕业设计项目值不值得被收藏,不在于它用了多少框架,而在于它有没有完整地展示出你理解数据从哪来、怎么算、怎么用。把这篇文章里提到的链路和坑都过一遍,你的项目就不只是一个能交差的作业,而是可以写进简历、拿到面试里讲半天的作品。

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

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

立即咨询