ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

基于Python与Spark的豆瓣电影爬虫与数据分析可视化实践

基于Python与Spark的豆瓣电影爬虫与数据分析可视化实践 简介一份基于PythonSpark的豆瓣电影爬虫与数据分析可视化系统源码及数据库文件面向毕业设计、课程设计与期末大作业场景。项目实现从豆瓣抓取电影数据、基于Spark分布式清洗统计、再以可视化图表展示的完整链路代码注释详尽新手也能快速理解并部署运行。压缩包共241个文件大小仅5.65MB组成上以165个XML配置文件、22个Java源码和6个Class编译文件为核心同时包含Python脚本、SQL数据库文件、HTML/CSS/JS前端页面、CSV结果数据及依赖JAR包覆盖开发、编译、运行、展示各环节。当前已有248人学习下载适合需要搭建高分毕设项目的同学参考。项目曾获98分评价导师认可度高Spark分析作业覆盖词频、电影类型、评分等级、评论数量、上映年份等维度输出分区文件可复用配套数据库与可视化页面能直接呈现电影市场分析效果兼具学习与二次开发价值。1. 基于Python与Spark的电影数据全链路项目怎么拆一提到“基于PythonSpark的豆瓣电影爬虫和数据分析可视化系统”多数人第一反应是“又是一个爬虫套壳项目”。但真正动手做过毕设或工程化练习的人会知道这个标题的含金量在于它把一条完整的数据链路串了起来数据采集、数据落地、分布式计算、指标分析、可视化呈现每一环都有独立的坑。爬虫只是入口Spark才是处理层的核心而可视化效果直接影响答辩或演示时别人对系统完成度的判断。这套项目的合理技术栈是Python负责爬虫采集和接口封装Spark负责对豆瓣电影评分、评论、类型、地区、上映时间等字段做清洗和聚合分析分析结果写入数据库最后用可视化框架把排行榜、评分分布、类型偏好等内容做成可交互图表。适合的人群是正在做毕设的本科生、想转行数据工程但缺一个完整项目经验的求职者以及想把手头零散数据集变成可查询、可展示系统的开发人员。下面按“环境搭建与架构设计 → 爬虫实现 → Spark分析 → 可视化与查询”这条主线展开最后补一节如何把这套系统从“能跑”打磨成“能讲”。2. 环境准备与数据管道设计Python、Spark与数据库的选型组合2.1 Python与Spark版本匹配比“装好”更重要的是“配好”标题里同时出现Python和Spark第一个关键问题是版本匹配。Spark虽然提供了PySpark接口但它对Python版本、Java版本、Hadoop生态的兼容性都有明确约束。以目前生产环境中使用最广泛的Spark 3.x系列为例它依赖Java 8/11/17Python 3.8及以上版本都支持但如果你用的是Spark 2.4.x那Python 3.7就是上限Python 3.8会出现兼容性告警。安装建议使用Anaconda管理Python环境单独为这个项目创建虚拟环境不要把系统自带的Python环境污染掉。创建完环境后用pip安装pyspark安装包会自动拉取对应版本的JVM依赖。这里有个常见误区很多人以为安装pyspark就是安装了完整的Spark集群其实pyspark包自带的是standalone模式的Spark运行环境可以让你在单机上提交Spark任务和训练分布式计算逻辑但要用到HDFS或YARN集群还必须单独下载Spark发行版并配置环境变量。conda create -n douban-spark python3.9 -y conda activate douban-spark pip install pyspark3.3.0 java -version这段命令的核心思路是先把Python解释器固定到3.9因为3.9同时兼容Spark 3.3和大多数爬虫库。安装pyspark时指定版本号是强烈建议的做法——不指定版本会拉到最新的Spark版本有时候会超出你本机Java的兼容范围。最后的java -version是验证Spark运行环境的必要步骤Spark启动时会先检查Java环境如果你本机没有装JDK或装的是OpenJ9而不是HotSpot启动时会直接报错。2.2 数据库选型为什么不用MySQL而选SQLite MongoDB组合标题提到了“数据库文件”这说明项目交付物里带着一个已经落地的数据库。豆瓣电影数据分析场景下数据量级通常在几万到几十万条之间这个量级MySQL完全扛得住但毕设项目要考虑的是“演示方便”和“结构清晰”。推荐组合结构化数据电影基本信息、评分、导演、演员存放SQLite一个单文件就能持久化全部数据评审老师打开项目就能直接运行不需要额外安装数据库服务评论详情、用户观影记录这类半结构化数据存放MongoDB因为豆瓣评论的字段不固定有人写了长评有人只给短评MongoDB的文档模型更贴合。CREATE TABLE movies ( id INTEGER PRIMARY KEY, title TEXT NOT NULL, douban_id TEXT UNIQUE, rating REAL DEFAULT 0, director TEXT, actors TEXT, genres TEXT, release_date TEXT, runtime INTEGER, language TEXT, region TEXT, create_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP );这个SQL建表语句突出了几个设计点douban_id字段设置为UNIQUE约束防止爬虫重复采集同一部电影genres字段设计成TEXT而不是独立的关联表是对数据分析场景的妥协——做技术演示时用逗号分隔类型比JOIN关联表更直观如果你后续要做复杂推荐系统再拆多对多表也来得及。2.3 数据管道设计图从抓取到可视化的四层架构整个系统的数据流通路径可以概括为四个层次采集层、存储层、计算层、展示层。采集层由Python爬虫构成负责按调度策略抓取豆瓣页面的电影元数据和短评内容拿到原始HTML或JSON后做字段解析存储层将解析后的数据写入本地文件或数据库计算层由Spark承担任务是完成数据清洗、缺失值填充、评分区间划分、类型热度统计等分析工作展示层读取Spark的输出结果生成可视化图表。这套架构的核心价值不是“技术新颖”而是每一层都可以独立替换。比如采集层从requests换成Scrapy存储层从SQLite换成PostgreSQL展示层从ECharts换成PowerBI都不会影响其他层的逻辑。这种可替换性是毕设答辩时的高频加分点也对应了实际数据团队里“职责单一、接口明确”的工程原则。3. 豆瓣电影爬虫实现从requests到Scrapy的完整演化路径3.1 反爬策略分析与请求头伪装从基础headers到登录态处理豆瓣在反爬上的策略经历了多次升级。早期只需要在请求头里带一个正常的User-Agent就能抓取Top250榜单后来增加了Cookie校验和访问频率限制现在对高频访问的IP会做封禁处理。爬虫模块的设计重点在于如何让请求看起来像一个正常人用浏览器浏览。先看基础的requests实现import requests from fake_useragent import UserAgent ua UserAgent() headers { User-Agent: ua.random, Accept: text/html,application/xhtmlxml,application/xml;q0.9,*/*;q0.8, Accept-Language: zh-CN,zh;q0.8,en-US;q0.5,en;q0.3, Connection: keep-alive } def fetch_page(url): session requests.Session() session.headers.update(headers) try: resp session.get(url, timeout10) resp.raise_for_status() resp.encoding utf-8 return resp.text except requests.RequestException as e: print(f请求失败: {url}, 错误: {e}) return None这个实现采用了一个关键组件requests的Session对象。Session会自动保存请求后的Cookie状态并且快于直接使用requests.get连续请求。使用fake_useragent库随机生成User-Agent是为了降低同一浏览器标识连续请求的识别概率。真实的豆瓣爬虫还需要处理Cookie常见做法是先在浏览器里登录豆瓣并复制Cookie到配置文件中这只适用于毕设演示场景下的低频率采集。3.2 数据解析策略XPath比BeautifulSoup更适配豆瓣页面结构解析HTML页面有两个主流选择BeautifulSoup和lxml的XPath。针对豆瓣电影这类结构非常规整的列表页面XPath是更好的选择因为它可以直接通过路径定位到节点不需要像BeautifulSoup那样先find再find_all二次筛选。from lxml import etree def parse_movie_list(html): tree etree.HTML(html) items tree.xpath(//div[classitem]) movies [] for item in items: title item.xpath(.//span[classtitle]/text()) rating item.xpath(.//span[classrating_num]/text()) quote item.xpath(.//p[classquote]/span/text()) movie { title: title[0] if title else , rating: float(rating[0]) if rating else 0.0, quote: quote[0] if quote else } movies.append(movie) return moviesXPath解析的关键在于理解豆瓣页面的DOM结构。一个div[classitem]节点代表一部电影里面嵌套了标题、评分、引言等多个信息节点。使用XPath的.//写法而不是//表示从当前节点开始搜索子节点这是避免跨条目抓取错位的重要细节。评分字段做float类型转换前需要判断列表是否为空否则空列表会抛出索引越界异常。3.3 爬虫并发设计线程池与请求频率的平衡爬虫的并发设计是影响抓取效率和安全性的最大变数。盲目拉高线程数量比如直接上100个线程会在几分钟内触发豆瓣的封IP策略导致一小时甚至一天内无法继续采集。在多次实测和调研里对豆瓣而言5~8个并发线程是相对安全的范围。from concurrent.futures import ThreadPoolExecutor, as_completed def crawl_with_threads(urls, max_workers6, delay1.0): results [] with ThreadPoolExecutor(max_workersmax_workers) as executor: future_to_url {executor.submit(fetch_page, url): url for url in urls} for future in as_completed(future_to_url): url future_to_url[future] try: html future.result() if html: movies parse_movie_list(html) results.extend(movies) else: print(f页面为空: {url}) except Exception as e: print(f任务执行异常: {url}, {e}) time.sleep(delay) return results这个实现中值得关注的参数是delay1.0它让每个任务完成后线程自动休眠1秒目的是把请求频率控制在一个安全区间内。即使配置了max_workers6的并发数实际每秒请求量也只有6次左右远低于单线程但足够支撑几万条数据的采集周期。对于几千条数据量级的毕设项目这样的设计在效率和安全性上是平衡的。3.4 分布式爬虫扩展从单机到多机调度的升级方向如果你不满足于单机采集想挑战“大规模数据采集”的架构能力Scrapy就是比requests更完善的选择。Scrapy底层集成了Twisted异步网络框架单个Spider实例就能支撑较高的并发量同时它的中间件机制可以用来扩展代理IP切换、自定义User-Agent轮换等功能。Scrapy还支持通过scrapy-redis组件实现分布式调度核心思想是让所有爬虫节点共享同一个Redis任务队列。主节点负责往Redis里推送待抓取的URL从节点从Redis队列里取URL、抓取、把解析结果写入公共存储这样多台机器就能协作完成同一批数据的采集。# settings.py 配置 SCHEDULER scrapy_redis.scheduler.Scheduler DUPEFILTER_CLASS scrapy_redis.dupefilter.RFPDupeFilter REDIS_URL redis://localhost:6379 SCHEDULER_PERSIST True这个配置让Scrapy从默认的单机任务队列切换到Redis共享队列。通过Redis持久化任务进度即使爬虫进程崩溃重启也不会丢失已经完成的任务状态。团队里不同角色的协作可以拆分为调度节点和采集节点而数据格式与解析逻辑完全复用同一套Spider代码。4. Spark数据分析从DataFrame清洗到评分预测模型4.1 数据加载与初始化SparkSession的配置参数优化Spark接入数据的起点是构建SparkSession负责统一管理SQL、DataFrame、Streaming等模块。在该项目中需要将数据库中的电影数据加载到Spark分布式的DataFrame上。from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(DoubanMovieAnalysis) \ .master(local[4]) \ .config(spark.driver.memory, 2g) \ .config(spark.sql.shuffle.partitions, 8) \ .config(spark.serializer, org.apache.spark.serializer.KryoSerializer) \ .getOrCreate() df spark.read \ .format(jdbc) \ .option(url, jdbc:sqlite:douban_movies.db) \ .option(dbtable, movies) \ .load()这里的参数设计值得逐项说明master(local[4])表示在本地用4个线程模拟分布式执行如果你机器是8核可以适当调大spark.driver.memory控制驱动程序的内存大小爬取的数据量超过50万条时需要至少4g更重要的是spark.sql.shuffle.partitions参数它决定了shuffle阶段产生的分区数量默认值是200但如果你只有几十万条数据200个分区只会带来任务调度开销。调整为8个分区后聚合操作的速度会有明显提升。4.2 数据清洗与特征工程去重、缺失值填充与类型拆分真实爬回来的数据会有很多脏数据常见问题包括评分字段空值、上映日期格式不统一、类型字段包含多个类型如“剧情/爱情/动画”。数据分析和清洗环节就是处理这些不一致。from pyspark.sql.functions import col, split, when, isnan, isnull # 去除douban_id重复的记录 df_clean df.dropDuplicates([douban_id]) # 过滤评分为空或为NaN的记录 df_valid df_clean.filter( col(rating).isNotNull() ~isnan(col(rating)) (col(rating) 0) ) # 拆分genres字段 df_genres df_valid.withColumn(genre_list, split(col(genres), /)) # 填充缺失的region字段 df_filled df_genres.fillna({region: 未知, language: 未知})以上代码中值得注意的几点dropDuplicates的参数是指定列名只根据douban_id去重而不是整行去重这是为了应对爬虫重跑场景下部分字段更新但ID不变的情况isnan和isnull是两种不同的判断方式isnan针对的是数值类型的NaN值isnull针对SQL语义下的NULL值评分字段两种类型都可能出现所以并在一起处理更稳妥最后一个fillna直接指定了每列的默认填充值在实际工程里这样逐列指定比全局统一填充更精确——比如region填充“未知”合理但如果把title也填了默认值反而会污染数据。4.3 核心指标计算评分分布、类型热度与年度趋势的SQL化实现Spark分析的核心不是写复杂的算法而是把业务需求翻译成分布式的聚合计算。常用做法是直接注册临时视图用SQL写分析逻辑因为SQL语义更清晰也好跟评审老师解释。df_filled.createOrReplaceTempView(movie_view) # 评分分布统计 rating_dist spark.sql( SELECT CAST(rating AS INT) as rating_group, COUNT(*) as cnt FROM movie_view GROUP BY CAST(rating AS INT) ORDER BY rating_group ) # 类型热度统计 genre_rank spark.sql( SELECT genre, COUNT(*) as movie_cnt, ROUND(AVG(rating), 2) as avg_rating FROM ( SELECT explode(genre_list) as genre, rating FROM movie_view ) GROUP BY genre ORDER BY movie_cnt DESC ) # 年度上映数量与均分趋势 year_trend spark.sql( SELECT SUBSTR(release_date, 1, 4) as year, COUNT(*) as total, ROUND(AVG(rating), 2) as avg_score FROM movie_view WHERE release_date IS NOT NULL AND LENGTH(release_date) 4 GROUP BY SUBSTR(release_date, 1, 4) ORDER BY year )三个SQL分别服务于三类可视化图表的需求。第一个评分分布用取整分组因为豆瓣评分区间是0到10按整数分组后能直观呈现“高分电影是少数、中间分数是大多数”的长尾分布第二个类型热度统计里使用了Spark SQL的explode函数它的作用是把数组类型的每个元素展开成一行——这正好应对了genres字段被拆分成genre_list数组后的统计需求第三个年度趋势用SUBSTR截取前四位作为年份如果字段里混入了“未知”文本LENGTH(release_date) 4的条件会尽量排除非法数据干扰。4.4 基于ALS的评分预测模型用电影协同过滤补全评分缺失值爬虫抓取的数据中有一个常见现象评分字段缺失但电影的评论数和想看人数都有。如果你在展示阶段能补充一个“预测评分”的功能——基于已有的评分数据预测缺失评分会显著提升系统完整度。Spark MLlib中内置了ALS交替最小二乘算法可以直接用于协同过滤。它的核心原理是将用户-电影评分矩阵分解为两个低维矩阵的乘积一个代表用户特征一个代表电影特征通过迭代优化最小化真实评分与预测评分之间的误差。from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator # 准备训练数据用户ID、电影ID、评分 als_data df_filled.select( col(id).alias(movieId), col(rating).alias(rating) ).withColumn(userId, lit(1)) train, test als_data.randomSplit([0.8, 0.2], seed42) als ALS( maxIter10, regParam0.1, userColuserId, itemColmovieId, ratingColrating, coldStartStrategydrop ) model als.fit(train) predictions model.transform(test) evaluator RegressionEvaluator( metricNamermse, labelColrating, predictionColprediction ) rmse evaluator.evaluate(predictions) print(fRMSE: {rmse})在这个演示代码里userId使用常量1是因为毕设场景下往往没有真实的用户行为数据这是一种过渡设计更完整的ALS模型需要从爬取的评论者打分记录中构建真实的userId列。coldStartStrategydrop参数很关键——对于测试集中存在的、但训练时未出现过的电影IDALS无法给出预测默认行为是抛异常设置成drop后会自动丢弃这些行不中断整个评估流程。ALS模型训练出的拟合结果只是辅助信息不要篡改原始评分预测列单独展示更能体现模型的独立性。5. 可视化系统设计从ECharts图表到Flask交互查询页面5.1 可视化组件选型ECharts是数据分析项目的最佳落点可视化层的技术选型中ECharts因为支持关系图、词云、地图热力等多种复杂图表类型并且配置简单已经成为数据分析类项目的首选也常被认为是数据可视化工具的典型代表。这套豆瓣电影系统推荐用ECharts原因有三个其一ECharts在MIT开源协议下可以免费商用其二国内社区和文档非常完善评分分布、折线图等常用配置直接抄官方示例改数据即可其三ECharts支持Canvas和SVG双渲染模式数据量大时可以切到Canvas。一个评分分布的可视化配置如下var chart echarts.init(document.getElementById(ratingChart)); var option { title: { text: 豆瓣电影评分分布, left: center }, tooltip: { trigger: item }, xAxis: { type: category, data: [1分, 2分, 3分, 4分, 5分, 6分, 7分, 8分, 9分, 10分] }, yAxis: { type: value }, series: [{ name: 电影数量, type: bar, data: [2, 8, 15, 42, 88, 156, 342, 521, 267, 45], itemStyle: { color: #FF7F50 } }] }; chart.setOption(option);这段配置的逻辑很清楚x轴定义为评分区间series里配置条形图数据。注意itemStyle中显式指定了一个颜色因为评分分布这种图表配色统一比突出某根柱子更耐看。实际数据你从Spark分析的结果中导出后填充到data数组里就行。5.2 Flask后端搭建将Spark分析结果封装为API接口可视化页面需要数据源推荐用Flask搭建一个轻量级Web应用负责接收前端请求、从数据库读取分析结果、返回JSON数据。这样前端页面可以保持纯粹只负责数据渲染。from flask import Flask, jsonify, request import sqlite3 app Flask(__name__) DB_PATH douban_movies.db def query_db(sql, params()): conn sqlite3.connect(DB_PATH) conn.row_factory sqlite3.Row cur conn.cursor() cur.execute(sql, params) rows cur.fetchall() conn.close() return [dict(row) for row in rows] app.route(/api/rating_dist) def rating_dist(): data query_db( SELECT CAST(rating AS INT) as rating_group, COUNT(*) as cnt FROM movies WHERE rating IS NOT NULL GROUP BY rating_group ORDER BY rating_group ) return jsonify(data) app.route(/api/genre_rank) def genre_rank(): data query_db( SELECT genre, COUNT(*) as movie_cnt, AVG(rating) as avg_rating FROM ( SELECT TRIM(value) as genre FROM movies, json_each([ || REPLACE(genres, /, ,) || ]) ) WHERE genre ! GROUP BY genre ORDER BY movie_cnt DESC ) return jsonify(data) if __name__ __main__: app.run(debugTrue, port5000)这里特别说一下genre_rank接口的实现SQLite没有内置explode函数所以用json_each函数加字符串拼接的方式模拟了Spark中explode的效果。先把喜剧/爱情/动画替换成[喜剧,爱情,动画]的JSON格式再用json_each把它展开成多行——过程和Spark的数组展开逻辑完全一致只是语法不同。这个思路也说明了一个工程原则不是所有分析都必须在Spark里做轻量级查询交给数据库重量级聚合才交给Spark。5.3 可视化大屏页面不用写代码的方式搭建展示面板如果你不想为每个图表单独写HTML页面也可以把Flask和ECharts结合成一个可视化仪表盘页面但如果你没有React这类前端框架的经验有一个更省力的思路把Spark分析结果导出为JSON文件放到任意一个支持数据替换的Vue或React模板中甚至也可以用Dashboard类的可视化大屏搭建工具来加速。操作方式很直接# 将Spark分析结果导出为JSON供前端直接加载 df_result.write.json(./output/rating_dist.json, modeoverwrite)用这种方式将分析结果保存下来之后前端项目只负责读取JSON通过Ajax请求打断点调试避免了每次调试都要同时启动Flask后端和Spark任务的麻烦。对毕设演示来说Spark计算一次、结果落盘、后续所有展示环节都读快照这个效率是最高的。6. 让项目从“能跑”升级为“能讲”代码质量与答辩技巧的最后一公里很多人写完这种项目跑通了就以为结束了直到答辩被问到“你这个分析结果和爬虫策略有什么关系”才卡壳。最后一节给三个可落地的优化方向是把项目等级从及格拉到优秀的核心技巧。第一个方向是爬虫采集全过程的完整性记录。在爬虫代码里增加一个crawl_log表记录每次采集的起始时间、结束时间、成功条数、失败条数、耗时秒数。这一块数据在答辩时作用巨大它能说明你对爬虫的容错、重试、调度策略有过真实的考量而不是只写了个一次性脚本。实现上也不复杂在爬虫抓取结束的回调里追加一行日志即可。第二个方向是参数调优经验的可视化对比。以第4章中ALS模型的regParam来说一个有效的做法是写一个小循环让regParam从0.01到0.5依次取值记录每个值对应的RMSE把结果画成折线图。这张图能直观地告诉评委我调过参且知道调参的逻辑依据。reg_params [0.01, 0.05, 0.1, 0.2, 0.5] results [] for rp in reg_params: als ALS(maxIter10, regParamrp, userColuserId, itemColmovieId, ratingColrating, coldStartStrategydrop) model als.fit(train) pred model.transform(test) rmse evaluator.evaluate(pred) results.append((rp, rmse)) print(fregParam{rp}, RMSE{rmse:.4f})这段代码没有引入额外的Python依赖只是把模型训练封装在循环里最后输出一组非递增的RMSE曲线这就能直观地体现正则化参数与过拟合程度之间的关系。最后一个方向是数据血缘说明。在项目的README里用表格标明每个字段的采集来源、清洗规则、落库位置这种看似文档层面的工作实际上推动自己梳理了整个数据链条的合理性也让系统后续的扩展有了清晰的起点。本文还有配套的精品资源点击获取
返回列表