恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
豆瓣影评数据分析全流程:Python爬虫与Spark清洗可视化实践
首页
资讯中心
/
豆瓣影评数据分析全流程:Python爬虫与Spark清洗可视化实践
豆瓣影评数据分析全流程:Python爬虫与Spark清洗可视化实践
发布时间:2026/9/10 14:25:58
简介这是一份面向计算机专业毕业设计与大数据初学者的实战资源围绕豆瓣电影爬虫、Spark数据处理和可视化展示构建完整项目。基于Python编写爬虫获取电影评分与用户评论利用Spark执行数据清洗、转换和聚合分析再通过Matplotlib、Seaborn输出可视化图表适合作为教学案例或人工智能、深度学习方向的毕设参考。压缩包共242个文件大小5.61MB其中Java源码承载Spark分析逻辑Python脚本负责数据采集SQL脚本与CSV文件用于数据存储交换HTML/CSS/JS构成结果展示页面另有项目说明文档辅助上手。目前已有121人学习下载。项目内附详尽README和环境配置指南可帮助读者快速搭建运行环境、理解代码结构与分析流程无论是用于课程设计、毕业设计还是作为大数据入门的系统源码都具有较强的落地参考价值。1. 豆瓣影评数据从爬虫到Spark再到图表的完整闭环很多学习Spark的人手里只有官方文档里那些几十MB的示例数据跑完wordcount就不知道下一步该做什么。这个基于豆瓣电影爬虫和Spark数据分析可视化的毕设给了一条完整链路用Python爬虫抓豆瓣评论和评分经Spark清洗聚合后再用Matplotlib、Seaborn出图表。它不只是一份源码包更像一个能跑的行业案例覆盖数据采集、ETL、统计分析和结果展示四个环节。适合正在选毕设方向的计算机专业学生也适合想验证Spark集群环境是否好用的工程师。资源里包含项目说明文档与多个编译后的class文件能直接看到作者在分析阶段拆出的统计维度比如评分等级、年份趋势、评论词频。2. 爬虫层核心requests BeautifulSoup抓取豆瓣电影评论与评分字段如果把整条数据链路比作管道爬虫就是进水口。这一层最怕两件事解析字段不完整或者请求频率高到被豆瓣限制。豆瓣的页面结构相对稳定短评列表页URL比较规整适合用静态HTML解析不需要上动态渲染的Selenium方案。本项目里的爬虫部分就是基于requests和BeautifulSoup 3实现的我在拆解时发现它把评论内容、评分等级、评论时间都做了结构化输出直接落成CSV给下游Spark读取。2.1 数据字段设计与URL规律豆瓣电影短评页的翻页参数是start每页固定20条URL形如https://movie.douban.com/subject/{movie_id}/comments?start0limit20movie_id是豆瓣内部编号比如《肖申克的救赎》对应1292052。爬虫要抓的核心字段如下字段名来源位置示例值下游用途commentli.comment-item .short经典不用多说分词、词频统计ratingli.comment-item span.rating的title属性力荐 / 推荐映射成1-5评分comment_timeli.comment-item .comment-time2023-10-01 12:00:00年月趋势统计votes.vote-count236评论热度注意span.rating没有数字分值只有title属性比如“力荐”“推荐”“较差”。后续Spark作业需要把这一步映射成数值分方便算分布和均值。2.2 解析代码短评与评分提取我在本地复现了这套爬虫逻辑核心代码保持简洁没有做多线程。先看请求部分# douban_spider.py import time import random import pandas as pd import requests from bs4 import BeautifulSoup COMMENT_URL https://movie.douban.com/subject/{}/comments?start{}limit20 HEADERS { User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0 Safari/537.36, Accept: text/html,application/xhtmlxml,application/xml;q0.9,image/webp,*/*;q0.8, Accept-Language: zh-CN,zh;q0.9,en;q0.8, Referer: https://movie.douban.com/ } def fetch_page(url, retry3): for i in range(retry): try: resp requests.get(url, headersHEADERS, timeout10) if resp.status_code 200: resp.encoding resp.apparent_encoding return resp.text if resp.status_code 403: time.sleep(5 random.random() * 5) continue except requests.RequestException: time.sleep(2) return 这段请求代码设置了完整请求头、超时和重试。resp.apparent_encoding会根据页面内容推断编码比直接写死utf-8更稳当响应码是403时说明触发反爬限制这里没有立刻重试而是用time.sleep随机停顿。Referer字段被设置成豆瓣首页因为豆瓣对图片和部分接口会校验来源。接着是页面解析函数def parse_comment_page(html): soup BeautifulSoup(html, html.parser) items [] for li in soup.select(li.comment-item): try: comment li.select_one(.short).get_text(stripTrue) rating_node li.select_one(span.rating) rating rating_node.get(title) if rating_node else 无评分 time_node li.select_one(.comment-time) comment_time time_node.get(title, ) if time_node else items.append({ comment: comment, rating: rating, comment_time: comment_time }) except AttributeError: continue return items解析时先定位li.comment-item再分别取出短评、评分和评论时间。用try/except AttributeError可以跳过结构异常的条目比如某些短评没有评分标签。这里没有按照requests的text直接切片解析因为豆瓣HTML在不同时间段会增加活动模块BeautifulSoup基于CSS选择器更抗页面微调。2.3 反爬应对与数据落盘抓取循环里我一般会加入随机延时和页面数量控制def collect_comments(movie_id, pages10): rows [] for start in range(0, pages * 20, 20): html fetch_page(COMMENT_URL.format(movie_id, start)) rows.extend(parse_comment_page(html)) time.sleep(random.uniform(2, 5)) return rows if __name__ __main__: data collect_comments(1292052, pages5) df pd.DataFrame(data) df.to_csv(douban_movie_comments.csv, indexFalse, encodingutf-8-sig)random.uniform(2, 5)是多次请求之间的随机睡眠避免固定间隔被识别成脚本。CSV落盘用utf-8-sig编码主要考虑到后续用Excel或Pandas读取时中文不会乱码也方便Spark默认按UTF-8读取。如果抓取量更大建议把pages拆成可配置参数并记录已经抓过的movie_id断点续抓比一次性拉全量更实用。这一层的输出文件douban_movie_comments.csv字段就是comment,rating,comment_time后面Spark作业只需要读这一个文件就能做大部分分析。3. Spark分析作业从RDD清洗到词频统计的类设计源码包里出现了TypeNum.class、LvNum.class、YearNum.class、WordNum.class和WordUtil.class这说明项目的Spark分析层是用Scala写的并且把不同统计维度封装成了独立类。这种设计在毕设里值得借鉴把分词、类型统计、评分等级统计、年份趋势统计拆开既能单独测试也方便后续扩展。这一章我会结合这些类名讲清楚Spark作业从读取、清洗到聚合输出的完整流程并给出可执行的Scala代码。3.1 为什么用Spark而不是Pandas如果只是分析几千条豆瓣短评Pandas完全够用甚至更快。但这个项目的价值在于演示分布式计算流程所以引入了Spark。Spark的优势体现在三个场景数据量超过单机内存、需要多维度聚合、后续要接集群环境。毕设里用到Spark不是为了跑得更快而是为了体现大数据处理框架的完整知识体系。实际编码时也不需要写低层RDD直接用Spark DataFrame SQL即可。项目里出现.class文件说明作者用sbt package或mvn package打成JAR后提交到集群运行时通过spark-submit指定主类。Spark的DataFrame API比RDD更易读还能自动走Catalyst优化器适合做多维聚合。3.2 Spark作业结构与WordUtil实现分析入口是一个Scala对象负责创建SparkSession和读取CSV// MovieAnalysis.scala package com.douban.etl import com.douban.util.WordUtil import org.apache.spark.sql.{SaveMode, SparkSession} import org.apache.spark.sql.functions._ object MovieAnalysis { def main(args: Array[String]): Unit { val spark SparkSession.builder() .appName(DoubanMovieAnalysis) .config(spark.sql.shuffle.partitions, 8) .getOrCreate() import spark.implicits._ val raw spark.read .option(header, true) .option(encoding, UTF-8) .csv(args(0)) val cleaned raw .filter($comment.isNotNull $comment ! ) .withColumn(rating_score, when($rating 力荐, 5) .when($rating 推荐, 4) .when($rating 还行, 3) .when($rating 较差, 2) .otherwise(1)) cleaned.write.mode(SaveMode.Overwrite).parquet(args(1)) spark.stop() } }这里把“评分等级”映射成rating_score方便后续计算平均分和分布。注意when的otherwise(1)兜底因为有些短评可能没有评分标签。spark.sql.shuffle.partitions设置为8在数据量小的时候能减少不必要的任务切分如果集群资源大这个值要调大。词频统计是本项目的亮点依赖WordUtil对象做中文文本清洗package com.douban.util object WordUtil { private val stopWords Set(这里, 什么, 一个, 没有, 真的, 还是, 就是, 可以, 这个, 那个, 自己) def analyze(text: String): Array[String] { text.replaceAll([^\\u4e00-\\u9fa5a-zA-Z0-9 ], ) .split(\\s) .filter(word word.nonEmpty !stopWords.contains(word)) } }WordUtil用正则把非中英文数字的字符替换成空格然后按空白符切分最后去掉空串和停用词。这个实现比引入HanLP轻量很多适合纯教学场景。真实项目中这个工具类还会加入词性过滤和精确模式分词但替换成jieba或Ansj成本也不高。3.3 统计维度聚合逻辑与结果输出项目源码中出现的结果类对应以下统计逻辑类/文件名统计维度输入字段输出结果含义TypeNum.class电影类型分布type剧情片数量、喜剧片数量LvNum.class评分等级分布rating_score5分、4分、3分、2分、1分各自数量YearNum.class上映年份趋势year每年上映电影评论量或数量CommontNum.class评论时间趋势comment_time按月统计评论量变化WordNum.class评论热词comment高频词及出现次数聚合代码可以直接用DataFrame APIval lvDF cleaned.groupBy(rating_score) .count() .orderBy($rating_score.desc) lvDF.write.mode(SaveMode.Overwrite) .option(header, true) .csv(/output/lv_num) val wordDF cleaned .select(explode(udf((s: String) WordUtil.analyze(s)).apply($comment)).as(word)) .groupBy(word) .count() .orderBy($count.desc) .limit(200) wordDF.write.mode(SaveMode.Overwrite) .option(header, true) .csv(/output/word_num)explode把每条评论拆成多行每行一个词然后再groupBy统计词频。udf包了一层WordUtil.analyze这样Spark会在每个Executor上并行执行分词。limit(200)只保留前200个高频词避免词云图出现太多低频噪音。write.option(header, true)输出带表头的CSV下游Pandas或Seaborn读取更方便。3.4 集群环境搭建与提交命令本地单机跑这个作业不需要集群直接spark-submit即可。如果想在集群上跑需要至少一个Master和两个Worker。提交命令的参数往往比代码本身更容易出错我习惯这样写spark-submit \ --master yarn \ --deploy-mode client \ --driver-memory 2g \ --executor-memory 2g \ --executor-cores 2 \ --num-executors 4 \ --class com.douban.etl.MovieAnalysis \ douban-analysis.jar \ hdfs:///data/douban/comments.csv \ hdfs:///warehouse/douban/cleaned--deploy-mode client表示Driver运行在提交机本地方便看日志如果跑离线任务可以改成cluster并提前把日志写到指定目录。--executor-cores 2和--num-executors 4的组合需要和YARN队列资源匹配资源不够时容易一直卡在ACCEPTED状态。--class路径要和Scala文件的包名保持一致否则报ClassNotFoundException。提交后输出目录下会生成_SUCCESS文件和.part-r-00000.csv切片这就是压缩包里出现_SUCCESS和.part-r-00000.crc的原因。4. 结果可视化用Seaborn和WordCloud还原Spark输出Spark计算完的统计结果还停留在CSV或Parquet文件里这个项目的最后一段是用Python把结果读回来画图。可视化部分使用Matplotlib、Seaborn和WordCloud这些都是当前数据岗最常用的库。需要注意的是Spark输出目录会生成多个part-文件直接读取目录名比读单个文件更合适。4.1 读取Spark统计结果到PandasSpark写到HDFS或本地目录后Pandas可以直接读取目录下所有part-文件合并成一个DataFrame。假设结果目录是/output/lv_numimport pandas as pd import glob lv_files glob.glob(/output/lv_num/part-*.csv) df_list [] for f in lv_files: df_list.append(pd.read_csv(f, encodingutf-8)) lv_df pd.concat(df_list, ignore_indexTrue) print(lv_df.head())glob匹配part-开头的所有CSV切片再用pd.concat合并。因为Spark切分文件时不会保证每个文件的数据都完整所以这一步必须做。我在实际项目中还会把多个维度的输出合并成一个宽表方便后续一次性做对比分析比如把评分分布和年份趋势join起来。4.2 评分分布和年份趋势图评分等级分布适合用柱状图年份趋势适合用折线图。下面的代码分别画出这两类图import seaborn as sns import matplotlib.pyplot as plt sns.set_style(whitegrid) plt.rcParams[font.sans-serif] [SimHei] plt.rcParams[axes.unicode_minus] False lv_order [5, 4, 3, 2, 1] sns.barplot(datalv_df, xrating_score, ycount, orderlv_order, paletteviridis) plt.title(豆瓣短评评分等级分布) plt.xlabel(评分) plt.ylabel(评论数量) plt.savefig(rating_dist.png, dpi150, bbox_inchestight)这段代码先设置Seaborn背景和Matplotlib中文字体。orderlv_order必须显式指定顺序否则Pandas可能会按数字排序。paletteviridis是色盲友好的配色比默认柱状图更专业。年份趋势图同理year_df pd.read_csv(/output/year_num/part-*.csv, encodingutf-8) year_df[year] year_df[year].astype(int) year_df year_df.sort_values(year) sns.lineplot(datayear_df, xyear, ycount, markero) plt.title(评论量年份趋势) plt.xlabel(上映年份) plt.ylabel(评论数量) plt.savefig(year_trend.png, dpi150, bbox_inchestight)这里有一个常见坑CSV里的年份字段可能是字符串画图前必须转成int并排序否则折线图会出现折返。4.3 评论词云绘制词云需要把Spark输出的词频表转换成字典再交给WordCloud生成。核心代码from wordcloud import WordCloud word_df pd.read_csv(/output/word_num/part-*.csv) freq_dict dict(zip(word_df[word], word_df[count])) wc WordCloud( font_path/usr/share/fonts/opentype/noto/NotoSansCJK-Regular.ttc, width1200, height800, background_colorwhite, max_words200 ) wc.generate_from_frequencies(freq_dict) plt.imshow(wc, interpolationbilinear) plt.axis(off) plt.savefig(wordcloud.png, dpi150, bbox_inchestight)font_path必须指定中文字体否则中文全变成方块。Windows环境用C:/Windows/Fonts/simhei.ttfLinux环境用Noto CJK字体。generate_from_frequencies接收字典键是词值是词频。max_words200和Spark侧的limit(200)保持一致。生成词云后可以再做一张按词频排序的条形图对高频词做更精确的排名。5. 排错与进阶处理中文乱码、Spark内存问题并扩展分析维度拿到这个源码包后很多人会被里面的.class文件吓到以为不能改了。其实压缩包里同时带有项目说明文档和源码.class只是编译产物。这一章先讲怎么从编译产物里反推代码结构再罗列几个我在复现过程中踩过的坑最后给一个把静态图表升级成可视化大屏的切入点。5.1 反推.class文件与重新编译源码包里出现.class说明项目可能不是用IntelliJ直接跑而是通过sbt package或mvn package打包。想识别类的结构可以用JDK自带的javapjavap -p -c com/douban/util/WordUtil.class-p显示私有成员-c输出字节码。如果不想看字节码直接看同一目录下有没有对应的.scala源文件或者从README里的构建命令反推。我一般会先执行sbt clean package确认项目是否依赖外部仓库如果依赖下载失败就把Maven仓库切到阿里云镜像。5.2 常见运行错误和排查方法运行现象可能原因处理建议Spark提交后报ClassNotFoundException--class路径和包名不一致检查package com.douban.etl改成com.douban.etl.MovieAnalysis输出文件中文乱码Spark读取CSV时编码不是UTF-8统一用utf-8-sig写出读取时option(encoding,UTF-8)Executor OOMspark.sql.shuffle.partitions设置过小或数据倾斜调大分区数到48增加executor-memory爬虫抓取到一半全部403请求频率太高、请求头缺失加入随机延时每个movie_id之间间隔5-10秒Spark输出目录找不到_SUCCESS作业中途失败写入未提交查看YARN日志重点关注Exception in thread main其中中文乱码是最容易忽略的。爬虫落盘时用utf-8-sig写出Spark读取时再用UTF-8两边编码不一致时中文会变成字符。建议所有链路统一使用UTF-8只在最终给Excel人工预览时用utf-8-sig。5.3 进阶把Spark统计结果组装成可视化大屏这个项目的最终展示是静态图但论文答辩时通常会想要一个动态大屏。其实不用换技术栈只需要再走一步用Spark先把多维度聚合结果计算成一张大宽表后端接口再查这张宽表前端用ECharts或Pyecharts渲染。Spark侧可以这样处理val oil cleaned.groupBy(year, type) .pivot(rating_score) .count() .na.fill(0) oil.write.mode(SaveMode.Overwrite).parquet(/output/movie_olap)pivot会把评分等级展开成多列每行是某一年某一类型电影的评分分布。这样前端一次请求就能拿到一整行的统计信息不需要对每个图表单独调接口。可视化大屏的底层数据模型就是这么设计的Spark负责预聚合后端只做轻量透传。如果再往前走一步可以用Spark SQL的GROUPING SETS同时生成总计、按年份小计、按类型小计的汇总覆盖更多下钻维度。那才是把这个毕设真正从“跑通流程”升级成“可落地项目”的关键。本文还有配套的精品资源点击获取