拓冰建站拓冰建站
首页 / 资讯中心 / 正文

基于Hadoop的豆瓣电影数据分析:爬虫、MapReduce与可视化实践

简介一套基于Hadoop的豆瓣电影数据分析系统源码与技术文档适合计算机科学、数据科学及大数据专业的高年级本科生与研究生可用于课程设计、毕业设计或学术研究参考。项目全程覆盖了从数据采集、分布式存储、并行计算到结果可视化的完整工程流程爬虫环节利用Python语言中的Requests请求库和BeautifulSoup解析库定向抓取豆瓣高分影片信息随后将数据交由Hadoop平台完成存储与MapReduce分析任务并输出多种维度的统计结果。压缩包内含32个文件总大小约2.04MB主要包含Java编写的MapReduce计算程序、Python爬虫及可视化脚本、txt格式的结果数据、png格式的统计图表、md格式的技术文档以及若干备份文件结构清晰、便于按模块检索。目前已有43人学习浏览项目代码经过完整功能测试在毕业设计答辩中得到认可可帮助学习者快速体会Hadoop在大数据场景下的真实用法。借助提供的源码、文档与结果数据既可深入研读各模块的关联逻辑也可针对评分、类型、地区等分析维度进行二次扩展是理解大数据技术栈的实用参考。资源仅作为学习交流用途请勿用于商业传播。1. 说个可能和直觉相反的结论这套基于 Hadoop 的豆瓣电影大数据分析系统看上去文件一大堆真正的主线只有三个——用 Python 把豆瓣高分影片抓下来用 MapReduce 在 Hadoop 上按国家、类型、年份、演员、导演做维度统计再用 matplotlib 把统计结果画成能放进论文的图。作为一个毕业设计级别的完整源码包它比网上很多只有一个 README 的“项目”实在得多所有环节都能跑通而且每个环节恰好踩在 Hadoop 入门者最常犯错的几个点上。适合两类人一类是正在做课程设计或毕业设计、需要一份完整参考实现的学生另一类是已经工作、想用最快速度把 Mapper/Reducer 运行模型捡起来的人。把这个项目读懂比刷二十道 Hadoop 面试题更能理解 Shuffle 和 Combiner 的实际行为。2. 代码包结构从爬虫到 Hadoop 再到可视化的数据链路2.1 源码包里的三类文件分别对应哪个环节解压 zip 后src 下大体是三个 Java 目录对应 MapReduce 各类、python 目录爬虫与可视化和一堆*.txt结果文件。第一次看这个包不要急着读代码先把spider.py生成movies.txt、Movies.java 等生成type.txt/country.txt/time.txt/actors.txt/directors.txt、actor.py 与 sandiantu.py 消费这些 txt 的链路理顺。用一张表概括最直观文件 / 目录所处环节作用spider.py数据采集从豆瓣抓取高分影片字段清洗后写入movies.txtMovies.javaHadoop 计算统计影片类型分布输出type.txtCountry.javaHadoop 计算统计制片国家/地区数量输出country.txtActors.javaHadoop 计算统计演员出演次数输出actors.txtDirector.javaHadoop 计算统计导演作品数量输出directors.txtLong.javaHadoop 计算对时长/评分等数值字段做区间统计输出time.txtactor.py/directors.py/sandiantu.py可视化读 txt 结果绘制条形图、散点图这种结构在答辩时非常好讲评委问数据怎么流转的对着这张表说三句话就清楚。值得注意的是Long.java不是 Hadoop 里的LongWritable而是项目自己的一个类名第一次读源码时容易被这个名字带偏它的实际定位在第五章单独讲。2.2 数据怎么组织五张统计输出表的字段设计movies.txt是整条链路的中间产物每一行一条影片记录字段之间用制表符分隔。不同提交版本的字段顺序可能不一样动手前先执行head -5 movies.txt确认当前版本的字段顺序不要拿着别人代码里的下标硬套。我拆这套源码时遇到最常见的顺序是年份、片名、类型、国家、导演、演员、评分、时长即year\tname\ttype\tcountry\tdirector\tactors\trating\tduration。设计上有一个值得借鉴的点MapReduce 不直接产出一张统一的大宽表而是按维度分开输出五个小文件。这样做的直接好处是后续 Python 可视化时不需要反复解析同一行里的多个字段。缺点是有同学会问“为什么不一次性统计多个维度”理由其实很实际——单个 Mapper 解析一行后如果同时往多个 Context 写会导致输出结构混乱分开跑五个 Job 虽冗余但每个 Job 只干一件事排错和答辩都容易讲。2.3 先上传 HDFS再逐条提交 Job跑 Hadoop 之前把movies.txt放进 HDFS。伪分布式环境里常见路径是/user/root/douban/input上传后立刻用-cat验证hdfs dfs -mkdir -p /user/root/douban/input hdfs dfs -put movies.txt /user/root/douban/input/ hdfs dfs -cat /user/root/douban/input/movies.txt | head -5这里有几个参数细节值得说明-mkdir -p会递归创建父目录避免目录不存在时-put直接报错-cat后面接head -5是为了只输出五行防止数据量大时终端刷屏。如果上传后发现文件头带有 BOM 或者第一行是表头不要慌两种常见处理方式表头用sed -i 1d movies.txt直接删掉BOM 则参照第三章的字段清洗逻辑处理。文件就位后按第五章的方式逐个hadoop jar提交。目录结构保持input与多个output分离是为了避免 MapReduce 输出目录已存在时报FileAlreadyExistsException。3. 用 Requests BeautifulSoup 抓取豆瓣高分影片字段3.1 列表页拿索引详情页拿完整字段爬虫部分用了 Requests 拿页面、BeautifulSoup 解析这也是当下 Python 课程里最主流的组合。豆瓣的列表页只展示片名、评分等少量信息导演、演员和更多元数据必须进详情页拿所以spider.py的爬取策略是两段式先遍历列表页收集详情页链接再逐条请求详情页抽取字段。核心代码大致形如import time import requests from bs4 import BeautifulSoup headers { User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) } def fetch_detail(url): resp requests.get(url, headersheaders, timeout10) resp.encoding resp.apparent_encoding soup BeautifulSoup(resp.text, html.parser) name soup.select_one(h1 span).get_text(stripTrue) rating soup.select_one(strong.rating_num).get_text(stripTrue) return {name: name, rating: rating} for page in range(10): list_url fhttps://movie.douban.com/top250?start{page * 25} resp requests.get(list_url, headersheaders, timeout10) soup BeautifulSoup(resp.text, html.parser) for item in soup.select(.hd a): time.sleep(1) # 控制请求频率避免触发反爬 data fetch_detail(item[href]) print(data)代码里的select_one是 BeautifulSoup 的 CSS 选择器方法返回第一个匹配节点.hd a这个选择器在豆瓣列表页中能同时拿到片名和详情链接比逐层find_all更简洁。resp.encoding resp.apparent_encoding这一行很关键豆瓣页面在部分环境下会被错误识别为 ISO-8859-1导致中文乱码显式设置编码可以避免中文片名变成乱码再写进movies.txt。爬虫的节奏控制在每页之间和每部详情之间各睡 1 秒这个参数在课程设计里够用不要为了追求速度盲目调小间隔项目是演示性质稳定比速度重要。3.2 字段清洗全角符号、空值和多值字段抓下来的字段不能直接进 HDFS因为 MapReduce 默认按制表符切分字段而豆瓣页面里的类型、国家、演员字段常常是“剧情 / 爱情 / 冒险”这种用斜杠分隔的多值。常见的清洗策略有两条一是统一把分隔符替换成|二是只保留主值。我一般建议课程项目保留多值原因很实际——Actors.java要看哪几位演员出现次数最多如果把多值砍成单值就失去了分析意义。相应的清洗函数一般写成def clean_text(value): value value.replace( / , |) value value.replace( , ).strip() return value这里strip()去掉首尾空白replace( / , |)把列表页常见的间隔符统一成管道符到 MapReduce 阶段再用split(\\|)展开。注意全角空格在豆瓣页面上并不少见strip()默认不会去掉它所以需要显式替换掉再处理。空值字段在写文件时补N/A不要让两个制表符连在一起否则后面 Java 侧split(\\t)会得到空字符串导致数组越界。3.3 入库前的数据校验清洗结束后顺手做一次体检再落盘能省掉后面好几次排错。检查维度一般是三个行数是否等于预期影片数、每行字段数是否一致、评分列能否转成 float。第一项直接wc -l movies.txt第二项可以用一行 awk 统计异常行awk -F \t NF 8 {ok} NF ! 8 {print NR: NF} movies.txt-F \t指定制表符为分隔符NF 8判断当前行是否正好八个字段不符合时把行号打出来。字段数不一致最常见的原因是简介或演员列表里残留了换行符。课程项目不必做得太重但这一步能保证后续五个 MapReduce Job 不会因为脏数据反复失败。4. MapReduce 实现以 Country.java 为例拆 Mapper 与 Reducer4.1 Mapper 侧怎么拆字段最稳五类里面Country.java逻辑最简单统计每个国家的影片数量。它很适合当第一道门来理解 Map 端代码。下面按movies.txt的year\tname\ttype\tcountry\tdirector\tactors\trating\tduration顺序来写注意第四个字段才是国家public class CountryMapper extends MapperObject, Text, Text, IntWritable { private Text outKey new Text(); private IntWritable one new IntWritable(1); Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); String[] fields line.split(\\t); if (fields.length 4) { return; } String country fields[3].trim(); if (country.isEmpty()) { return; // 空国家不参与统计 } outKey.set(country); context.write(outKey, one); } }line.split(\\t)里\\t在 Java 字符串中是转义后的制表符这一步把一行记录拆成字段数组然后取第四个字段作为国家。这里有两个容易被忽视的细节第一split默认不会过滤前导空字符串如果行首有意外空格国家名称会变成带空格的脏 key所以先trim()再判断空值。第二fields.length 4的防御性检查不能省上一章清洗后如果还有意外空行不加这个判断就会出现ArrayIndexOutOfBoundsException而 Hadoop 日志对这种异常只会给出很长的堆栈排错效率很低。4.2 Reducer 求和与结果输出Reducer 侧逻辑不复杂把所有相同 key 的 value 累加即可。这个类完全可以同时承担 Combiner 的角色因为sum操作满足交换律和结合律public class CountryReducer extends ReducerText, IntWritable, Text, IntWritable { private IntWritable result new IntWritable(); Override protected void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int sum 0; for (IntWritable val : values) { sum val.get(); } result.set(sum); String[] parts key.toString().split(\\|); for (String part : parts) { context.write(new Text(part), result); } } }Reducer 里多做了一个很实际的处理如果 Mapper 输出的国家是“美国|英国”这种用管道符拼接的多值就在这里拆开再分别输出这样后续读country.txt时每个国家单独成行不会出现一个 key 下面挂着多国共享的计数。sum val.get()中的get()方法把IntWritable转成 int 做加法最后用set装回可序列化的 Writable 类型。整个类没有自定义数据类型也没有使用复杂 API便于课程设计阶段阅读。Driver 部分的提交参数值得留意我拆这个源码包时看到他们把五个作业的参数几乎重复写了五遍这在初期没问题但后面想调整并行度或 Combiner 时要改多个地方。一个更省事的写法是用一个辅助方法接收输出路径和主类名把setJarByClass、setMapperClass、setReducerClass、FileInputFormat.addInputPath、FileOutputFormat.setOutputPath统一收进函数代码量直接从每个类二十行缩到每个类八行。4.3 Combiner 的适用边界与 reducer 个数设置setCombinerClass是 MapReduce 课程里最常被误用的参数。一句话结论是只要合并操作满足交换律和结合律就能用同一个 Reducer 当 Combiner求和、最大最小值都满足求平均值不满足因为部分平均的再平均不等于全局平均。下表把这几个典型场景列出来方便答辩时直接答统计类型Combiner 可用原因求和、计数可用满足结合律求最大/最小可用部分最大值的最大仍是全量最大求平均值不可用部分平均的均值不等于全局平均TopN 全局排序不可全量combine 后无法判断全局前 NReducer 个数同样讲究。默认mapreduce.job.reduces1在数据量很小的课程项目中没问题输出只有一个part-r-00000读起来方便但当数据量增大、单 reduce 处理不过来时可以设置成job.setNumReduceTasks(2)以上。需要注意 reduce 数大于 1 后结果会分散到多个part-r-00000、part-r-00001后续可视化脚本读取时要用hadoop fs -getmerge合并处理这点到第六章还会验证。5. Long.java 与多 Job 串联数值字段的桶统计和结果排序5.1 桶统计map 阶段直接做区间映射Long.java在项目里负责的是time.txt这类数值型字段的统计和官方示例里的LongWritable没有任何关系。时长、评分这种连续数值在 MapReduce 里很少直接作为 key 输出——直接输出的结果是每个唯一值一个 key数据四散无法分析。更常见的做法是在 Map 端把连续值切成区间这一步叫分桶bucketing逻辑直接放在 map 阶段进行。一个典型实现片段如下public class LongMapper extends MapperObject, Text, Text, IntWritable { private IntWritable one new IntWritable(1); Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { String[] fields value.toString().split(\\t); if (fields.length 8) { return; } int minutes Integer.parseInt(fields[7].replaceAll([^0-9], )); String bucket; if (minutes 60) bucket 0-60; else if (minutes 90) bucket 60-90; else if (minutes 120) bucket 90-120; else if (minutes 150) bucket 120-150; else bucket 150; context.write(new Text(bucket), one); } }Integer.parseInt(fields[7].replaceAll([^0-9], ))是处理“128分钟”这类带中文单位字符串的常见技巧先用正则把所有非数字字符替换掉再转成 int。^在方括号内的含义是“取反”所以[^0-9]匹配所有非数字字符。这样的写法在字段里混入全角空格的场景下比单纯trim()更保险。分桶时最需要注意的是区间边界 60会把恰好 60 分钟的影片归入下一桶这个边界选择没有对错但要在文档里写明白否则答辩时被问起会显得不严谨。5.2 多 Job 串联的两种常见编排方式五个统计作业不能只提交一次就全部跑完因为 Actor 和 Country 各自要独立的 Mapper/Reducer。源码包里最常见的是在 Linux 命令行写一个按顺序执行的脚本每次执行一个hadoop jar。伪分布式环境下的完整提交命令长这样hadoop jar douban.jar com.douban.count.CountryJob /user/root/douban/input/movies.txt /user/root/douban/output/country hadoop jar douban.jar com.douban.stat.LongJob /user/root/douban/input/movies.txt /user/root/douban/output/time hdfs dfs -rm -r /user/root/douban/output/country注意hadoop jar的第二个参数是包含 main 方法的类名同一个 jar 包里每个作业一个入口这是常见打包方式输出目录必须不存在于 HDFS 上否则 Job 启动阶段直接报文件已存在。所以我习惯在脚本头部放一个-rm -r保证重复迭代时不会因为旧目录卡住。比命令行更稳的编排方式是写一个 ChainDriver在 main 方法里依次调用job.waitForCompletion(true)后面 Job 的输入路径直接引用前面 Job 的输出路径省去手工清理的麻烦。这个方式适合稍微进阶的用法比如统计完国家维度后紧接着做一次按数量倒序的 TopN 排序。5.3 TopN 排序与结果目录隔离如果需要绝对 TopN就不能依赖 reduce 端的默认字典序排序。Text类型的 key 按字典序排数值结果是字符串会出现 100 排在 20 前面的问题。常见的规范化方案是把数字 pad 到等宽字符串例如String.format(%06d, count)然后利用默认排序取得正确顺序或者在 Reducer 里用 TreeMap 累积 TopN最后在cleanup阶段输出。课程项目我更推荐前者因为它不需要自定义RawComparator可以在只改一行的前提下达到目标outKey.set(String.format(%06d, sum) _ country); context.write(outKey, result);String.format(%06d, sum)把数值 count 输出成至少六位的定宽字符串前导不足补零再拼上国家名作为一个复合 key。这样 reduce 输出在 HDFS 上会先按定宽数字排序后面做可视化时用sort -t_ -k1或直接读入 Python 后按完整字符串解析都能得到正确顺序。要提醒的是这种复合 key 只适合演示和验证生产环境里更严谨的 TopN 需要二次排序或 TreeMap但本项目的定位决定了第一种方案已经足够。6. 把统计结果拖回本地用 matplotlib 验证分析结论6.1 读取part-r-00000再绘图的两种姿势MapReduce 的输出在 HDFS 上是一个part-r-00000文件可视化脚本需要先把它弄回本地。两种常见姿势直接hadoop fs -cat重定向到本地 txt或者在脚本里调subprocess读取。务实一点课程项目用第一种即可hdfs dfs -getmerge /user/root/douban/output/country ./country_result.txt-getmerge会把该目录下的所有part-*合并成一个本地文件避免多个 reducer 输出时分别下载。如果前面设置了多个 reducer这个过程就非常关键否则actor.py读不到完整数据。合并后用head -5 country_result.txt先看一眼内容格式再写绘图逻辑这是整个项目里最快定位上游错误的方法——经常有同学上来就画图结果图是空白的回来查才发现 HDFS 上的结果根本没落下来。6.2 散点图与条形图的中文显示处理actor.py、directors.py的图都属于横向条形图sandiantu.py是年份和评分组成的散点图。matplotlib 第一次画中文几乎必现方框原因是默认字体不含中文字形。常见修法是一行代码指定字体族plt.rcParams[font.sans-serif] [SimHei]更保险的做法是引入font_manager定位系统字体文件适合 Linux 服务器上没有 SimHei 的环境。条形图的典型逻辑如下import matplotlib.pyplot as plt plt.rcParams[font.sans-serif] [SimHei] plt.rcParams[axes.unicode_minus] False pairs [] for line in open(country_result.txt, encodingutf-8): parts line.rstrip(\n).split(\t) if len(parts) 2: pairs.append((parts[0], int(parts[1]))) pairs.sort(keylambda x: x[1], reverseTrue) top pairs[:10] plt.figure(figsize(10, 6)) plt.bar([x[0] for x in top], [x[1] for x in top], colorsteelblue) plt.xticks(rotation45) plt.tight_layout() plt.savefig(country_top10.png, dpi150)split(\t)按制表符切列int(parts[1])把计数字段转成数值reverseTrue决定降序排列取[:10]是只画前十个国家。figsize(10, 6)控制画布大小rotation45防止国家名重叠dpi150保证导出图片在论文里足够清晰。验证一张图是否合理最直接的方法是把图画完后再对照 HDFS 上country.txt的前五行数字看图形中的 top 是否和数值线性排序一致——如果图形最高的国家不是数值最大的那个说明数据加载字段位置有误优先检查split下标而不是去改绘图代码。本文还有配套的精品资源点击获取
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门