
简介这份答辩PPT围绕潮流美妆大数据分析可视化系统的设计与实现展开主要面向高校计算机、数据科学及相关专业学生适用于课程设计、毕业设计或项目答辩等场景。内容完整覆盖课题背景与研究意义、国内外应用现状、系统总体设计、功能模块划分、爬虫数据抓取与清洗、Hadoop数据存储、Spark高性能处理、Django框架核心特性以及管理员功能界面展示并附带总结和参考文献可为准备相关主题答辩的同学提供清晰的演示框架和内容参考。资源包为1个pptx演示文稿大小24.49MB目前已有129人学习。PPT结构完整、重点突出既能帮助理解美妆数据从采集、清洗、存储、分析到可视化呈现的完整链路也可作为答辩提纲、页面排版和讲述节奏的参考样本有效节省自主梳理和制作时间。1. 从爬虫到图表的完整闭环才是美妆大数据项目的真正门槛先泼一盆冷水如果只写一个爬虫把数据导到 Excel 算完就交付这件事实在太简单。真正让“爬虫HadoopSparkDjango”这组技术栈产生价值的是后面那条完整的数据链路Scrapy 负责把美妆商品、价格、评论抓下来Hadoop 负责把噪声很大的原始数据按可回溯的方式存进 HDFSSpark 负责在内存里完成聚合和特征分析最后 Django 再把结果以可视化界面呈现给运营人员或答辩评委。这套结构在课程设计、毕业设计和中小型数据平台里非常常见每个环节都有替代品但组合起来之后链路口径、性能瓶颈和排错思路会变得完全不一样。我拆过不少类似项目发现最常出问题的往往不是爬虫本身而是“原始数据到底有多少条、清洗后剩多少条、页面上图表里的数字对应哪一层结果”。这篇文章按数据流向逐个讲清楚。2. Scrapy 爬虫层定义 Item、清洗入库、控制采集速度2.1 用 Scrapy 而不用 requests核心差异在并发和管道很多新手会先选requests BeautifulSoup因为它写起来直接发请求、拿 HTML、解析、存 CSV。但爬到几万条商品评论后就会发现请求调度、失败重试、去重、限速全要自己手写代码会从 80 行膨胀到 800 行。Scrapy 把这些能力内置化了Scheduler 负责请求队列Downloader Middleware 负责代理和重试Item Pipeline 负责数据后处理Spider 只需要关心“从页面里提取什么”。两者的选择其实不复杂我一般按这个表判断对比维度requests BeautifulSoupScrapy请求调度手写循环控制粒度粗内置 Scheduler去重和队列开箱即用失败重试自己封装 requests 异常RetryMiddleware配置RETRY_TIMES即可抓取速率手写time.sleep()容易抖动DOWNLOAD_DELAY或 AutoThrottle 自动限速数据处理解析后手动调用写入函数Spider 返回 ItemPipeline 统一清洗和入库扩展性单机脚本进程内串行可配合 scrapy-redis 做分布式抓取对美妆这种品类多、分页深、评论量大的场景Scrapy 的异步并发优势很明显。它基于 Twisted 事件循环几十个请求可以在同一时间内交替等待响应而不需要手动开多线程。2.2 先定义好 Item再写 Spider字段才不会散在写 Spider 之前一定要先在items.py里把数据模型定下来。这个项目的数据字段包含商品 ID、名称、价格、评论内容、评分、创建时间。字段定义不统一后续 HDFS 和 MySQL 都会跟着乱。# items.py import scrapy class BeautyCommentItem(scrapy.Item): product_id scrapy.Field() # 商品唯一标识用于去重和关联 product_name scrapy.Field() # 商品名称 price scrapy.Field() # 成交价格原始字符串 comment scrapy.Field() # 评论内容 comment_score scrapy.Field() # 评分1-5 的整数 create_time scrapy.Field() # 评论时间ISO 格式Spider 里的做法是每解析到一个评论节点就yield一个 Item。yield意味着把对象交给框架后续由 Engine 转发给 Item Pipeline当前请求可以立即结束从而腾出资源处理下一个请求。# spider.py import scrapy from beauty.items import BeautyCommentItem class BeautySpider(scrapy.Spider): name beauty start_urls [https://example.com/beauty/list] def parse(self, response): for sel in response.css(div.comment-item): item BeautyCommentItem( product_idsel.css(a::attr(data-pid)).get(), product_namesel.css(h3.product-name::text).get(), pricesel.css(span.price::text).get(), commentsel.css(p.comment-content::text).get(), comment_scoresel.css(span.score::text).get(), create_timesel.css(time::attr(datetime)).get(), ) yield item需要注意response.css(...).get()只会取第一个匹配结果使用..get()而不校验空值会导致管道层拿到None。所以清洗逻辑必须放到 Pipeline 里统一处理不在 Spider 里做判断这样 Spider 保持“只负责提取”的单一职责。2.3 Pipeline 清洗入库重复校验与批量写 MySQL抓回来的评论经常有重复记录比如同一用户在不同分页被重复抓取。入库时我习惯用INSERT ... ON DUPLICATE KEY UPDATE做幂等更新前提是product_id和create_time构成唯一索引。Pipeline 里的另一个重点是批量提交每 200 条executemany一次避免逐条 commit 导致写入性能急剧下降。# pipelines.py import pymysql class MySQLPipeline: def open_spider(self, spider): self.conn pymysql.connect( hostlocalhost, port3306, userbeauty_user, password******, databasebeauty, charsetutf8mb4, ) self.cursor self.conn.cursor() self.buffer [] def process_item(self, item, spider): self.buffer.append(( item.get(product_id), item.get(product_name), item.get(price), item.get(comment), item.get(comment_score), item.get(create_time), )) if len(self.buffer) 200: self.insert_batch() self.buffer [] return item def insert_batch(self): sql INSERT INTO comment (product_id, product_name, price, comment, comment_score, create_time) VALUES (%s, %s, %s, %s, %s, %s) ON DUPLICATE KEY UPDATE comment_score VALUES(comment_score) self.cursor.executemany(sql, self.buffer) self.conn.commit() def close_spider(self, spider): self.insert_batch() self.cursor.close() self.conn.close()参数说明utf8mb4是必须要用的字符集评论里经常有 emoji普通utf8会直接报错VALUES(comment_score)是 MySQL 在重复键时更新的旧字段写法。如果业务上要求“重复数据直接丢弃”把 ON DUPLICATE 改成INSERT IGNORE就行但要注意IGNORE也会吞掉其他类型的 SQL 错误排错时容易一头雾水。2.4 限速、重试与 robots 合规过去见过有人为了追求速度把CONCURRENT_REQUESTS调到 64结果对方服务器直接返回 403整批 IP 被拉黑。这里的原则是抓取频率要商量着来能慢则慢。# settings.py ROBOTSTXT_OBEY True DOWNLOAD_DELAY 1.2 # 同一域名下两次请求间隔 1.2 秒 CONCURRENT_REQUESTS 8 # 全局并发不要超过 16 RETRY_ENABLED True RETRY_TIMES 2 # 重试 2 次即可过多会加重目标站压力 DEFAULT_REQUEST_HEADERS { User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64), Accept-Language: zh-CN,zh;q0.9, }DOWNLOAD_DELAY的单位是秒它的作用是均匀化请求间隔比time.sleep()更优雅因为它由下载中间件统一调度。RETRY_TIMES不要设置成 5 或更高很多反爬系统对“短时间重复请求同一 URL”的惩罚力度远大于“慢速但持续”。需要特别提醒爬虫一定要遵守 robots 协议和当地法律法规。政企项目通常有授权范围毕业设计也应使用公开可访问的数据集或允许爬取的站点。爬虫技术本身中性但在真实场景里要优先考虑数据来源的合规性这在答辩时也是加分项。3. Hadoop 存储与清洗HDFS 分区布局与 MapReduce 去噪3.1 HDFS 目录为什么按层级和时间分区爬虫把数据写入 MySQL 后是否还需要 Hadoop答案是需要。MySQL 适合点查询但你要做全量历史数据的清洗、分析和回溯MySQL 的存储成本和计算吞吐都跟不上。把原始评论数据落到 HDFS等于给项目加了一层“不可变的事实来源”路径数据内容生命周期/app/raw/beauty/dt2025-03-01原始 JSON/CSV未清洗永久保留用于问题追溯/app/clean/beauty/dt2025-03-01清洗后的结构化文本保留最近 30 天可归档/app/agg/beauty/brand_salesSpark 聚合结果覆盖写每次跑批覆盖只留最新按dtyyyy-MM-dd做目录分区好处是清晰某天数据有问题只需要重跑那一天的清洗任务不用动其它历史数据Spark 做增量分析时还可以直接用input_path/dt2025-03-01精确指到某一天。HDFS 的底层存储策略简单流式读取性能远强于 MySQL 这种行存储适合 MapReduce 和 Spark 顺序扫描。3.2 用 Hadoop Streaming 跑 Python 清洗任务清洗任务是典型的“读原始文件 → 过滤脏数据 → 写新目录”。这项任务用 Java 写太重我习惯用 Hadoop Streaming 配合 Python直接用hadoop jar命令提交。hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -input /app/raw/beauty/dt2025-03-01 \ -output /app/clean/beauty/dt2025-03-01 \ -mapper python3 mapper.py \ -file mapper.py先写一个 mapper读入每行 JSON过滤缺少关键字段的数据同时做简单的价格格式化#!/usr/bin/env python3 # mapper.py import sys, json for line in sys.stdin: line line.strip() if not line: continue try: data json.loads(line) except json.JSONDecodeError: continue product_id data.get(product_id) comment data.get(comment) if not product_id or not comment: continue price data.get(price, 0) cleaned json.dumps({ product_id: product_id, product_name: data.get(product_name, ), price: price.replace(¥, ).replace(,, ).strip(), comment: comment.strip(), comment_score: int(data.get(comment_score) or 0), create_time: data.get(create_time, ) }, ensure_asciiFalse) print(cleaned)这个 job 没有 reducer我故意不指定-reducer而是看情况加-D mapreduce.job.reduces0。这样 Map 输出的每一行会直接写入-output目录适合“行级过滤 字段整理”这种无需 shuffle 的清洗。如果要做跨行去重比如按用户 ID 去掉重复评论才需要引入 reducer把product_id create_time作为 key由 reducer 做语义去重。执行完成后检查输出行数可以直接用hdfs dfs -cat /app/clean/beauty/dt2025-03-01/* | wc -l这个数字在后期做数据一致性核对时非常关键。3.3 伪分布式常见的配置项与坑本地开发时我不会一上来就搞三台机器而是搭伪分布式所有 Hadoop 进程在一台机器上配置和真实集群一致只是单机版。第一次搭建最容易踩坑的文件是core-site.xml和hdfs-site.xml!-- core-site.xml -- property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/data/hadoop/tmp/value /property!-- hdfs-site.xml -- property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name value/data/hadoop/name/value /property property namedfs.datanode.data.dir/name value/data/hadoop/data/value /propertydfs.replication在伪分布式里只能设成 1否则 Datanode 会反复尝试复制副本但永远找不到第二台机器。经典的坑是格式化 NameNode 后启动 DataNode 失败检查日志发现data.dir里残留了旧版本的 VERSION 文件。解决办法是把/data/hadoop/data和/data/hadoop/name里的 current 目录清空再执行hdfs namenode -format。另外hadoop.tmp.dir不要放在/tmp系统重启后目录被清空集群状态就丢了。这里的参数都建议写在$HADOOP_HOME/etc/hadoop/下改完要重启集群先stop-dfs.sh再start-dfs.sh。如果namenode起不来第一时间看/data/hadoop/name/current是否生成了fsimage没有就说明格式化不完整。4. Spark 会话与特征处理从 HDFS 到 MySQL 的分析提速4.1 为什么 MapReduce 之后还要 Spark同一份清洗数据MapReduce 能算为什么非得 Spark原因很简单MapReduce 的每个 job 都会把中间结果写磁盘而美妆评论分析这个场景要反复迭代比如先算“品牌商品总数”再算“平均评分”再按时间窗口切一次MapReduce 会写多个中间目录每次都要重新从 HDFS 读磁盘 IO 成为瓶颈。Spark 走的是一条“内存计算”路线。RDD 和 DataFrame 可以常驻内存多个 transformation 组成 DAG 后Spark 会尽量在同一个 pipeline 里串起全部分区只有遇到 shuffle 才真正落盘。所以同样的清洗逻辑Spark 在几 G 到几十 G 数据量下跑得比 MapReduce 快得多而且 PySpark 的 DataFrame API 读起来更接近 SQL团队里的人更容易维护。4.2 DataFrame 实现聚合与简单情感分析假设清洗后的数据在 HDFS 的/app/clean/beauty/dt2025-03-01需要算每个商品的评论数、平均价格、平均评分。PySpark 的写法是from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, avg, udf from pyspark.sql.types import StringType spark SparkSession.builder \ .appName(beauty_analysis) \ .config(spark.executor.memory, 2g) \ .config(spark.sql.shuffle.partitions, 20) \ .getOrCreate() df spark.read \ .option(header, True) \ .option(charset, UTF-8) \ .csv(/app/clean/beauty/dt2025-03-01) df df.filter(col(product_id).isNotNull()) \ .withColumn(price, col(price).cast(double)) result df.groupBy(product_id, product_name) \ .agg( count(comment).alias(comment_cnt), avg(price).alias(avg_price), avg(comment_score).alias(avg_score) ) \ .orderBy(col(comment_cnt).desc()) result.show(20)说明一下withColumn(price, col(price).cast(double))是把字符串价格转成数值很多线上 CSV 里会出现199.00如果 map 阶段没清干净这里 cast 会得到 null。所以我在聚合前用filter把product_id为空的行剔掉实际生产里还可以加dropDuplicates([product_id, create_time])做全量去重。情感分析不需要太复杂用“词典打分”的方式足够应付答辩展示。定义一个 UDF把常见正向词和负向词写死在 Python 字典里positive_words [好用, 回购, 保湿, 显白, 推荐] negative_words [过敏, 刺激, 难用, 差评] def sentiment(comment): if not comment: return neutral for w in negative_words: if w in comment: return negative for w in positive_words: if w in comment: return positive return neutral sentiment_udf udf(sentiment, StringType()) df df.withColumn(sentiment, sentiment_udf(col(comment))) df.groupBy(product_id, sentiment).count().show()udf()会引入逐行 Python 调用的开销但数据量在百万级以内完全可接受。如果要处理亿级数据就应该改用when when表达式或者 Spark NLP避免 Python UDF 拖慢整个 stage。4.3 写回 MySQLJDBC 与 foreachPartition 的选择Spark 算完的结果最终要回到 MySQL 里给 Django 读取。最简单的做法是 JDBCresult.write \ .mode(overwrite) \ .jdbc( jdbc:mysql://localhost:3306/beauty?useSSLfalsecharacterEncodingutf8, brand_sales, properties{user: spark_user, password: ******, driver: com.mysql.cj.jdbc.Driver} )这种方法适合小结果集比如品牌销售汇总表几千行。但如果结果行数很多JDBC 连接会变成瓶颈。更推荐foreachPartition每个 executor 里的分区各自建一个连接分批写既避免了单点连接太多也保留了批量操作的能力。def write_partition(rows): import pymysql conn pymysql.connect( hostlocalhost, port3306, userspark_user, password******, databasebeauty, charsetutf8mb4 ) cursor conn.cursor() sql INSERT INTO brand_sales (product_id, product_name, comment_cnt, avg_price, avg_score) VALUES (%s, %s, %s, %s, %s) ON DUPLICATE KEY UPDATE comment_cnt VALUES(comment_cnt), avg_price VALUES(avg_price) cursor.executemany(sql, list(rows)) conn.commit() cursor.close() conn.close() result.foreachPartition(write_partition)参数上spark.sql.shuffle.partitions默认是 200对单机项目太大了改成 20 能减少很多空任务spark.executor.memory在本地环境给 2g 够用真实集群要按 executor 数除以总内存来估算。5. Django 可视化平台MTV 架构、查询优化与图表接口5.1 MTV 各层职责与浏览器的交互路径Django 官方把它的设计模式称作 MTV三个字母分别代表 Model、Template、View。这不是 MVC 的改名而是对职责切分的一种具体落地层对应文件职责Modelmodels.py通过 ORM 定义表结构把数据库表映射成 Python 类Templatetemplates/*.html负责页面渲染不要写业务逻辑Viewviews.py接收请求、查询数据、返回 HttpResponse 或 JsonResponseURLConfurls.py配置 URL 与 View 函数之间的路由映射浏览器访问/api/sales/时Django 先读 urls.py 找到对应 views 函数View 内走 ORM 查 MySQL拿到结果后再交给 Template 或直接序列化成 JSON。和 MVC 比Django 的“视图”更接近 Controller模板承担了传统 View 的职责这一点答辩时务必说清楚。5.2 ORM 查询既要快又要省values、annotate 与分页因为 Spark 已经把计算结果写进了 MySQLDjango 这边不需要做复杂聚合主要工作是查询单表和分页。以下是最常见的列表接口写法# views.py from django.http import JsonResponse from django.core.paginator import Paginator from .models import BrandSale def brand_sale_list(request): qs BrandSale.objects.all().order_by(-comment_cnt) paginator Paginator(qs, 20) page paginator.get_page(request.GET.get(page, 1)) items list(page.object_list.values( product_id, product_name, comment_cnt, avg_price, avg_score )) return JsonResponse({ page: page.number, total: paginator.count, items: items, })这里的values(...)会生成只包含指定列的 SQL而不是把整行对象加载到内存。Paginator(qs, 20)做的事情是对所有结果COUNT(*)再对当前页做LIMIT 20 OFFSET ...。200 万行数据下COUNT也会很慢性能瓶颈通常在index设计上brand_sales表的product_id、comment_cnt字段建议建立联合索引。如果在 Django 端直接做聚合可以用annotate但要注意它产生的 SQL 是GROUP BY数据量大时一定加索引别在视图里写 Python for 循环去统计。from django.db.models import Count, Avg from .models import Comment data Comment.objects.values(product_id, product_name) \ .annotate( comment_cntCount(id), avg_scoreAvg(comment_score) ) \ .order_by(-comment_cnt)[:100]5.3 图表渲染用 ECharts后端只负责给 JSON可视化页面我一般不在 Django 模板里堆字符串而是直接在前端用 ECharts。后端只提供 JSON 接口前端fetch获取数据后交给 ECharts 的setOption。这样前后端职责分明也方便答辩现场临时换图表类型。// 假设 static/js/chart.js fetch(/api/sales/?page1) .then(res res.json()) .then(data { const chart echarts.init(document.getElementById(chart)); chart.setOption({ xAxis: { type: category, data: data.items.map(i i.product_name) }, yAxis: { type: value }, series: [{ type: bar, data: data.items.map(i i.comment_cnt) }] }); });注意echarts.init的容器必须有显式高度styleheight:400px;必不可少否则图表高度会自适应为 0。展示多个图表时每个容器调用一次init对应不同setOption配置。5.4 导出 CSV 用 StreamingHttpResponse参数要设对答辩评委经常要求“现场导出数据”直接用 Django 的StreamingHttpResponse流式生成 CSV 是最稳的方案不会一次性把几千行数据全部塞进内存。# views.py from django.http import StreamingHttpResponse import csv def export_sales_csv(request): def generate_rows(): yield [product_id, product_name, comment_cnt, avg_price] for sale in BrandSale.objects.values( product_id, product_name, comment_cnt, avg_price ).iterator(chunk_size500): yield [sale[product_id], sale[product_name], str(sale[comment_cnt]), str(sale[avg_price])] response StreamingHttpResponse(generate_rows(), content_typetext/csv; charsetutf-8) response[Content-Disposition] attachment; filenamebrand_sales.csv return responsecontent_type决定浏览器按什么类型解析响应text/csv告诉浏览器这是可下载文件Content-Disposition里的attachment是触发下载的关键参数如果去掉浏览器会直接尝试在页面里渲染 CSV 内容。iterator(chunk_size500)也很重要它让 ORM 按 500 行一批从 MySQL 读取而不是把所有结果加载到内存对几十万行导出很关键。6. 答辩前做一次端到端数据核对与演示缓存6.1 双端行数核对脚本答辩最容易翻车的场景是评委问“你这个页面上显示的数据和爬虫抓的数据对得上吗”你支支吾吾答不上来。所以演示前一定要做一次“两端核对”写一个短脚本把三层数据数量对齐。核对思路MySQL 的comment表行数对应 HDFS 清洗后的文件总行数Spark 聚合后的记录数对应 Django 接口返回的记录数。先看第二层hdfs dfs -cat /app/clean/beauty/dt2025-03-01/* | wc -l再在 MySQL 里执行SELECT COUNT(*) FROM comment WHERE DATE(create_time) 2025-03-01;两个数字应该完全一致。如果 HDFS 行数大于 MySQL说明有评论还没入库或 Pipeline 中途失败如果小于说明 HDFS 清洗时过滤掉了太多有效记录。第三层核对是查看brand_sales表行数它应该等于 MySQL 评论表按product_id去重后的商品数量SELECT COUNT(DISTINCT product_id) FROM comment;6.2 演示现场的缓存与降级策略答辩现场网络不稳定时接口如果每次都要现算一旦踩到慢 SQL页面会长时间白屏。我的做法是给热点接口加一层缓存最简单的是用 Django 内置的 cache_page 装饰器from django.views.decorators.cache import cache_page cache_page(60 * 10) def brand_sale_list(request): ...cache_page(60 * 10)表示这个接口的响应缓存 10 分钟。开发环境用 LocMemCache生产环境换成 Redis加上 DummyCache 可以一键关闭缓存用于调试。演示前先手动访问一次接口让缓存生效现场刷新时就会快很多。还有一个隐藏技巧把 Spark 聚合结果写入 MySQL 时不要覆盖式更新整张表建议用INSERT ... ON DUPLICATE KEY UPDATE配合create_time字段做增量写入。这样演示时重跑某一天的 job只会影响那天的数据页面上的历史趋势不会跳动。整个系统的“可信度”就体现在这些数据能对上、能解释得通。本文还有配套的精品资源点击获取