
简介这是一套基于 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原始记录存储层解析后的字段MySQLmovie_facts 表计算层movie_facts 表Spark JDBC Spark SQLreport_* 结果表服务层report_* 结果表FlaskJSON API展示层JSON APIECharts可视化大屏这个流程里最容易做错的是直接用爬虫结果怼进 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) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; 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) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;rating用DECIMAL(3,1)而不是FLOAT因为评分只需要一位小数精确小数可以避免 Spark 读取后出现 9.1999999 这类误差。movie_id加唯一索引是写幂等插入的关键。quote字段可能包含引号或特殊符号所以数据库字符集使用utf8mb4兼容 emoji 和生僻字。crawl_log表不存业务数据只记录每一页爬取开始时间、结束时间、获取条数和状态。它的作用是回答答辩时最常见的问题“如果爬虫跑到一半断了你怎么知道从哪继续” 有了这张表就能直接统计上次成功页码和状态而不是靠日志文件翻找。3. 豆瓣电影爬虫的并发采集与反爬控制细节3.1 只抓 Top 250 的三类字段别上来就做全站爬虫豆瓣全站爬虫涉及用户登录态、评论分页、鉴权等多个复杂环节作为毕设没有必要还会因为请求量过大导致 IP 被封。所以最稳妥的抓取范围是 Top 25010 页、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 fhttps://movie.douban.com/top250?start{page * 25} headers HEADER_POOL[page % len(HEADER_POOL)] for attempt in range(3): try: resp requests.get(url, headersheaders, timeout10) 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(stripTrue) # 提取年份、评分、评价人数等字段后追加到 rows rows.append((movie_id, title)) return rows def run(use_threadTrue): pages range(10) results [] if not use_thread: for p in pages: results.extend(fetch_page(p)) return results with ThreadPoolExecutor(max_workers4) 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_workers4。豆瓣的限速是按出口 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 UPDATE4. 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?useSSLfalseserverTimezoneAsia/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?useSSLfalseserverTimezoneAsia/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-memory2gDriver 可用内存本地模式同时承担计算--executor-memory2g每个 Executor 的 JVM 堆内存spark.sql.shuffle.partitions4控制 Shuffle 后分区数默认 200小数据必须调低spark.sql.adaptive.enabledtrue开启动态合并与优化--jarsmysql-connector-j-8.0.33.jar提供 MySQL JDBC 驱动最小提交命令如下spark-submit \ --master local[4] \ --driver-memory 2g \ --executor-memory 2g \ --conf spark.sql.shuffle.partitions4 \ --conf spark.sql.adaptive.enabledtrue \ --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 结果表变成 JSONSpark 计算结果已经落到 MySQL可视化层就不需要再启动 Spark 作业了。这里用 Flask 提供一个只读 JSON 接口前端无论用 ECharts、Tableau 还是大屏工具都可以直接拉取数据。from flask import Flask, jsonify import pymysql app Flask(__name__) def query(sql): conn pymysql.connect( host127.0.0.1, userroot, passwordyour_password, databasedouban_movie, charsetutf8mb4 ) 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整个数据处理链路就闭环了。最后把这四个验证点做成脚本放进项目根目录评审提问时直接执行用实际输出回答比口头解释更有说服力。本文还有配套的精品资源点击获取