简介:这是一套基于 Python 与 Spark 的豆瓣电影爬虫和数据分析可视化系统,定位为高分毕业设计参考项目,也适用于期末大作业与课程设计。项目从豆瓣电影页面抓取数据,经清洗整理后存入数据库,再利用 Spark 完成词频、评分等级、评论数量、年份分布等统计,最终以图表形式呈现。资源共包含 241 个文件,压缩包大小约 5.65 MB,文件类型涉及 XML 配置、Java 辅助类、前端样式与脚本、Python 主程序、SQL 数据库脚本以及 Spark 输出结果文件,能够覆盖爬虫采集、数据清洗、分析计算、可视化展示的完整流程。代码附有详细注释,结构清晰,新手也能理解关键逻辑,下载后简单部署即可运行。目前已有 248 人学习下载,适合需要从零构建电影数据分析项目的同学直接作为模板,可大幅节省开发与调试时间。
1. 这套豆瓣电影系统不只是「爬虫+图表」的堆叠
豆瓣电影 Top 250 只有 10 页、250 条记录,但把它做成一个「Python 爬虫 + Spark 数据分析 + 可视化」完整系统时,大多数人卡住的不是爬虫,而是数据跑到一半没法复现。高分毕设或工程原型的判分点,从来不是抓了多少数据,而是网页到图表之间每一层是否可解释、可重算、可验证。我见过太多把 pandas 脚本硬写成 Spark 作业的代码,忽略了 JDBC 分区读取、结果回写和重复运行时的幂等性。这篇文章会按最常用的落地结构拆解:requests 并发抓豆瓣页面,MySQL 落库,Spark 做 ETL 和指标计算,再用 Flask 与 ECharts 完成可视化大屏。适合准备课程设计、毕业设计的学生,也适合想快速搭端到端数据管线的工程师。
2. 分层设计:Python 采集、Spark 分析、MySQL 存储
2.1 为什么采集层用 Python、分析层换 Spark
很多人第一反应是:250 条数据用 pandas 一分钟就算完,为什么还要引入 Spark?这不是性能问题,而是架构表达问题。采集层的瓶颈在网络 I/O 和反爬策略,Python 的 requests、BeautifulSoup 生态最顺手,天然适合做抓取和解析。分析层则不同,Spark DataFrame 的延迟计算、分区任务、失败重试、跨数据源能力,是可以放进答辩 PPT 的技术深度。
选型上有一个明确分工:采集层不碰 Spark,因为为 10 个页面启动一个 SparkSession 是浪费;分析层不碰 requests,因为指标计算需要的是确定性的批处理能力。如果抓的是全站几百万条短评,单机 pandas 很容易内存溢出,而 Spark 往集群一扔,代码不用大改。本地 Spark 对 250 条数据确实比 pandas 慢,但换来的是未来扩展时不需要推翻业务代码。
2.2 数据流:从 HTML 到可视化 JSON
常见做法是把流程拆成五个层次,每层只依赖上一层的产物,这样做的好处是每一层都可以单独重跑、单独验证。
| 层级 | 输入 | 主要工具 | 输出 |
|---|---|---|---|
| 采集层 | 豆瓣 Top 250 页面 | requests + BeautifulSoup | 原始记录 |
| 存储层 | 解析后的字段 | MySQL | movie_facts 表 |
| 计算层 | movie_facts 表 | Spark JDBC + Spark SQL | report_* 结果表 |
| 服务层 | report_* 结果表 | Flask | JSON API |
| 展示层 | JSON API | ECharts | 可视化大屏 |
这个流程里最容易做错的是直接用爬虫结果怼进 Spark,中间没有落库。不落库意味着爬虫失败一次,数据血缘就断了。我在实际项目里会坚持先入库,再通过 Spark JDBC 读取,因为 MySQL 既能当数据源,又能做断点续爬的去重依据。
2.3 先建这两张表:movie_facts 与 crawl_log
存储层建议只建两张表,一张存电影事实数据,一张存爬虫运行日志。事实表采用宽表冗余设计,把导演、演员、类型直接放在同一行,避免毕设里做复杂的关联查询。
CREATE DATABASE IF NOT EXISTS douban_movie DEFAULT CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci; USE douban_movie; CREATE TABLE movie_facts ( id BIGINT PRIMARY KEY AUTO_INCREMENT, movie_id VARCHAR(32) NOT NULL, title VARCHAR(255) NOT NULL, directors VARCHAR(500), actors VARCHAR(1000), year INT, country VARCHAR(255), genre VARCHAR(500), rating DECIMAL(3,1), votes INT, quote VARCHAR(255), UNIQUE KEY uk_movie_id (movie_id) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; CREATE TABLE crawl_log ( id BIGINT PRIMARY KEY AUTO_INCREMENT, page_no INT NOT NULL, start_time DATETIME, end_time DATETIME, fetched_rows INT, status VARCHAR(20) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;rating用DECIMAL(3,1)而不是FLOAT,因为评分只需要一位小数,精确小数可以避免 Spark 读取后出现 9.1999999 这类误差。movie_id加唯一索引,是写幂等插入的关键。quote字段可能包含引号或特殊符号,所以数据库字符集使用utf8mb4,兼容 emoji 和生僻字。
crawl_log表不存业务数据,只记录每一页爬取开始时间、结束时间、获取条数和状态。它的作用是回答答辩时最常见的问题:“如果爬虫跑到一半断了,你怎么知道从哪继续?” 有了这张表,就能直接统计上次成功页码和状态,而不是靠日志文件翻找。
3. 豆瓣电影爬虫的并发采集与反爬控制细节
3.1 只抓 Top 250 的三类字段,别上来就做全站爬虫
豆瓣全站爬虫涉及用户登录态、评论分页、鉴权等多个复杂环节,作为毕设没有必要,还会因为请求量过大导致 IP 被封。所以最稳妥的抓取范围是 Top 250,10 页、250 条,足够支撑评分分布、年代趋势、类型占比这些典型分析指标。
字段选择也有讲究。我在设计时坚持只抓三类:电影标识与名称,展示类的导演、演员、年份、国家、类型,以及可计算的评分和评价人数。quote属于可选加分项,用来做排行榜的展示辅助,抓不到就填空字符串,不影响分析。豆瓣网页结构相对稳定,但字段缺失时要能容忍,不要一遇到None就让整个爬虫崩溃。
3.2 requests + ThreadPoolExecutor 并发采集的最小实现
这是典型的 I/O 密集型任务,用多线程比多进程更合适。下面这个实现是常见的爬虫骨架,重点不在代码量,而在限速、重试、去重三个细节。
import random import threading import time import requests from bs4 import BeautifulSoup from concurrent.futures import ThreadPoolExecutor, as_completed HEADER_POOL = [ {"User-Agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36"}, {"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36"}, ] DELAY_RANGE = (1.2, 2.8) lock = threading.Lock() SEEN_URLS = set() def fetch_page(page): url = f"https://movie.douban.com/top250?start={page * 25}" headers = HEADER_POOL[page % len(HEADER_POOL)] for attempt in range(3): try: resp = requests.get(url, headers=headers, timeout=10) if resp.status_code == 200: return parse_page(resp.text) if resp.status_code == 418: time.sleep(3 + random.uniform(0, 1)) continue except requests.RequestException: time.sleep(2 ** attempt) return [] def parse_page(html): soup = BeautifulSoup(html, "html.parser") rows = [] for item in soup.select(".grid_view .item"): movie_id = item.select_one("a")["href"].split("/")[-2] title = item.select_one(".title").get_text(strip=True) # 提取年份、评分、评价人数等字段后追加到 rows rows.append((movie_id, title)) return rows def run(use_thread=True): pages = range(10) results = [] if not use_thread: for p in pages: results.extend(fetch_page(p)) return results with ThreadPoolExecutor(max_workers=4) as pool: futures = {pool.submit(fetch_page, p): p for p in pages} for fu in as_completed(futures): with lock: results.extend(fu.result() or []) return results这段代码里最关键的是max_workers=4。豆瓣的限速是按出口 IP 来的,线程开到 32 只会让请求更快撞上频率限制。延迟区间DELAY_RANGE控制在 1.2 到 2.8 秒随机分布,避免固定间隔被识别为机器行为。418是豆瓣常见的反爬返回码,看到 418 要增加延迟并重试,而不是继续硬闯。
3.3 断点续爬:如何保证重复运行不污染结果
爬虫最怕的不是跑崩,而是重启后产生重复数据。因为SEEN_URLS是内存里的集合,进程一结束就没了,所以我会用数据库里的movie_id来初始化去重集合。
def load_seen_from_db(): seen = set() # 伪代码示例:从 MySQL 查出所有 movie_id for mid in query("SELECT DISTINCT movie_id FROM movie_facts"): seen.add(mid) return seen这样即使只爬了 100 部就中断,重启后也会跳过已有数据。配合crawl_log的page_no,还能精确知道哪些页已经成功,哪些页需要重试。写入时使用 MySQL 的幂等更新,避免唯一索引冲突报错:
INSERT INTO movie_facts (movie_id, title, directors, actors, year, country, genre, rating, votes) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s) ON DUPLICATE KEY UPDATE rating = VALUES(rating), votes = VALUES(votes);这里的逻辑是:如果movie_id已存在,就更新评分和评价人数,不新增记录。这样同样一部电影的评分被刷新,但行数不会增长,后续 Spark 统计时不需要再花时间做全局去重。
反爬控制可以整理成一张自查表,答辩时考官问起来也能对答如流。
| 反爬现象 | 常见原因 | 处理策略 |
|---|---|---|
| 返回 418 | 请求频率过高 | 单线程或 2 线程 + 随机延迟 |
| 返回 403 | 缺少有效请求头 | 轮换 User-Agent,必要时保留 Cookie |
| 页面可访问但解析为空 | 被重定向到验证页 | 检查最终 URL 和页面标题特征 |
| 数据库唯一键冲突 | 重复运行 | 使用 ON DUPLICATE KEY UPDATE |
4. Spark 数据分析的 ETL 计算与提交参数
4.1 Spark JDBC 读 MySQL:分区与下推
爬虫把数据落库之后,Spark 的任务就开始了。这里最常见的错误是把整个表一股脑读进 Spark,然后靠filter去筛。正确的做法是让 JDBC Reader 在 MySQL 侧完成分区读取和下推过滤。
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("DoubanSparkETL") \ .config("spark.sql.shuffle.partitions", "4") \ .getOrCreate() df = spark.read \ .format("jdbc") \ .option("url", "jdbc:mysql://127.0.0.1:3306/douban_movie?useSSL=false&serverTimezone=Asia/Shanghai") \ .option("dbtable", "movie_facts") \ .option("user", "root") \ .option("password", "your_password") \ .option("driver", "com.mysql.cj.jdbc.Driver") \ .option("partitionColumn", "id") \ .option("lowerBound", 1) \ .option("upperBound", 250) \ .option("numPartitions", 4) \ .option("fetchSize", "500") \ .load()partitionColumn必须是一个数字列,这里用自增主键id。lowerBound和upperBound决定分区范围,numPartitions决定并行度。对 250 条数据来说,4 到 8 个分区足够;分区太多会产生大量空任务,反而拖慢作业。fetchSize控制每次从 MySQL 拉取的行数,数据量大时可以减小到 200,避免单次返回过多占用 executor 内存。
4.2 用 Spark SQL 算三个指标:榜单、年代趋势、类型分布
读进来的 DataFrame 注册成临时视图,就可以用纯 SQL 完成分析。下面这三个指标是豆瓣电影分析里最常用、也最能体现 Spark 能力的三板斧。
df.createOrReplaceTempView("movie_facts") # 指标一:评分 + 评价人数的综合榜单 rank_top = spark.sql(""" SELECT title, rating, votes FROM movie_facts WHERE votes > 0 ORDER BY rating DESC, votes DESC LIMIT 20 """) # 指标二:每年上映数量和平均评分 year_trend = spark.sql(""" SELECT year, COUNT(*) AS cnt, ROUND(AVG(rating), 2) AS avg_rating FROM movie_facts WHERE year IS NOT NULL AND year > 1900 GROUP BY year ORDER BY year """) # 指标三:类型占比(类型字段按逗号分隔) genre_df = spark.sql(""" SELECT trim(genre) AS genre, COUNT(*) AS cnt FROM ( SELECT explode(split(genre, ',')) AS genre FROM movie_facts WHERE genre IS NOT NULL AND genre != '' ) t GROUP BY trim(genre) ORDER BY cnt DESC """)指标一用votes > 0过滤掉没有评价人数的脏数据。指标二里ROUND(AVG(rating), 2)保证输出到数据库后不会出现过长小数。指标三用了explode(split(...))把逗号分隔的类型拆成多行,这是 Spark SQL 处理一对多字段的标准姿势。注意先trim再分组,否则 “剧情 爱情” 和 “剧情 爱情 ” 会被当成两个类型。
算出结果后,把结果表写回 MySQL,供可视化层查询。每个指标独立写入一张结果表,避免前端直接查询原始宽表。
rank_top.write.mode("overwrite") \ .format("jdbc") \ .option("url", "jdbc:mysql://127.0.0.1:3306/douban_movie?useSSL=false&serverTimezone=Asia/Shanghai") \ .option("dbtable", "report_rank_top") \ .option("user", "root") \ .option("password", "your_password") \ .save()mode("overwrite")会先删除目标表再重新写入,保证了结果表始终与最新原始数据一致。实际开发中我更习惯把原始表、结果表分库或加前缀report_,这样分析师不会误把聚合结果当原始数据。
4.3 spark-submit 参数对照表与内存调优
Spark 参数是面试和答辩的高频考点,不需要背全部,但下面这几个必须能解释清楚。
| 参数 | 示例值 | 作用 |
|---|---|---|
--master local[4] | local[4] | 本地 4 线程运行,不是 4 个 Executor |
--driver-memory | 2g | Driver 可用内存,本地模式同时承担计算 |
--executor-memory | 2g | 每个 Executor 的 JVM 堆内存 |
spark.sql.shuffle.partitions | 4 | 控制 Shuffle 后分区数,默认 200,小数据必须调低 |
spark.sql.adaptive.enabled | true | 开启动态合并与优化 |
--jars | mysql-connector-j-8.0.33.jar | 提供 MySQL JDBC 驱动 |
最小提交命令如下:
spark-submit \ --master local[4] \ --driver-memory 2g \ --executor-memory 2g \ --conf spark.sql.shuffle.partitions=4 \ --conf spark.sql.adaptive.enabled=true \ --jars /opt/jars/mysql-connector-j-8.0.33.jar \ etl.py这里有个容易踩的坑:local[4]只代表本地启动 4 个线程,不涉及--executor-memory。只有在 standalone 或 YARN 集群模式下,executor-memory才有实际意义。spark.sql.shuffle.partitions默认是 200,意味着即使只有 250 条数据,Shuffle 也会分成 200 个任务。小数据集不改这个参数,你会发现日志里刷出大量空任务,运行时间全耗在任务调度上。
内存调优的通用思路是:Driver 负责规划和收集结果,不需要给太大;Executor 才是真正执行分组、排序的地方。本地调试时通常各给 2g 足够,真正跑集群时再根据数据量等比放大,并保留一定比例内存给 JVM 堆外和系统开销。
5. 可视化 API 与答辩前的一致性验证
5.1 Flask 接口:把 Spark 结果表变成 JSON
Spark 计算结果已经落到 MySQL,可视化层就不需要再启动 Spark 作业了。这里用 Flask 提供一个只读 JSON 接口,前端无论用 ECharts、Tableau 还是大屏工具,都可以直接拉取数据。
from flask import Flask, jsonify import pymysql app = Flask(__name__) def query(sql): conn = pymysql.connect( host="127.0.0.1", user="root", password="your_password", database="douban_movie", charset="utf8mb4" ) with conn.cursor() as cur: cur.execute(sql) cols = [d[0] for d in cur.description] rows = [dict(zip(cols, r)) for r in cur.fetchall()] conn.close() return rows @app.route("/api/trend") def trend(): return jsonify(query( "SELECT year, cnt, avg_rating FROM report_year_trend ORDER BY year" ))这个设计的核心是:查询不要走 Spark,否则每次打开页面都要经历一次作业启动的几十秒延迟。批计算一次,服务层查 MySQL,这是最容易拿到答辩加分点的架构决策。
5.2 ECharts 异步加载数据,30 行内完成大屏雏形
前端只需要在页面加载时调用接口,填入setOption。以年代趋势折线图为例:
fetch('/api/trend') .then(res => res.json()) .then(data => { myChart.setOption({ xAxis: { type: 'category', data: data.map(d => d.year) }, yAxis: { type: 'value' }, series: [{ type: 'line', data: data.map(d => d.avg_rating), smooth: true }] }); });其他图表的逻辑完全一致,换一下 API 路径和series.type即可。这种前后端分离的方式,比在 Python 里直接拼 HTML 字符串更清晰,也更容易扩展图表类型。
5.3 答辩前自检:三条 SQL 验证数据血缘与幂等
交代码之前,我会把下面三条 SQL 存成verify.sql,当场跑给评审看,效果比任何截图都有说服力。
-- 1. 原始表与结果表数量核对 SELECT (SELECT COUNT(*) FROM movie_facts) AS source_cnt, (SELECT COUNT(*) FROM report_year_trend) AS result_cnt; -- 2. 评分异常值检查 SELECT COUNT(*) AS abnormal_cnt FROM movie_facts WHERE rating <= 0 OR rating > 10 OR votes < 0; -- 3. 重复 movie_id 检查 SELECT movie_id, COUNT(*) AS cnt FROM movie_facts GROUP BY movie_id HAVING COUNT(*) > 1 LIMIT 5;第一条验证 ETL 过程有没有丢数据;第二条验证字段约束,评分必须在 0 到 10 之间;第三条验证爬虫幂等,如果movie_id重复,说明写入逻辑有问题。如果crawl_log里上一轮抓取 250 行、这一轮仍抓 250 行,数据库也是 250 个movie_id,整个数据处理链路就闭环了。最后把这四个验证点做成脚本放进项目根目录,评审提问时直接执行,用实际输出回答,比口头解释更有说服力。
本文还有配套的精品资源,点击获取