资讯动态

基于Hadoop的国产电影数据分析与可视化项目全流程实战

发布时间:2026/9/9 21:46:17 来源:尧图企业网站定制
最近完整跑了一遍“基于Hadoop的国产电影数据分析与可视化”这个项目从数据采集、存储清洗、分布式计算到最终大屏展示整条链路都走通了。这个项目技术栈非常典型Scrapy负责采集Hadoop生态HDFS、Hive、MapReduce负责存储和计算最后用可视化把分析结果呈现出来。写这篇博客就是把整个过程做个复盘包括架构思路、关键代码、踩过的坑和选型理由给准备做类似大数据项目的朋友一份可以直接参考的实操笔记。1. 项目背景与整体设计思路1.1 为什么我要做一个电影数据分析项目国产电影市场这些年变化非常大每年上映几百部片子有票房爆款也有口碑佳作数据里藏着很多有意思的规律。我之前发现很多刚入门大数据的人手里有Hadoop、Hive这些技术但不知道拿什么数据来练手所以我决定自己造一个完整的数据分析项目爬取国产电影的基础信息、评分、票房、上映日期、题材类型等然后放到Hadoop生态里做存储与分析最后用可视化大屏把结论直观地呈现出来。这种项目的优势在于数据量适中既不会小到体现不出分布式计算的价值也不会大到单机跑不动的程度最适合用来学习Hadoop大数据技术栈的完整流程。而且电影数据对大部分人而言都有基本认知字段含义容易理解分析结果也能直接看出业务含义。1.2 技术栈选型Hadoop、Scrapy、ECharts 为什么搭在一起项目核心技术栈是“Scrapy Hadoop生态 可视化工具”的组合。选型逻辑要掰开来说Scrapy做爬虫是Python生态里最成熟的方案之一它对并发请求的处理做得比较好自带Item Pipeline、Downloader Middleware后续要扩展分布式爬虫也有成熟组件。Hadoop HDFS做存储层原因在于我预期数据量能达到几万到几十万条记录包含大量文本字段电影简介、导演、演员等这类半结构化数据放进HDFS可以保留原始格式后续处理弹性更大。Hive做数据仓库把HDFS上的结构化数据映射成表用SQL做清洗过滤比写MapReduce处理这一步省力得多。MapReduce做核心计算引擎作业提交到YARN上执行能体现Hadoop分布式计算的真实工作场景。可视化层选了ECharts扛大屏完全没有问题图表类型丰富交互流畅做国产电影可视化大屏属于经典方案。这套组合的另一层意义在于它覆盖了大数据的四个核心环节数据采集、数据存储、数据处理、数据展示全部是岗位面试里会被反复追问的知识点。1.3 整体架构与数据流向整个项目的数据流向可以从上往下分成五层采集层Scrapy爬虫从公开数据源抓取电影信息输出结构化CSV文件存储层CSV文件通过命令或API上传到HDFS指定目录数仓层Hive建立原始表和清洗表通过SQL完成去重、格式转换、字段补全计算层Hive的清洗结果作为输入使用MapReduce任务完成核心指标统计应用层统计结果导出为JSON或CSV前端ECharts读取数据渲染大屏这里有个设计上的关键决策哪些清洗逻辑放Hive哪些统计逻辑放MapReduce。我的原则是能上SQL的清洗绝不写MapReduce比如剔除空值、评分格式标准化、枚举字段映射这些工作Hive一条SQL就搞定了省时省力。而像“统计不同题材在不同年份区间的评分趋势”这种需要自己设计输出结构的复杂聚合再交给MapReduce来处理。2. 数据采集层Scrapy 爬虫的设计与实现2.1 爬虫目标与字段设计爬虫要解决的核心问题只有一个我要拿到哪些字段以及这些字段能否支撑后续的分析目标。我最终确定的字段清单包括电影名称、导演、主演、类型剧情/喜剧/动作等、制片地区国产电影这里主要用于区分大陆/港台、上映年份、片长、豆瓣评分、评分人数、票房数据如果能获取到。这里有个很关键的经验字段宁多勿少后续清洗可以删但爬完发现少字段就得重新跑一遍数据。我一开始没有采集评分人数后来想看“高票房是否等于高热度”时发现没有这种数据只能回头补爬浪费时间。2.2 Scrapy 核心代码实现标准的Scrapy项目包含items.py、spiders目录、pipelines.py、settings.py几个核心模块。给大家看一下我实现的核心代码结构。先定义Item# items.py import scrapy class MovieItem(scrapy.Item): title scrapy.Field() director scrapy.Field() actors scrapy.Field() genre scrapy.Field() region scrapy.Field() year scrapy.Field() duration scrapy.Field() rating scrapy.Field() rating_count scrapy.Field() box_office scrapy.Field() summary scrapy.Field()这里我把summary电影简介也存下来了虽然分析阶段不一定用到但作为数据仓库的原始数据层保留它有助于之后做自然语言处理方向的项目扩展。Spider的写法是典型的列表页详情页模式# spiders/movie_spider.py import scrapy from scrapy.selector import Selector from ..items import MovieItem class MovieSpider(scrapy.Spider): name movie_spider def start_requests(self): base_url https://example-movie-site.com/films?year{}page{} for year in range(2015, 2024): for page in range(1, 50): url base_url.format(year, page) yield scrapy.Request(urlurl, callbackself.parse_list, meta{year: year}) def parse_list(self, response): sel Selector(response) # 列表页里提取每个电影详情页的链接 detail_links sel.xpath(//div[classmovie-item]/a/href).extract() for link in detail_links: yield scrapy.Request(urlresponse.urljoin(link), callbackself.parse_detail) def parse_detail(self, response): sel Selector(response) item MovieItem() item[title] sel.xpath(//h1[propertyv:itemreviewed]/text()).get() item[director] sel.xpath(//a[relv:directedBy]/text()).get() item[actors] ,.join(sel.xpath(//a[relv:starring]/text()).extract()) item[genre] ,.join(sel.xpath(//span[propertyv:genre]/text()).extract()) item[year] response.meta.get(year) duration sel.xpath(//span[propertyv:runtime]/content).get() item[duration] duration.replace(分钟, ) if duration else item[rating] sel.xpath(//strong[propertyv:average]/text()).get() item[summary] sel.xpath(//span[propertyv:summary]/text()).get() yield item采集完成后通过Pipeline写入CSV# pipelines.py import csv class CsvPipeline: def open_spider(self, spider): self.file open(movies_raw.csv, w, newline, encodingutf-8-sig) self.writer csv.DictWriter(self.file, fieldnames[ title, director, actors, genre, region, year, duration, rating, rating_count, box_office, summary ]) self.writer.writeheader() def close_spider(self, spider): self.file.close() def process_item(self, item, spider): self.writer.writerow(item) return item有一个细节需要提醒CSV编码我用了UTF-8-SIG而不是单纯UTF-8这是为了避免用Excel或部分工具打开CSV时中文乱码。写文件不要每写入一条就close一次性能会很差正确做法是在open_spider时打开文件close_spider时统一关闭。2.3 反爬策略与数据质量保障爬虫阶段遇到过几个问题分享下解决方案。第一是请求频率控制。目标站点对频繁请求限制比较严格我的处理方法是在settings.py里设置DOWNLOAD_DELAY 1.5和CONCURRENT_REQUESTS 8让请求平缓发出效率和安全性之间取了平衡。第二是动态页面。有些数据源是动态渲染的Scrapy直接请求拿不到内容我配合Selenium做了一层渲染抓取。使用Selenium要注意下载chromedriver时匹配Chrome版本否则创建浏览器对象直接报错。第三是数据去重。爬虫重跑好几次CSV里会有重复行。我在爬虫端对title year生成指纹写入到一个set里做去重后续Hive清洗阶段还会再做一次分组去重作为兜底。数据质量是分析的基础这一步不要省。3. 数据入仓HDFS 存储与 Hive 数仓清洗3.1 为什么不用 MySQL非得上 HDFS Hive这是新手最容易问的问题其实核心在于数据的“姿态”。爬虫产出的原数据是半结构化文本有大量空值和脏数据直接导入MySQL再根据业务需求改表结构灵活性很差。而HDFS天然适合存放原始数据文件不管你是CSV、JSON还是Parquet先原样放进去再用Hive或Spark按需加工成不同的表结构这才是大数据架构常见的数据湖思路。另外从学习角度讲如果不把Hadoop引进来整个项目就退化成了一个普通的Python爬虫数据库项目在大数据面试里几乎没有亮点。HDFS和Hive是Hadoop技术栈里最核心的两个组件用这个项目把他们串起来面试官面到相关技术时你能拿出真实案例说服力完全不一样。3.2 Hive 建表与清洗 SQL原始CSV文件先上传到HDFShdfs dfs -mkdir -p /user/hive/data/movies hdfs dfs -put movies_raw.csv /user/hive/data/movies/然后创建Hive外部表。我特意使用外部表而不是内部表原因在于原始数据文件是独立于Hive的外部表删除表结构时不会把HDFS上的数据文件一起删掉安全性更好。CREATE EXTERNAL TABLE movies_raw ( title STRING, director STRING, actors STRING, genre STRING, region STRING, year INT, duration INT, rating DOUBLE, rating_count INT, box_office DOUBLE, summary STRING ) ROW FORMAT SERDE org.apache.hadoop.hive.serde2.OpenCSVSerde WITH SERDEPROPERTIES ( separatorChar ,, quoteChar \ ) STORED AS TEXTFILE LOCATION /user/hive/data/movies;这里用OpenCSVSerde而不是默认的LazySimpleSerDe是因为电影简介summary字段里可能包含逗号或引号如果直接用逗号分隔解析会被截断产生错乱。OpenCSVSerde能正确处理带引号的CSV字段这个坑当时排查了很久。接下来做清洗。清洗目标删除个别无法解析的分隔错乱记录空评分、空题材的记录剔除年份字段标准化为YYYY票房单位统一改成“亿”没有票房数据的置NULL清洗后的数据写入新表CREATE TABLE movies_clean AS SELECT title, director, actors, genre, region, year, duration, rating, rating_count, box_office FROM movies_raw WHERE title IS NOT NULL AND rating IS NOT NULL AND genre IS NOT NULL AND year 2000;这一条SQL执行起来非常快因为Hive会把任务转化成MapReduce并提交给YARN执行这也是我第一次直观感受到“写SQL就是在写MapReduce”这句话的含义。3.3 清洗后的数据规范说明清洗完成后数据质量检查非常关键我习惯跑几条简单查询确认-- 检查数据总量 SELECT COUNT(*) FROM movies_clean; -- 检查每年数据量分布 SELECT year, COUNT(*) AS cnt FROM movies_clean GROUP BY year ORDER BY year;这个环节有一类典型问题某一年份的电影数量明显偏少大概率是爬虫阶段该年份的列表页没翻完这就要回头补爬。宁可早期多花时间检查不要等分析阶段发现规律不对再回来查数据。4. MapReduce 分析从原始数据到洞察4.1 分析维度设计从哪些角度解读国产电影在动手写MapReduce之前先要明确业务问题否则就会变成“为了跑任务而跑任务”。我最终选了四个分析主题按年份统计国产电影产量与平均评分回答“国产电影是不是越拍越多、口碑是涨是跌”按题材统计电影数量与评分分布回答“什么题材的国产电影最容易出高分”高票房电影与高评分电影的重合度分析回答“叫好和叫座到底是不是一回事”导演产量与作品评分的关联分析找出哪些导演“产量又高口碑又稳”这几个维度覆盖了时间趋势、类型分布、市场关系、人物分析四个角度做出来的可视化图表也比较有层次感。4.2 Python Hadoop Streaming 实现 MapReduce我用的是Hadoop Streaming的方式用Python写Mapper和Reducer相比直接写Java代码调试成本低很多也方便不熟悉Java的读者直接复现。任务需求是统计“每年、每个题材”的电影数量和平均评分。Mapper端代码#!/usr/bin/env python # mapper.py import sys for line in sys.stdin: line line.strip() if not line: continue parts line.split(\t) if len(parts) 5: continue title, director, actors, genre, region, year parts[0], parts[1], parts[2], parts[3], parts[4], parts[5] rating parts[6] try: year int(year) rating float(rating) except ValueError: continue # genre 字段可能是 喜剧,爱情 这种多题材拆开分别统计 for g in genre.split(,): key f{year}\t{g.strip()} print(f{key}\t{rating}\t1)Reducer端代码#!/usr/bin/env python # reducer.py import sys current_key None current_sum 0.0 current_count 0 for line in sys.stdin: line line.strip() if not line: continue parts line.split(\t) if len(parts) ! 4: continue key f{parts[0]}\t{parts[1]} rating float(parts[2]) count int(parts[3]) if current_key key: current_sum rating current_count count else: if current_key: avg_rating current_sum / current_count print(f{current_key}\t{avg_rating:.2f}\t{current_count}) current_key key current_sum rating current_count count if current_key: avg_rating current_sum / current_count print(f{current_key}\t{avg_rating:.2f}\t{current_count})运行命令hadoop jar /path/to/hadoop-streaming-*.jar \ -input /user/hive/warehouse/movies_clean \ -output /user/analysis/genre_year_rating \ -mapper python3 mapper.py \ -reducer python3 reducer.py \ -file mapper.py \ -file reducer.py \ -numReduceTasks 4友情提示Streaming的方式虽然方便但Map阶段的输出默认以Tab键分隔key和value如果你处理的数据本身包含Tab就要特别注意。我在设计原始表的时候特意留意过这个问题电影简介字段如果有Tab会导致分割错乱清洗阶段要单独处理。4.3 分析任务调度与结果落盘MapReduce跑完的结果在HDFS上可以用命令直接查看hdfs dfs -cat /user/analysis/genre_year_rating/part-r-00000 | head -50为了让可视化前端能方便读取我通常会把结果从HDFS导出为本地JSON文件或者直接落成ClickHouse如果公司有的话。这个项目规模不大我把多个分析任务的结果合并成一个analysis_result.json文件放在Web服务目录下前端异步请求就可以了。这里还有一个重要的经验多个分析任务之间是有依赖关系的比如“年度产量统计”和“题材平均评分”互不依赖可以并行跑“导演产量与评分关联分析”依赖先统计导演产量再关联评分需要分两步。学习阶段可以手动一条条提交生产环境建议用Oozie或Apache Airflow做工作流调度不然任务多了容易乱。5. 可视化呈现搭建电影数据大屏5.1 可视化方案选型与数据接口设计可视化部分我没有选择直接在Jupyter里画图草草了事而是搭建了一套带交互的大屏页面。前端框架选了ECharts Vue 3后端用了一个轻量级的Flask服务来读取JSON数据。数据接口设计成一份聚合后的JSON文件更实用{ year_trend: [ {year: 2015, count: 36, avg_rating: 6.2}, {year: 2016, count: 42, avg_rating: 6.4} ], genre_rating: [ {genre: 喜剧, count: 85, avg_rating: 6.1}, {genre: 科幻, count: 22, avg_rating: 7.3} ], top10_movies: [ {title: 某电影, rating: 9.0, box_office: 56.7} ] }这种“一份JSON喂所有图表”的方式对中小型项目最省事可视化层和后端解耦大屏页面纯静态部署都可以。5.2 大屏核心图表实现大屏布局我采用经典的三列式结构左侧放“年度产量与评分趋势”折线图、中间放“题材分布”饼图和“Top10高分电影”柱状图、右侧放“票房与评分关系”散点图。ECharts的配置有点长但核心内容并不复杂给一个折线图的例子// yearTrend.js const chart echarts.init(document.getElementById(yearTrend)); const option { tooltip: { trigger: axis }, legend: { data: [电影数量, 平均评分] }, grid: { left: 10%, right: 10%, top: 15%, bottom: 10% }, xAxis: { type: category, data: years }, yAxis: [ { type: value, name: 电影数量, axisLabel: { formatter: {value} 部 } }, { type: value, name: 平均评分, min: 0, max: 10, axisLabel: { formatter: {value} 分 } } ], series: [ { name: 电影数量, type: bar, data: counts, itemStyle: { color: #3b82f6 } }, { name: 平均评分, type: line, yAxisIndex: 1, data: avgRatings, itemStyle: { color: #ef4444 } } ] }; chart.setOption(option);双Y轴图在这里特别好用因为电影数量和评分的量纲差异很大如果共用Y轴趋势会被“压扁”或“拉平”。这也是做可视化时一个容易忽略的细节。5.3 可视化常见的性能问题与优化如果你的数据量确实很大比如几十万条记录前端直接渲染所有数据点会非常卡。我的优化手段有三个聚合下推可视化展示的数据永远是在计算层已经聚合好的指标比如按年聚合、按题材聚合前端绝不处理原始明细数据。后端接口分页加载如果确实需要展示明细Table通过接口做分页每页最多50条。ECharts的sampling配置散点图数据量大时开启sampling: lttb在不影响趋势形状的前提下显著减少渲染点数。我这个项目聚合后的数据只有几百条渲染非常流畅但这套优化思路在真实业务场景里很重要。6. 踩坑实录与项目复盘6.1 Scrapy 爬虫阶段的坑爬虫阶段最大的坑是动态页面接口和数据源的选择。第一次写爬虫时直接对着某个站点的静态页面爬结果发现页面结构用了Ajax异步加载Scrapy请求返回的HTML里根本没有电影列表数据。后来换了方案在开发者工具里找到真实的数据接口直接请求接口返回JSON解析JSON提取字段比处理HTML省力得多。另一个坑是IP被限制。短时间频繁请求导致IP被临时封禁解决方案是设置下载延时、使用代理池轮换。不过需要特别提醒在学习和测试阶段尽量控制请求频率选择公开开放的公开数据集作为练习对象不要给目标站点造成过大压力。6.2 Hadoop 集群搭建与运行时的坑Hadoop这块踩的坑就比较经典了。伪分布式模式下最容易出问题的是启动之后NameNode没有正常起来。原因往往是没有配置JAVA_HOME或者core-site.xml、hdfs-site.xml里的路径配置有误。每次启动后建议马上检查进程列表和日志jps # 应该有 NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager如果NameNode没起来去看$HADOOP_HOME/logs/hadoop-user-namenode-host.log的报错信息。YARN运行时MapReduce任务卡在ACCEPTED状态不执行多半是ResourceManager的内存调度配置问题。在yarn-site.xml里把yarn.nodemanager.resource.memory-mb和yarn.scheduler.maximum-allocation-mb调大一点就能解决。还有一个经验值得分享Hive执行SQL前检查数据库的默认文件格式早期Hive版本默认是TEXTFILE后续版本支持ORC、Parquet。对于分析查询ORC格式的查询效率比TEXTFILE高很多。我的清洗结果表特意开启了ORC存储CREATE TABLE movies_clean_orc STORED AS ORC AS SELECT * FROM movies_clean;6.3 项目整体复盘与后续可以扩展的方向做完整的项目复盘之后我发现这个项目的架构对于学习大数据技术栈的人来说是一个很标准的样板但还存在几个扩展方向引入Spark替代MapReduce做计算性能会有数量级提升增加实时数据流比如接入实时票房数据使用Kafka Flink做流处理在电影简介字段上做中文分词和情感分析分析观众对国产电影的评论倾向把调度改成Airflow管理让整个数据管道每周自动执行一次我个人在实际操作中的体会是做这类项目不要纠结于技术“高级不高级”关键是每个组件都真实起到作用、每个环节都能回答“为什么选它”和“它解决了什么问题”。把这条链路跑通之后再去看面试题里的HDFS写入流程、MapReduce Shuffle原理、Hive底层执行引擎这些题目就不再是死记硬背而是能结合项目经验去理解了。这套项目做完你不仅学到了技术更建立起了完整的大数据处理思维。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价