豆瓣电影爬虫与Spark数据分析可视化实战

发布时间:2026/9/13 20:11:04
豆瓣电影爬虫与Spark数据分析可视化实战 简介一份基于Python和Spark的豆瓣电影爬虫与数据分析可视化系统适合毕业设计、期末大作业和课程设计场景。项目完整覆盖从网页爬虫、数据清洗、Spark批量统计到前端可视化展示的整个流程面向想快速搭建大数据分析应用的Python和Spark初学者。资源包共包含241个文件压缩包大小约5.65MB主体以XML配置和Spark中间结果文件为主同时配有Java逻辑代码、HTML/CSS/JS页面、SQL数据库脚本以及Python爬虫与Spark处理脚本目录结构清晰便于按模块查阅。从处理输出可以看出系统已实现影评词频、电影类型、评分等级、评论数、年份分布等多个维度的统计与可视化能直接用于项目演示或二次开发。当前已有248人学习下载。代码注释详细新手也能看懂附带数据库文件简单部署即可运行作为作者手打98分的毕业设计兼具完整性和实用价值尤其适合作为课程设计、期末作业和毕设参考。1. 为什么不直接爬完就展示Spark在这套豆瓣电影系统里的位置如果只是把豆瓣电影Top250爬下来存到Excel那这项目最多算入门作业。真正的分水岭在于你能否回答近五年哪个类型的电影评分中位数最高、评论里出现频率最高的形容词是哪几个这类问题。这套基于PythonSpark的方案爬虫只负责收集真正的分析落在了Spark的分布式计算上。项目里出现的WordNum.class、TypeNum.class这类文件并不是Python源码而是Spark或MapReduce作业编译后的class。也就是说爬虫脚本和批处理分析代码是两套技术栈中间通过MySQL衔接。这种混合结构恰恰是真实数据工程里最常见的样子采集用Python批量计算用Spark。你拿到目录里除了Python文件还有数据库文件(.sql)和Spark作业源码意味着不需要从头建表抓数据导入数据库就能开始跑完整链路。适合期末大作业和课程设计作为参考也适合想搞清楚爬虫离线分析可视化三者怎么协作的人。2. 爬虫设计requests、代理池与反爬降级的取舍2.1 目标URL与请求头设计豆瓣电影需要抓两类页面列表页和详情页。列表页拿到电影ID详情页补全评分、评价人数、类型、语言这些字段。用requests而不是scrapy主要原因是项目规模不大requests配合Session足够控制请求频率也能更精细地处理异常。import requests import time from fake_useragent import UserAgent ua UserAgent() HEADERS { Accept-Language: zh-CN,zh;q0.9,en;q0.8, } def fetch_movie_detail(movie_id: int) - dict: url fhttps://movie.douban.com/subject/{movie_id}/ try: resp requests.get( url, headers{User-Agent: ua.random, **HEADERS}, timeout3 ) resp.raise_for_status() resp.encoding utf-8 return parse_detail(resp.text) except requests.RequestException as e: print(f[warn] {movie_id} failed: {e}) return {}这里的timeout3是为了防止一个电影详情页卡死整个爬虫线程resp.encoding utf-8强制编码避免豆瓣页面里偶尔的编码误判。ua.random每次请求随机一个User-Agent降低被规则拒绝的概率。解析函数parse_detail一般用BeautifulSoup按meta标签提取数据比如from bs4 import BeautifulSoup def parse_detail(html: str) - dict: soup BeautifulSoup(html, html.parser) info soup.select_one(div#info) if not info: return {} return { title: soup.select_one(h1 span).get_text(stripTrue), year: extract_year(info), rating: extract_rating(soup), }这段代码没有做全量字段解析而是保留extract_year和extract_rating两个辅助函数的位置。实际做毕设时你可以在里面分别用正则re.search(r(\d{4}), text)提取年份用propertyv:average取评分比一次写死更可维护。2.2 代理池和限速策略豆瓣对单IP的请求频率非常敏感连续请求超过几十次就会弹出验证码。常见做法是维护一个代理池每次请求前从池子里随机取一个代理。这里有一个很小的调度模块import random PROXIES_POOL [ http://user:passproxy1:8080, http://user:passproxy2:8080, ] def get_proxy(): return {http: random.choice(PROXIES_POOL), https: random.choice(PROXIES_POOL)} def crawl_with_retry(movie_id: int, retries: int 3): for attempt in range(retries): resp fetch_movie_detail(movie_id, proxiesget_proxy()) if resp: return resp time.sleep(2 attempt * 2) return {}重试间隔按2 attempt * 2递增第二次重试等4秒第三次等6秒尽量避开瞬时封禁。并发设计方面不要为了赶进度直接开50个线程否则封禁概率会指数上升。一般用ThreadPoolExecutor(max_workers5)加Semaphore控制并发量。真正追求速度的分布式爬虫通常会让代理池每分钟轮换几百个IP但那是有商业代理支撑的方案毕设里未必有必要。2.3 数据落库的表结构爬下来的数据最终要落到MySQL里Spark再读库分析。表结构不能设计成一把抓否则后面分析会很痛苦。常见设计是movie和movie_comment各一张表字段类型说明idint豆瓣电影编号titlevarchar(255)电影标题yearint上映年份ratingdecimal(3,1)豆瓣评分genresvarchar(255)类型逗号分隔countriesvarchar(255)制片国家/地区languagesvarchar(255)语言comment_countint评论数create_timedatetime入库时间评论表至少要有movie_id、comment_text、comment_time三个字段。这里的genres字段存成逗号分隔字符串会让后续Spark explode操作很顺手后面第四章会看到具体用法。主键就看电影编号如果重复爬取用INSERT ... ON DUPLICATE KEY UPDATE更新评分和评论数保证数据幂等。2.4 真实环境中的反爬变化豆瓣的反爬策略这些年变严格了很多。改User-Agent只能防住最简单的拦截真正要稳定抓取还得关注Cookie里的bid字段。很多爬虫第一次访问会先抓首页拿到Cookie再把Cookie塞进Session看起来更像真实用户。我一般会在requests.Session里维护同一个Cookie并把单次请求间隔调成1到2秒。如果发现页面返回登录跳转就说明当前代理已经被标记需要换IP。更保险的方式是用豆瓣的公开API比如搜索接口但那种接口也受频率限制和爬页面没有本质区别。这个项目里的爬虫部分重点在演示链路而不是挑战高强度反爬所以放在本地小批量跑就完全够用。3. Spark作业读取MySQL从原始数据到可分析的DataFrame3.1 为什么是Spark而不是Pandas几万条电影记录Pandas完全能处理但毕设的评分要求里通常会强调分布式计算。Spark的价值在于当评论数据膨胀到几百万条时同样的聚合逻辑可以在多台机器上并行。更重要的是Spark SQL可以直接对MySQL建临时视图用纯SQL跑聚合这对不熟悉DataFrame API的人非常友好。另一个原因是项目中需要做中文分词统计Pandas要自己写并行很容易写出bugSpark的map、reduceByKey操作天然适合文本统计。如果你本机没有Spark集群常见做法是先搭standalone模式至少一个master一个worker本地用小分区数模拟分布式执行。3.2 JDBC读取与分区设置Spark读取MySQL不是把整表抽到内存而是通过JDBC连接器按分区读取。如果不设置分区Spark会对全表产生一个Task数据量大时直接OOM。设置分区字段和范围能让多个Task并行拉数据。val jdbcDF spark.read .format(jdbc) .option(url, jdbc:mysql://localhost:3306/douban) .option(dbtable, movie) .option(user, root) .option(password, your_password) .option(driver, com.mysql.cj.jdbc.Driver) .option(partitionColumn, id) .option(lowerBound, 1) .option(upperBound, 1000000) .option(numPartitions, 8) .load()partitionColumn必须是一个数值型字段一般用自增主键。Spark会把范围从lowerBound到upperBound切成numPartitions段每段对应一个WHERE id ? AND id ?的子查询。注意numPartitions不要大于数据库最大连接数否则每个分区建一个连接会把MySQL连接数打满。这里还有几个容易踩的坑MySQL驱动要放进Spark的jars目录serverTimezoneAsia/Shanghai要加到url里否则报日期时区错误。dbtable可以写成(SELECT id,title,rating,genres,comment_count FROM movie) tmp先投影再拉取减少网络传输。3.3 清洗和类型转换原始表里year可能是字符串rating也可能包含缺失值。Spark的强类型API要求我们在计算前把这些字段转成Int和Double。下面的代码处理了三种异常import org.apache.spark.sql.functions._ val cleanDF jdbcDF .filter(col(title).isNotNull) .withColumn(year, col(year).cast(int)) .withColumn(rating, col(rating).cast(double)) .withColumn(comment_count, col(comment_count).cast(int)) .na.fill(Map(rating - 0.0, comment_count - 0))cast(int)会把2009这种字符串转成数字但也会把None转成null。所以先用isNotNull过滤掉关键字段为空的记录。na.fill用来补评分和评论数为缺失的记录后面按评分排序时就不会出现null跑到最前面。这个顺序很关键先过滤再转换最后填充如果顺序错了填充完的类型又会被cast成null。3.4 中文分词与WordUtil项目中出现的WordUtil.class它的作用是把评论按中文分词拆开再统计每个词的出现次数。如果你用的是Spark常见实现是引入HanLP或者结巴分词。但HanLP的词典文件比较大在Spark集群上每个executor都加载一次会给内存增加负担。常见的做法是把分词封装成一个map函数让Spark对每个分区执行val seg new Segment() val commentDF spark.read.table(comment).select(comment_text) commentDF.map(row { val text row.getString(0) val words seg.segment(text).asScala.map(_.word.trim).filter(_.length 1).toList (words, 1) }).rdd.flatMap(pair pair._1.map(word (word, pair._2))) .reduceByKey(_ _)这里的_ .map和reduceByKey就是RDD里最常见的wordcount逻辑。有一个细节是分词词典需要把电影这类停用词过滤掉否则统计结果全是的、是这种无意义词。我一般会加载一个停用词表在filter的时候顺手排除。还有一次我在集群上跑发现输出里大量乱码后来确定是新节点没有安装系统级别的中文字体Spark executor打印日志时编码异常这不是业务bug而是环境问题。4. 指标计算评分分布、类型热度和年份趋势的实现4.1 评分区间统计评分分布通常是按9分以上、8-9、7-8、6-7、6以下分组统计每组电影数量。这种分桶逻辑用Spark SQL最清晰cleanDF.createOrReplaceTempView(movie) spark.sql( SELECT CASE WHEN rating 9 THEN 9 WHEN rating 8 THEN 8-9 WHEN rating 7 THEN 7-8 WHEN rating 6 THEN 6-7 ELSE 6- END AS rating_bucket, COUNT(*) AS cnt FROM movie GROUP BY rating_bucket ORDER BY rating_bucket )CASE WHEN的求值顺序是从上到下所以每个区间的下界是前一个区间的上界不需要再写rating 9 AND rating 8这种冗余条件。GROUP BY rating_bucket其实是对表达式分组但SQL里可以直接用别名Spark SQL支持这个语法。4.2 类型榜单分析genres字段存的是剧情, 爱情, 战争这种逗号分隔字符串要做类型榜必须先把每个电影按类型拆开。这里需要用split配合explodeval typeDF cleanDF .withColumn(type, explode(split(col(genres), ,))) .groupBy(type) .agg( count(*).alias(movie_count), avg(rating).alias(avg_rating) ) .orderBy(desc(movie_count))split(col(genres), ,)把字符串转成数组explode再把数组的一行拆成多行。关键在于如果genres为空explode会直接把这行数据丢掉而实际业务上确实有电影没标注类型。想保留这些电影可以换成explode_outer。类型名的空白也需要处理我常在split前用regexp_replace(col(genres), \\s, )去掉所有空格否则每个类型前面会带空格导致统计不准确。4.3 年份维度聚合年份趋势适合看的是每年平均评分或者每年电影数量。年份字段是int后直接groupBy(year)即可。比较实用的一个做法是把年份按十年分桶看长周期趋势spark.sql( SELECT CONCAT(FLOOR(year/10)*10, s) AS decade, ROUND(AVG(rating), 2) AS avg_rating, COUNT(*) AS movie_count FROM movie WHERE year 1970 GROUP BY FLOOR(year/10) ORDER BY decade )这里用FLOOR(year/10)*10得到1980、1990这种整十年起点再用CONCAT拼一个s后缀。ROUND(AVG(rating), 2)控制小数位数因为评分均值经常出现9.200000这样的长尾。不加WHERE year 1970的话一些年份为0的脏数据会单独成组整个趋势图最左边会莫名其妙多一列。4.4 评论情感与词频评论数据量大直接用完整文本做情感分析较慢。毕设场景里更常见的是统计评论中的高频词用词频变化代表关注度。我们可以在分词的基础上把每个词的频次和它出现的影评关联形成评论关键词表。这里有一个很关键的Spark内存问题reduceByKey会在shuffle前先做本地的merge所以不会把全量数据放在一个executor里。但如果中间出现了很大的key列表比如几万个电影id每组的评论数巨大就要考虑提高spark.default.parallelism和executor内存。通常我在spark-submit里这样设置spark-submit \ --master local[4] \ --executor-memory 2g \ --conf spark.sql.shuffle.partitions10 \ --conf spark.default.parallelism10 \spark.sql.shuffle.partitions决定shuffle后的分区数默认200在小数据集上会产生非常多空任务浪费调度开销。调成10可以明显减少CPU空转。如果你的集群内存只有4Gexecutor-memory不要超过3g给系统留一点余量。5. 可视化ECharts大屏怎么接Spark结果5.1 后端接口设计Spark算出的结果最终要落到MySQL或导出成JSON可视化层不能直接读MySQL大表。常见做法是让Flask提供一个统计接口接口内部只做读库和返回JSON具体SQL已经在Spark端算好前端拿到的就是聚合后的数据。from flask import Flask, jsonify import MySQLdb app Flask(__name__) db MySQLdb.connect(hostlocalhost, userroot, passwd123456, dbdouban, charsetutf8mb4) app.route(/api/rating_bucket) def rating_bucket(): cursor db.cursor() cursor.execute(SELECT rating_bucket, cnt FROM rating_stat ORDER BY rating_bucket) rows cursor.fetchall() return jsonify({categories: [r[0] for r in rows], data: [r[1] for r in rows]}) if __name__ __main__: app.run(port5000)这里charsetutf8mb4必须显式声明否则中文在前端很容易变成问号。接口返回的categories和data分离前端ECharts可以直接映射到x轴和y轴不需要再做二次处理。注意如果统计数据是Spark离线算好后写回一张rating_stat表那么这份接口代码里查询的表名就是它。5.2 前端的图表配置ECharts里最常用的三种图是柱状图、饼图和折线图。柱状图适合评分区间饼图适合类型占比折线图适合年份趋势。下面是一个典型的柱状图配置fetch(/api/rating_bucket) .then(response response.json()) .then(res { const chart echarts.init(document.getElementById(ratingChart)); chart.setOption({ title: { text: 豆瓣电影评分区间分布 }, tooltip: {}, xAxis: { type: category, data: res.categories }, yAxis: { type: value }, series: [{ type: bar, data: res.data, itemStyle: { color: #5470c6 } }] }); });fetch是异步请求ECharts必须在数据返回后再init否则表格的宽高是0图会画不出来。category类型的xAxis要求数据是字符串数组如果后端返回数字还需要String()转换一次。这类细节往往是毕设答辩现场报错的原因。5.3 从静态JSON到动态刷新如果只是做一个只读大屏把数据写死在JavaScript里也能交差。但更贴近真实项目的做法是让前端每隔一段时间重新拉一次接口。function loadData() { fetch(/api/rating_bucket) .then(r r.json()) .then(res chartRef.current.setOption(/* ... */)); } setInterval(loadData, 60 * 1000);setInterval(loadData, 60 * 1000)代表每分钟刷新一次。这样做的好处是Spark每天跑一次批处理生成新统计结果大屏跟着自动变化不需要手动刷新页面。但要注意时间间隔不要太短否则后端会被反复查询压力打满。配合setInterval在组件卸载时一定要clearInterval否则浏览器内存会一直挂着定时器。5.4 可视化大屏适配注意事项用过可视化大屏的人都知道分辨率不同布局会乱。最稳妥的方案是使用rem方案把设计稿宽度设为1920px前端所有尺寸按比例换算。ECharts的图表容器用width: 100%; height: 100%再由父容器控制宽高。另一个细节是ECharts图表的grid配置尤其是饼图和折线图。很多新手把图放在一个Flex容器里发现图被拉伸变形。遇到这种情况先检查echarts.init时容器是否已经可见、有没有固定宽高。还有一个更容易忽略的点是如果页面里有多个图表必须为每个图表的div分配独立的ref或id不能让多个chart实例共享同一个节点。6. 部署与验证用自带数据库文件快速跑通别在Spark环境上卡住6.1 环境准备这个项目给你省了不少事数据库文件已经准备好不需要自己建表。但环境还是得装齐。# Python 依赖 pip install requests beautifulsoup4 flask pymysql # Spark 环境 docker pull bitnami/spark:3.5Python爬虫和可视化基本不依赖Spark本机安装真正需要Spark的只有分析模块。如果不想手动搭集群用Docker跑一个spark standalone是最快的路径。记住docker pull之后要把spark的bin目录挂进容器或者直接用容器内的spark-submit命令执行。6.2 数据库导入项目自带的数据库文件通常是.sql导入MySQLmysql -uroot -p -e CREATE DATABASE douban DEFAULT CHARACTER SET utf8mb4 mysql -uroot -p douban douban.sqlDEFAULT CHARACTER SET utf8mb4非常重要如果创建库的时候用了默认latin1后面中文数据插入时会报编码错误。导入完成后先执行一条SELECT COUNT(*) FROM movie;看数量是否与文档一致能提前发现.sql文件被截断的问题。6.3 最小复现流程要把整个项目跑起来建议按这个顺序操作导入数据库后先运行爬虫脚本抓几条新数据验证数据库能写入。运行Spark分析脚本把统计结果写入rating_stat等结果表。启动Flask后端访问/api/rating_bucket确认返回JSON。打开前端页面看图表是否正常渲染。不要让爬虫大规模运行因为服务器早已不是当年文档里的环境一夜之间把IP封了很正常。先limit 20条跑通链路再决定是否扩展。6.4 如何查看Spark Job执行情况最后一个实用技巧是学会看Spark的监控页面。无论本地local模式还是standalone集群提交任务后都会在http://localhost:4040暴露一个Web UI。docker exec -it spark-master spark-submit --master local[2] /opt/analyse.jar打开localhost:4040/jobs/可以看到每个Stage的输入记录数和耗时。如果某个Stage的Input远远大于预期说明读取MySQL时没有限制分区字段全表扫描了。根据这个信息可以反向调整JDBC的lowerBound和upperBound而不是盲目调大executor内存。这是整个项目里最值得花时间研究的地方理解了Spark的任务分配才算真的把爬虫之后的计算链路吃透。本文还有配套的精品资源点击获取