简介:一套基于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.txt |
Movies.java | Hadoop 计算 | 统计影片类型分布,输出type.txt |
Country.java | Hadoop 计算 | 统计制片国家/地区数量,输出country.txt |
Actors.java | Hadoop 计算 | 统计演员出演次数,输出actors.txt |
Director.java | Hadoop 计算 | 统计导演作品数量,输出directors.txt |
Long.java | Hadoop 计算 | 对时长/评分等数值字段做区间统计,输出time.txt |
actor.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, headers=headers, timeout=10) resp.encoding = resp.apparent_encoding soup = BeautifulSoup(resp.text, "html.parser") name = soup.select_one("h1 span").get_text(strip=True) rating = soup.select_one("strong.rating_num").get_text(strip=True) return {"name": name, "rating": rating} for page in range(10): list_url = f"https://movie.douban.com/top250?start={page * 25}" resp = requests.get(list_url, headers=headers, timeout=10) 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 与 Reducer
4.1 Mapper 侧怎么拆字段最稳
五类里面Country.java逻辑最简单:统计每个国家的影片数量。它很适合当第一道门来理解 Map 端代码。下面按movies.txt的year\tname\ttype\tcountry\tdirector\tactors\trating\tduration顺序来写,注意第四个字段才是国家:
public class CountryMapper extends Mapper<Object, 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 Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); @Override protected void reduce(Text key, Iterable<IntWritable> 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 后无法判断全局前 N |
Reducer 个数同样讲究。默认mapreduce.job.reduces=1在数据量很小的课程项目中没问题,输出只有一个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 Mapper<Object, 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", encoding="utf-8"): parts = line.rstrip("\n").split("\t") if len(parts) == 2: pairs.append((parts[0], int(parts[1]))) pairs.sort(key=lambda x: x[1], reverse=True) top = pairs[:10] plt.figure(figsize=(10, 6)) plt.bar([x[0] for x in top], [x[1] for x in top], color="steelblue") plt.xticks(rotation=45) plt.tight_layout() plt.savefig("country_top10.png", dpi=150)split("\t")按制表符切列,int(parts[1])把计数字段转成数值,reverse=True决定降序排列,取[:10]是只画前十个国家。figsize=(10, 6)控制画布大小,rotation=45防止国家名重叠,dpi=150保证导出图片在论文里足够清晰。验证一张图是否合理,最直接的方法是把图画完后再对照 HDFS 上country.txt的前五行数字,看图形中的 top 是否和数值线性排序一致——如果图形最高的国家不是数值最大的那个,说明数据加载字段位置有误,优先检查split下标,而不是去改绘图代码。
本文还有配套的精品资源,点击获取