简介面向计算机相关专业毕业设计的美妆大数据分析可视化系统论文完整覆盖Hadoop、爬虫与Spark三大技术栈并融合Python、MySQL、PyCharm、Scrapy及数据可视化等工具结构完整是可直接参考的项目文档。资源压缩包仅含1个docx文件大小3.85MB内含摘要、中英文关键词、目录、参考文献及完整正文正文从背景意义、开发环境搭建、Scrapy爬虫编写到数据清洗整合、MySQL存储、可视化图表生成均有详细展开系统实现路线清晰。已有535人学习适合作为毕业设计选题、系统架构设计及论文写作的模板。通过阅读该论文可快速理解网络爬虫原理、Hadoop与Spark的大数据处理流程以及数据采集到可视化分析的全链路实现方法对于需要完成类似课题或准备答辩的学生还能在章节组织、技术路线和系统设计层面获得有效参考。1. 美妆大数据分析系统为什么要把爬虫、Hadoop、Spark 串成一条线美妆商品的评论、价格和成分数据单条看起来很小但把商品详情、评论列表、促销快照各存一份数据量会指数级膨胀。口红的文字评论平均只有几十个字图片链接和成分列表却是非结构化内容传统关系库在单表几千万行时还能靠索引维持一旦要按品牌、价格段、时间窗口做聚合全表扫描会直接拖垮线上库。这个项目的核心思路是把采集、存储、计算、展示拆成四个独立层爬虫负责拿到原始文本Hadoop 负责横向扩展存储Spark 负责批量清洗和聚合可视化系统负责把结果变成运营能看的趋势图。这个链路适合两类人一类是正在做毕业设计或课程设计想搞清楚这些组件如何真正协作另一类是工作几年后准备从“脚本加 MySQL”升级到“集群加分析引擎”的工程师。下面按实现顺序讲每一步都给出可运行的命令和参数少讲概念多给操作。2. 爬虫层美妆评论与价格数据的采集、清洗和落盘格式2.1 先定字段再写爬虫美妆数据的最小语义模型写爬虫最大的误区是看到页面就开始解析标签结果抓下来的数据到分析阶段才发现缺品牌、缺促销价。面对美妆这种属性复杂的数据我先建议固定一个最小字段集至少包含商品 ID、品牌、价格、评分、评论内容、评论时间和来源渠道。价格要单独存原始价和促销价评论内容保留换行符和制表符因为清洗阶段要根据它们判断文本质量。下面是一个可以直接套用的字段表字段名类型来源说明product_idstring商品详情页关联产品和评论brandstring详情页品牌栏聚合维度缺失值要标记 unknownpricedouble价格区块促销价优先原始价兜底salesint列表页销量抓取时间点快照注意单位comment_idstring评论页去重主键的一部分ratingint评论页1 到 5 分非法值转 0comment_textstring评论页保留原始文本后续清洗comment_timestring评论页统一成 yyyy-MM-dd HH:mm:sssourcestring爬虫配置区分不同站点防止混合统计字段设计完成后再写代码。常见做法是用 requests 加 BeautifulSoup 做轻量采集抓成功后再交给分布式爬虫框架扩展。下面是一个可以跑通的采集脚本只抓评论页每次请求间隔 0.8 秒避免高频访问触发反爬。import requests, random, time, json from bs4 import BeautifulSoup UA_POOL [ Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36, Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36, ] def crawl_reviews(product_id, pages5): s requests.Session() out_path freviews_{product_id}.jsonl with open(out_path, w, encodingutf-8) as fp: for page in range(1, pages 1): url fhttps://example.com/product/{product_id}/reviews?page{page} resp s.get( url, headers{User-Agent: random.choice(UA_POOL)}, timeout10 ) soup BeautifulSoup(resp.text, html.parser) for r in soup.select(.review-item): record { product_id: product_id, comment_id: r.select_one(.review-id).text.strip(), rating: int(r.select_one(.rating).get(data-score)), comment_text: r.select_one(.comment).text.strip(), comment_time: r.select_one(.time).text.strip(), source: demo } fp.write(json.dumps(record, ensure_asciiFalse) \n) time.sleep(0.8) return out_path crawl_reviews(MZY001, pages2)这里必须解释几个参数Session会复用 TCP 连接比反复创建 requests 对象快也能保存部分 Cookietimeout10防止某个评论页卡住整个任务time.sleep(0.8)是采集节奏不是固定值如果目标站点对频率不敏感可以调到 0.3 秒但不要取消。输出用 JSON Lines 而不是 CSV因为评论里可能包含逗号、换行和引号CSV 在后续被 Spark 读取时需要额外处理转义而 JSON Lines 每一行是一个完整 JSON 对象天然支持多行文本。2.2 清洗阶段处理表情符号、重复评论和非法评分爬虫只能保证字段完整不能保证数据干净。做旅游网站大数据分析或招聘网站可视化项目时前期一半时间都花在清洗数据美妆数据也不例外。常见问题包括表情符号乱码、评分字段出现百分号、同一条评论因分页被重复抓取。这里用 Pandas 做一次快速清洗把结果再写成 JSON Lines。import pandas as pd df pd.read_json(reviews_MZY001.jsonl, linesTrue) df[comment_text] df[comment_text].str.replace( r[\U0001F300-\U0001FAFF], , regexTrue ) df[comment_text] df[comment_text].str.replace( r[\n\t\r], , regexTrue ) df[rating] pd.to_numeric(df[rating], errorscoerce).fillna(0).astype(int) df df.drop_duplicates(subset[comment_id], keepfirst) df.to_json(reviews_clean.jsonl, orientrecords, linesTrue, force_asciiFalse)\U0001F300-\U0001FAFF是 Unicode 表情符号区间直接替换成空字符串避免后续 Spark 读写时出现编码异常errorscoerce会把非数字评分转成 NaN再用fillna(0)填充保留该条评论而不是直接删除这样后续分析可以过滤掉 rating 等于 0 的数据。drop_duplicates(subset[comment_id], keepfirst)以评论 ID 为准去重保留第一次抓到的那条记录。到这里数据已经可以进入 Hadoop。通常我会在本地保留一份原始抓取文件再把清洗文件作为后续入仓的输入这样即使分析和可视化阶段出了问题还可以回到原始数据重新洗。3. Hadoop 存储层HDFS 文件组织、伪分布式搭建与数据入仓3.1 为什么选择 HDFS美妆原始数据是写多读少的非结构化数据评论数据的特征是不更新、只追加商品价格和销量则是一段时间内的快照。HDFS 的模型正好匹配这种场景一次写入、多次读取不支持随机修改。它把文件切成块分布在多个 DataNode 上适合存储几十 GB 以上的非结构化文本。如果数据只有几万条评论用 MySQL 完全够用一旦每天新增几百万条评论MySQL 的大表聚合就会变得吃力而且评论中没有固定的字段结构关系库的 schema 约束反而成了负担。使用 HDFS 还有一个实际原因Spark 可以直接读取 HDFS 上的 JSON Lines 或 Parquet不需要经过 MySQL 导出导入这道工序。爬虫清洗后的数据直接上传到 HDFS 专门目录分析任务就能开始跑省掉中间的文件搬运。3.2 Hadoop 伪分布式搭建三个配置文件和一次格式化搭建 Hadoop 环境是这条链路里最容易劝退人的一步。伪分布式不是 Hadoop 的标准形态但适合单机学习和跑通链路。我一般用 Hadoop 3.x提前配置好 JAVA_HOME 后再改三个文件。第一个是core-site.xml指定 NameNode 地址第二个是hdfs-site.xml指定副本数第三个是环境变量hadoop-env.sh指定 JDK 路径。下面是常见的配置方式。export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64 export HADOOP_HOME/opt/hadoop export PATH$HADOOP_HOME/bin:$HADOOP_HOME/sbin:$PATH cat $HADOOP_HOME/etc/hadoop/core-site.xml EOF configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/data/hadoop/tmp/value /property /configuration EOF cat $HADOOP_HOME/etc/hadoop/hdfs-site.xml EOF configuration property namedfs.replication/name value1/value /property /configuration EOFfs.defaultFS是整个集群的默认文件系统地址后面执行hdfs dfs -ls /时补全的根路径就来自这个配置。hadoop.tmp.dir指定元数据存放位置默认放在系统临时目录重启会丢。dfs.replication在伪分布式下必须设为 1因为只有一台 DataNode无法保存三份副本等以后扩展成真集群时改成 3并且需要引入 ZooKeeper 辅助 NameNode 做高可用。这里要特别注意很多搜索教程会把“hadoop 和 zookeeper 整合实战”直接混进单机安装导致你还没确认 HDFS 是否启动就先被 ZooKeeper 的选举状态卡住。伪分布式阶段不需要 ZooKeeper。完成配置后执行一次格式化并启动hdfs namenode -format start-dfs.sh hdfs dfs -mkdir -p /warehouse/beauty/reviews hdfs dfs -ls /namenode -format会初始化文件系统的元数据执行一次就够了多次格式化会导致 NameNode 与 DataNode 的集群 ID 不一致启动时报错。start-dfs.sh会同时拉起 NameNode 和 DataNode启动后可以通过http://localhost:9870查看分布式存储状况。如果 DataNode 启动失败优先看/opt/hadoop/logs目录下的日志不要反复格式化大部分问题来自hadoop.tmp.dir没有创建或者端口被占用。3.3 数据入仓按日期分区上传解决小文件隐患清洗好的reviews_clean.jsonl需要传入 HDFS。入仓规则建议按日期建目录因为增量更新和删除都围绕日期进行比如每天只用上传当天抓到的评论。上传命令如下DATE_PART2025-06-01 hdfs dfs -mkdir -p /warehouse/beauty/reviews/date$DATE_PART hdfs dfs -put reviews_clean.jsonl /warehouse/beauty/reviews/date$DATE_PART/ hdfs dfs -du -h /warehouse/beauty/reviews/date$DATE_PART/date$DATE_PART是 Hive/Spark 常用的分区目录格式分析引擎能够靠这个路径自动裁剪数据。du -h查看该目录下文件的实际大小当单文件明显小于 128MB 时建议先做一次合并再上传。HDFS 不适合存大量小文件NameNode 内存中要记录每个文件、目录和块的元数据几千个小评论文件会把元数据空间打满。常见做法是使用hdfs fs -appendToFile把小 JSONL 文件合并成一个或者在爬虫端定时把多条记录写进同一个文件保证每个文件大小接近一个块。下表列出我在配置 Hadoop 时最常核对的几个参数适合伪分布式场景配置项配置文件推荐值说明fs.defaultFScore-site.xmlhdfs://localhost:9000集群入口地址hadoop.tmp.dircore-site.xml/data/hadoop/tmp持久化元数据目录dfs.replicationhdfs-site.xml1单机必须为 1dfs.blocksizehdfs-site.xml134217728块大小默认 128MBdfs.namenode.name.dirhdfs-site.xml/data/hadoop/nameNameNode 镜像目录到这里原始评论已经进入 HDFS。用得多了以后会意识到Hadoop 安装与配置并不难难点是搞清楚每个参数到底会影响哪一层ZooKeeper、副本数、块大小这些概念如果不拆开理解后面会在集群扩展时反复踩坑。4. Spark 分析层Spark SQL 做数据清洗、日期聚合和特征统计4.1 Spark 集群搭建与本地调试模式的选择Spark 是计算引擎负责读取 HDFS 上的数据并完成清洗和聚合。这里先说集群搭建。真集群通常选择 Spark on YARN让 YARN 统一调度资源单机学习阶段可以直接用 Spark Standalone 或者最简单的方式在提交任务时指定--master local[*]让 Spark 跑在本地线程上。我一般先写 PySpark 脚本在本地跑通逻辑再用spark-submit提交到 YARN。资源参数设置比安装本身更容易出错。常见的错误是给 executor 内存设置过大比如机器只有 8GB--executor-memory 6g会直接导致集群没有额外内存给操作系统和 DataNode。下面是一个适用于 4 节点配置的提交命令spark-submit \ --master yarn \ --deploy-mode client \ --driver-memory 2g \ --executor-memory 4g \ --num-executors 4 \ --executor-cores 2 \ --conf spark.sql.shuffle.partitions20 \ analysis.py参数说明--master yarn表示资源调度交给 YARN--deploy-mode client保持客户端进程运行方便看日志--executor-memory 4g是单个 executor 可用的堆内存不是单机总内存--num-executors 4配合--executor-cores 2让每个 executor 同时执行两个任务。spark.sql.shuffle.partitions控制 shuffle 后的分区数量默认 200如果数据量不到百万行可以降到 20 到 50减少没有意义的空任务和输出文件数量。实际集群搭建时不要照抄这个命令要依据数据量和节点内存重新分配。4.2 Spark SQL 读 JSON Lines先确认 schema 再做日期处理Spark 读取 HDFS 上的 JSON Lines 时会自动推断字段类型但推断结果不一定符合分析需求。price 可能被识别成 stringcomment_time 也会变成 string。所以读取后的第一步不是统计分析而是查看 schema 并统一类型。下面的代码演示怎么读入数据并创建临时视图。from pyspark.sql import SparkSession, functions as F spark SparkSession.builder.appName(beauty_analysis).getOrCreate() df spark.read.json(/warehouse/beauty/reviews) df.printSchema() df.createOrReplaceTempView(reviews)read.json(/warehouse/beauty/reviews)会递归读取该目录下所有 JSON 文件Spark 直接利用date分区目录做路径裁剪实际读取时只会进到对应日期。printSchema()能立刻发现类型问题比如看到comment_time: string是很正常的后面我们必须用to_date或to_timestamp转换。这里不建议用samplingRatio1.0去推断全量数据代价太高通常生产数据只要观察几万行就能确定 schema。4.3 日期加减和聚合好评率、价格分组和时间窗口美妆数据里最常用的分析维度是时间、品牌、价格带和评分。Spark SQL 的日期处理是重点尤其“spark sql 日期加减”这类问题经常出现在实际需求里。比如要看三个月前到现在的评分趋势需要把评论时间减三个月要看未来一年的复购窗口需要把评论时间加一年。下面的 SQL 先清洗时间字段再生成一个聚合结果。CREATE OR REPLACE TEMP VIEW reviews_clean AS SELECT product_id, brand, CAST(price AS DOUBLE) AS price, rating, comment_text, TO_DATE(comment_time, yyyy-MM-dd HH:mm:ss) AS comment_date, ADD_MONTHS(TO_DATE(comment_time, yyyy-MM-dd HH:mm:ss), 12) AS next_year_date FROM reviews WHERE rating IS NOT NULL AND rating 0 AND comment_text IS NOT NULL;TO_DATE转换时指定了两个格式参数第一个参数是待转换字段第二个是解释字符串模板Spark 里的yyyy是年MM是月dd是日大小写不能错。ADD_MONTHS(..., 12)表示把日期加 12 个月得到下一年的对应日期如果想做最近 30 天则使用DATE_SUB(comment_date, 30)。过滤条件里的rating 0对应清洗阶段填充的非法评分避免把脏数据带进统计。接下来做分品牌、分月的聚合SELECT brand, DATE_FORMAT(comment_date, yyyy-MM) AS month, COUNT(*) AS review_count, ROUND(AVG(rating), 2) AS avg_rating, SUM(CASE WHEN rating 4 THEN 1 ELSE 0 END) AS positive_count, ROUND(SUM(CASE WHEN rating 4 THEN 1 ELSE 0 END) / COUNT(*), 4) AS positive_rate FROM reviews_clean GROUP BY brand, DATE_FORMAT(comment_date, yyyy-MM) ORDER BY month DESC, review_count DESC;DATE_FORMAT(comment_date, yyyy-MM)把日期折叠成月份positive_rate用好评数除以评论总数得到好评率这个指标比平均评分更能反映用户对美妆产品的真实接受度。GROUP BY 后面必须同时写brand和DATE_FORMAT(comment_date, yyyy-MM)SELECT 里的别名不能再出现在 GROUP BY 中。4.4 用 PySpark 写回结果Parquet 是比 CSV 更合适的落盘格式聚合结果如果每次都要用 SQL 控制台导出效率太低。正确方式是把聚合结果写到 HDFS用 Parquet 格式保存。Parquet 是列式存储后续只查某几个字段时能大幅减少 IO。下面是一个完整的 PySpark 脚本片段daily_stats spark.sql( SELECT product_id, DATE_FORMAT(comment_date, yyyy-MM-dd) AS day, COUNT(*) AS review_count, ROUND(AVG(rating), 2) AS avg_rating, SUM(CASE WHEN rating 4 THEN 1 ELSE 0 END) AS positive_count FROM reviews_clean GROUP BY product_id, DATE_FORMAT(comment_date, yyyy-MM-dd) ) daily_stats.coalesce(1).write.mode(overwrite).parquet(/warehouse/beauty/stats/daily)coalesce(1)的作用是让最后生成一个数据文件方便查数和启动下一个流程但生产环境下它会迫使所有数据经过单分区造成任务瓶颈。更好的做法是用repartition(4, day)按天分区输出每一个日期一个文件。mode(overwrite)会覆盖已有结果适合重复跑批如果业务要求保留历史结果应该使用mode(append)并配合日期分区条件。Spark 分析层做完后HDFS 上的stats目录已经是非常清晰的聚合结果。下面是可视化层要做的事把这些结果导出到 MySQL并给前端提供查询接口。5. 可视化层Spark 结果写 MySQL再用 Flask 和 ECharts 出图5.1 Spark 结果到 MySQL 的写库策略Spark 写 MySQL 要避免小事务。如果直接对几百个分区同时写库MySQL 连接数会被瞬间打满常见的提示是Too many connections。建议在写库前重新分区控制并发连接数。用 JDBC 写入时URL 里的参数比想象中更影响性能。下面是一个标准写法daily_stats.repartition(3).write.jdbc( urljdbc:mysql://localhost:3306/beauty_data?useSSLfalseserverTimezoneAsia/ShanghairewriteBatchedStatementstrue, tabledaily_stats, modeoverwrite, properties{ user: report, password: your_pass_here, driver: com.mysql.cj.jdbc.Driver, batchsize: 1000, truncate: true } )repartition(3)把 DataFrame 压缩到 3 个分区这样最多只有 3 个并发写入线程。rewriteBatchedStatementstrue是 MySQL JDBC 的批处理优化参数能把多条 INSERT 打包成一条语法执行大幅度降低网络往返。truncatetrue配合mode(overwrite)会先清空表再写入避免重复数据累积。写库之前还需要确认 MySQL 表结构。不建议把原始评论导到 MySQL只导聚合结果表格字段与 Spark 输出保持一致字段名类型说明daydate统计日期product_idvarchar(32)商品 IDreview_countint当日评论总数avg_ratingdecimal(3,2)当日平均评分positive_countint当日好评数评分大于等于 4建表时给(day, product_id)建联合主键这样 Spark 重复写入时能够通过主键去重避免可视化面板出现同一天两条记录。5.2 Flask 后端返回 JSON 接口参数校验和防 SQL 注入可视化系统不直接读 HDFS而是通过后端接口查询 MySQL。Flask 在这里只做一层轻量封装把 SQL 查询结果转成 JSON 给前端。接口设计通常分成趋势接口和排行榜接口下面是一个最简单的商品趋势接口支持product_id和days参数。from flask import Flask, request, jsonify import mysql.connector app Flask(__name__) DB_CONFIG { host: localhost, user: report, password: your_pass_here, database: beauty_data, } app.route(/api/trend) def trend(): product_id request.args.get(product_id, ) if not product_id: return jsonify({error: product_id is required}), 400 days request.args.get(days, 30, typeint) if days 0 or days 365: return jsonify({error: days must be between 1 and 365}), 400 conn mysql.connector.connect(**DB_CONFIG) cursor conn.cursor(dictionaryTrue) cursor.execute( SELECT day, review_count, avg_rating, positive_count FROM daily_stats WHERE product_id %s AND day DATE_SUB(CURDATE(), INTERVAL %s DAY) ORDER BY day , (product_id, days)) rows cursor.fetchall() cursor.close() conn.close() return jsonify(rows) if __name__ __main__: app.run(host0.0.0.0, port5000)request.args.get(days, 30, typeint)把查询字符串里的days强制转成整数这样传入字符串时不会报错。SQL 中的%s是参数占位符具体的值由cursor.execute的第二个参数传入这是防止 SQL 注入最基本的手段。INTERVAL %s DAY可以在 MySQL 原生语法里动态生成日期范围比在 Python 侧拼好字符串再传进去更安全。dictionaryTrue让游标返回字典格式的列表jsonify直接把列表转成 JSON 数组返回。如果可视化页面需要同时展示多个指标可以在这个接口后面再扩展一个分组接口但思路完全一致。5.3 前端 ECharts 折线图把 JSON 直接映射到 series后端接口就绪后前端只需要调用fetch拿到数据数组再映射到 ECharts 的xAxis和series。可视化部分不追求高难度重点是让图表能够直接表达 Spark 算出来的趋势。下面是一段可以在管理后台页面直接使用的 JavaScriptfetch(/api/trend?product_idMZY001days30) .then(res res.json()) .then(data { const chart echarts.init(document.getElementById(trendChart)); chart.setOption({ xAxis: { type: category, data: data.map(d d.day) }, yAxis: { type: value, min: 0, max: 5 }, series: [{ name: 平均评分, type: line, smooth: true, data: data.map(d d.avg_rating) }, { name: 好评数, type: bar, yAxisIndex: 1, data: data.map(d d.positive_count) }] }); });xAxis.data和series.data都直接来自接口返回的day和avg_rating字段。smooth: true让折线更圆滑适合展示评分趋势。温度计式的双 Y 轴在这里可有可无如果评论数和评分数量级差距太大可以像示例一样为positive_count单独指定yAxisIndex: 1。可视化层的定时更新也需要处理。常见做法是写一个 crontab每天凌晨跑一次 Spark 分析任务完成后再触发一次写库脚本前端看到的数据就是昨天的完整数据。如果要求小时级刷新只需把DATE_FORMAT的精度改成yyyy-MM-dd HH并且调整调度频率。到这里从爬虫到前端展示的链路已经完整但系统能跑起来和能稳定跑是两回事。6. 进阶验证三个自查指标和一处必调的写入调优6.1 用三个指标判断链路是否正确系统上线前先不要急着做复杂图表用最直接的数量一致性验证。第一个指标是评论去重后的总条数在 HDFS 上用hdfs dfs -cat /warehouse/beauty/reviews/date2025-06-01/*.jsonl | wc -l统计原始行数在 MySQL 里执行SELECT COUNT(*) FROM daily_stats WHERE day 2025-06-01两个数字如果不匹配说明 Spark 清洗时过滤条件写错了。第二个指标是时间范围一致性查询聚合结果里的最小日期和最大日期确认没有把2024-06-01的数据混进2025-06-01分区分区目录经常因为上传时路径写错导致这种错位。第三个指标是空值比例原始评论里品牌字段为空的数量不应该超过 5%超过则要回看爬虫字段解析逻辑。这三个指标都能用脚本定期执行看到异常时再逐层排查比盲目调优靠谱得多。6.2 写 MySQL 前强制减少分区避免连接数打满Spark 写 MySQL 最常见的坑是分区数等于并发连接数。前面代码里用了repartition(3)但实际任务里有些人会误用coalesce(3)。区别在于repartition会打乱数据重新分配coalesce只合并分区不会重新 shuffle。如果某个分区数据量特别大比如某个爆款美妆商品评论占了全量的 80%coalesce之后会导致单个任务承担全部写入压力效率反而下降。建议写库前统一使用repartition并按主键字段做repartition(3, day)让同一天的评论尽量落在同一分区MySQL 端也能通过主键顺序提升插入效率。另一个值得配置的是 Spark 的动态合并参数。在提交任务时加上--conf spark.sql.adaptive.enabledtrue --conf spark.sql.adaptive.coalescePartitions.enabledtrueSpark 可以在 shuffle 结束后自动合并过小的分区减少生成的文件数。这两个参数在 Spark 3 之后默认开启但显式写出来能让后续维护的人知道你的调优意图。真正上线前先用 1 万条样本跑通 HDFS 分区、Spark 日期清洗、MySQL 主键这三条链路确认没有比发现数据不一致更崩溃的线上反馈。本文还有配套的精品资源点击获取