尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

基于Spark与Hadoop的新闻数据分析可视化系统实战

基于Spark与Hadoop的新闻数据分析可视化系统实战 做大数据方向的课程设计、毕业设计或者找工作投简历时的项目积累很多人一看到“基于 Spark 的新浪网数据分析可视化系统”这种题目第一反应是这不就是爬点新闻、做个图表吗真正做下去才发现这里的水比想象中深不少。我自己完整跑通过一个 Hadoop Spark Django 的版本数据源用新浪网的公开内容做演示最终在页面上呈现可视化大屏包括词云、趋势折线、分类占比、热榜排行这些常见模块。这个项目最大的价值不是某个算法有多深而是它把“数据采集 - 存储 - 清洗 - 分析 - 后端接口 - 大屏展示”整条链路串起来了并且每一层都用到了对应场景下的流行工具。如果你正在准备课程设计、毕业设计或者想在简历上写一个能讲清楚的大数据项目这篇内容可以帮你把每个环节的关键点理顺选型怎么考虑、数据怎么处理、Spark 分析怎么写、Django 怎么跟大屏对接、哪些坑可以提前避开。1. 项目整体设计与技术选型思路1.1 需求梳理这个系统到底要做什么很多同学拿到题目后会直接想“我要写多少代码”但正确的顺序是先拆需求。基于 Spark 的分析可视化系统不管界面叫什么名字目标其实非常固定对一批和“新浪网”相关的数据完成从原始文本到统计结果的加工最后用网页展示分析结论。既然项目名字里带了“新浪网”数据源起码得有来源属性比如新闻标题、发布时间、来源板块、阅读/评论数据否则分析结果和主题贴合不上。我在动手之前给自己列了五个可交付的功能点一套可运行的存储和计算环境HDFS 作为数据底座Spark 负责分析一个数据采集或构造脚本产出规范的 CSV/JSON 数据Spark 清洗和聚合任务输出关键词统计、分类占比、时间趋势、热度排行Django Web 后端提供 JSON 接口给前端使用可视化大屏用 ECharts 展示统计结果。这里特别提醒一下如果拿不到新浪网的实时数据完全可以用历史公开数据集或者自己构造一批带“新浪新闻分类风格”的模拟数据来演示重点在于链路完整和分析逻辑正确老师验收时关心的也是你能不能把这个流程解释清楚。1.2 技术选型为什么是 Hadoop Spark Django而不是其他组合项目叫“基于 Spark”但题目里同时绑了 Hadoop所以存储层用 HDFS 基本是必须的。HDFS 在这里的角色是分布式文件系统Spark 从 HDFS 上读取原始数据分析完成后再把结果写回。很多同学刚开始不理解“先上 Hadoop 的意义”觉得单机读文件就行。但从课程设计的展示角度HDFS 能体现分布式存储思想从数据处理角度如果后续把数据量扩展到 GB 级别Spark 加 HDFS 的价值才能体现出来。Spark 与 Hadoop MapReduce 相比优势在于中间结果可以尽量放在内存里写代码也友好得多。尤其做词频统计、分类聚合这类操作用 Spark SQL 或者 RDD 算子都很直观。Django 的选择则是因为它的 ORM、模板和后台管理非常稳在“需要搭建一个带可视化大屏的 Web 系统”这种场景下比 Flask 的“自由”更适合结构化分工。如果一定要换方案比如把 Django 换成 FastAPI把可视化换成 Vue 大屏模板也可以但那意味着你要同时管前端工程化和后端接口两套东西复杂度会明显上升。课设和简历项目讲究“能用最小的成本跑通最完整的链路”所以 Django 直接渲染页面、内嵌 ECharts反而是相对省力的做法。2. 数据采集与预处理存储层的实际落地2.1 数据源与字段设计数据源我选择了新浪网公开的新闻类内容。需要注意做课程设计时不要试图抓取全量数据也不要去碰有访问限制的接口。我用 requests BeautifulSoup 抓取了几大类新闻的标题和基础信息比如国内、国际、体育、娱乐、科技这几个分类每个分类抓取了几百条最终数据量控制在几千到几万条之间。采集频率上加了随机延时避免对目标站点造成压力这也是做爬虫的基本礼仪。存储到本地前我先规划了字段方便后续 Spark 清洗和分析news_id唯一标识title新闻标题category所属分类source来源名称publish_time发布时间comment_count评论数url原文链接content: 正文摘要用于分词这些字段看着简单但每个都有用途。比如 category 用来做“分类占比”饼图publish_time 用来做“发布时间趋势”comment_count 用来做“热度排名”content 或 title 用来做“关键词词频”。字段设计做好了后面五个分析指标都能直接对号入座。2.2 原始数据上传 HDFS数据抓下来之后我先统一保存成 CSV 格式编码务必使用 UTF-8。很多同学在做中文分析时出现乱码问题基本都出在这里Windows 下 Excel 默认可能用 GBK但线上环境和 Spark 读取时更多默认 UTF-8两边不一致就会出乱码。数据准备好后把它传进 HDFShadoop fs -mkdir -p /user/sina/data hadoop fs -put ./sina_news.csv /user/sina/data/传完后用hadoop fs -ls确认文件确实存在再去写 Spark 任务。有些同学会把 HDFS 命令忽略掉直接在本地路径上跑 Spark这样最后的项目演示里“Hadoop”就只是个摆设面试或者答辩时容易被追问。2.3 Spark 数据清洗要点原始数据不能直接用来统计否则会出现重复数据、空值、格式不统一等问题。清洗这一步我用 Spark 的 DataFrame API 来做。主要逻辑有四个去重、过滤空值、格式统一、字段类型转换。from pyspark.sql import SparkSession from pyspark.sql.functions import col, to_date, hour, when spark SparkSession.builder \ .appName(sina_news_clean) \ .getOrCreate() df spark.read.option(header, True).csv(hdfs://localhost:9000/user/sina/data/sina_news.csv) df df.dropDuplicates([news_id, title]) df df.filter(col(title).isNotNull()) df df.filter(col(publish_time) ! ) df df.withColumn(comment_count, col(comment_count).cast(int)) \ .withColumn(publish_date, to_date(col(publish_time), yyyy-MM-dd HH:mm:ss)) df.write.mode(overwrite).parquet(hdfs://localhost:9000/user/sina/clean_data)清洗完成后我存成了 Parquet 格式。Parquet 是列式存储比 CSV 更省空间而且 Spark 读取时效率更高。对于课设数据来说两者差异不大但从“加分项”角度来看能说出“为什么用 Parquet”是明显的亮点。3. 基于 Spark 的核心分析实现3.1 分析任务拆解四类业务指标怎么算清洗完成后Spark 的真正重头戏是算指标。我这里的分析目标对应大屏上的几个可视化组件关键词 Top N对标题或正文摘要做分词统计高频词各类别新闻数量占比统计不同分类的条数计算比例按小时/日期的发布趋势从 publish_time 中提取时间维度统计数量变化热度新闻 Top 榜按评论数排序取评论最多的一批新闻。第一个指标用 RDD 算子比较好理解因为分词结果本身是一个个词。我当时用了 jieba 分词库在 Spark 里通过 map 操作对每条新闻做分词再用 reduceByKey 累加词频。import jieba from pyspark.sql import Row def segment(text): words jieba.cut(text) return [w.strip() for w in words if len(w.strip()) 1 and w.strip() not in stopwords] rdd df.select(title).rdd.map(lambda row: row[title]) word_rdd rdd.flatMap(segment).map(lambda word: (word, 1)) word_count word_rdd.reduceByKey(lambda a, b: a b) top_words word_count.sortBy(lambda x: x[1], ascendingFalse).take(50)这里有个细节停用词表非常关键。如果不过滤“我们”“可以”“一个”“什么”这类词词云里最大的一定是这些没营养的连接词而不是有业务含义的关键词。停用词表可以从 GitHub 上找开源的中文停用词表再结合自己数据里的实际情况补充。第二、三、四个指标用 Spark SQL 反而更简洁。把清洗后的 DataFrame 注册成临时视图直接写 SQL 做聚合和排序df.createOrReplaceTempView(news) # 分类占比 category_stats spark.sql( SELECT category, COUNT(*) AS cnt FROM news GROUP BY category ORDER BY cnt DESC ) # 发布时间趋势 hour_stats spark.sql( SELECT HOUR(publish_time) AS hour, COUNT(*) AS cnt FROM news GROUP BY HOUR(publish_time) ORDER BY hour ) # 评论数 Top 榜 top_news spark.sql( SELECT title, category, comment_count, publish_time, url FROM news ORDER BY comment_count DESC LIMIT 10 )这种“RDD Spark SQL”混合使用的方式非常贴合实际开发。RDD 适合处理无结构或半结构的文本迭代计算Spark SQL 适合处理结构化数据的过滤和聚合两个工具各有适用场景硬要用一个打天下反而会让代码变得很别扭。3.2 分析结果的落库方式Spark 计算出来的结果最终要给 Django 用但 Django 不可能直接去 HDFS 上读文件更不可能直接在 Web 请求里启动一个 Spark 任务。数据链路需要再接一段把分析结果导出到关系型数据库Django 再通过 ORM 查询。我在项目里是先把 Spark 分析结果写成 CSV 文件再导入 MySQL。还有一种方案是直接用 Spark 的 JDBC 连接器写 MySQL但这需要额外处理驱动依赖对初学阶段来说有点绕。如果做第二版我也考虑过把结果存成 JSON 文件放在某个固定目录Django 那边直接读 JSON 也能满足页面展示需求但数据规范性不如 MySQL 好。实际导入时我用 Django 的bulk_create批量插入几万条数据导入很快。导入完成后Django 的活就是纯粹的读表、聚合、出接口。分析任务执行完毕后日志里会有每个 stage 的耗时和 shuffle 数据量。这里有一个值得养成的习惯记录一下“读取了多少条原始数据清洗后剩多少条每个分析结果多少条”这个数据在写文档和答辩时很有用。词频统计结果的保存格式(keyword, cnt)分类统计结果的保存格式(category, cnt)趋势统计结果的保存格式(hour, cnt)新闻热度榜的保存格式(title, category, comment_count, publish_time, url)在清洗和统计分析之前一定要看数据总量。通常几万条以内的数据在单机 Spark 上跑是秒级完成的如果发现异常慢先看是不是分区数设得太大或 GC 频繁。3.3 关于分区数和内存设置的提醒Spark 刚上手最典型的两个问题一个是内存溢出一个是任务特别慢。数据量本身不大但默认配置可能导致性能异常。如果使用spark-submit或pyspark提交任务客户端模式下会继承 JVM 默认内存堆内存通常不够用。可以在提交时加参数比如spark-submit \ --master local[2] \ --executor-memory 2g \ --driver-memory 2g \ analyze.py假如是单机 8G 内存的笔记本分配给 Spark 的总内存不要超过物理内存的一半否则操作系统本身和 Django 都会受影响。分区数方面可以用spark.sql.shuffle.partitions控制spark.conf.set(spark.sql.shuffle.partitions, 10)默认值是 200对几万条数据来说太浪费了。把一份小数据切成 200 份再 shuffle 聚合就像把一张纸撕碎成两百片再拼起来徒增开销。数据量小就调小分区数这是大数据分析里很实在的一条经验。4. Django 后端与可视化大屏对接4.1 Django 工程结构与数据读取Django 部分我按常规方式创建了一个项目和一个 app名称就叫dashboard。ORM 模型与 Spark 分析结果对应# dashboard/models.py from django.db import models class CategoryStat(models.Model): category models.CharField(max_length50) count models.IntegerField() class Meta: db_table category_stat class KeywordStat(models.Model): keyword models.CharField(max_length50) count models.IntegerField() class Meta: db_table keyword_stat class TrendStat(models.Model): hour models.IntegerField() count models.IntegerField() class Meta: db_table trend_stat class NewsRank(models.Model): title models.CharField(max_length500) category models.CharField(max_length50) comment_count models.IntegerField() publish_time models.CharField(max_length50) url models.URLField() class Meta: db_table news_rank四个模型对应四张大屏图表结构非常简单。注意 MySQL 表名不要用 Django 默认生成的长名字手动指定db_table可以方便后续调试。4.2 JSON 接口实现方式Django 端不建议直接返回 HTML 片段给前端图表更规范的方式是设计统一 JSON 接口。大屏页面启动时通过 AJAX 拉取数据再渲染 ECharts。在视图文件里写一个 JSON 响应的公共函数可以简化代码from django.http import JsonResponse from dashboard.models import CategoryStat, KeywordStat, TrendStat, NewsRank def api_category(request): data list(CategoryStat.objects.values(category, count)) return JsonResponse({code: 0, data: data}) def api_keyword(request): data list(KeywordStat.objects.order_by(-count)[:50].values(keyword, count)) return JsonResponse({code: 0, data: data}) def api_trend(request): data list(TrendStat.objects.order_by(hour).values(hour, count)) return JsonResponse({code: 0, data: data}) def api_top_news(request): data list(NewsRank.objects.order_by(-comment_count)[:10].values( title, comment_count, publish_time, category, url )) return JsonResponse({code: 0, data: data})实际项目中为了省事可以把四个接口合并成一个/api/overview一次性把大屏所有数据都返回。考虑性能时就分开。演示场景下分开更清晰方便单独排查数据问题。关于跨域问题如果你的页面直接放在 Django 模板里由 Django 渲染那不存在跨域。如果你把前端拆出去用 Vite 单独起一个服务再用 Nginx 部署就需要配置django-cors-headers同时要注意请求地址写的是http://127.0.0.1:8000而不是相对路径。课程设计里最简单的方式就是让页面从 Django 的静态文件或模板中加载。4.3 可视化大屏的布局与 ECharts 接入可视化大屏的“大屏感”来自三个方面深色背景、信息分区、图表丰富度。我做的是 16:9 页面分三栏布局中间栏放核心指标或标题左侧放分类占比饼图和词云右侧放小时趋势折线图和新闻热度滚动条。顶栏放大标题整体视觉风格类似驾驶舱 Dashboard。大屏页面引入 ECharts 时推荐下载 echarts.min.js 放到项目静态目录不要用 CDN。因为演示时如果现场网络不好图表会加载失败。词云图需要额外的 echarts-wordcloud 插件要下载与 ECharts 版本兼容的版本建议用 echarts 4.9 搭配 wordcloud 1.1.3兼容性最稳。页面模板里每个图表容器是一个 div通过 JavaScript 在初始化时请求接口然后 setOption。下面是一个典型写法div idchartCategory stylewidth: 100%; height: 300px;/divfetch(/api/category) .then(response response.json()) .then(res { let chart echarts.init(document.getElementById(chartCategory)); chart.setOption({ series: [{ type: pie, radius: [40%, 70%], data: res.data.map(item ({ name: item.category, value: item.count })) }] }); });实际开发时可以封装一个initChart函数把每个图表的初始化、AJAX 加载、异常处理统一起来避免每个图表写一大段重复逻辑。如果大屏需要定时刷新给 AJAX 部分包一层setInterval即可但注意刷新时间不要小于 10 秒否则图表会闪烁反而影响演示效果。大屏的视觉细节不要低估。背景色建议用深蓝或深灰主标题加亮色边框图表之间保持间距。如果页面只是白底加普通图表很难叫“可视化大屏”。我是先写一个简单的body { background: #0f1c3a; }再在图表容器外面加一圈半透明边框整体效果立刻不同。5. 环境搭建与部署踩坑实录5.1 Hadoop、Spark 与 Django 的版本匹配版本匹配是大数据项目最容易踩坑的地方没有之一。我建议的稳定组合Python 3.8 或 3.9Java 8Hadoop 3.3.xSpark 3.3.x自带 Hadoop 3 clientDjango 4.xMySQL 5.7 或 8.0这里再次提醒如果你的机器内存只有 8G不建议开三节点集群单机 Hadoop 伪分布式模式更合适。伪分布式并不是“缩水版”它仍然会启动 NameNode、DataNode、ResourceManager、NodeManager 等完整角色只是都在一台机器上运行演示和教学完全足够。Hadoop 安装时的几个步骤容易出错。第一是JAVA_HOME必须指向 JDK 安装目录不能只设置到/usr/lib/jvm上层第二是启动前要执行hdfs namenode -format格式化只需要一次不要每次启动都格式化否则 HDFS 会掉数据。第三是修改core-site.xml和hdfs-site.xml时确认fs.defaultFS的地址与代码中读写的 HDFS 路径前缀一致。Django 部署我直接放在本机开发和演示。如果没有独立服务器用python manage.py runserver 0.0.0.0:8000作为成品演示也够了。如果部署到 Linux 服务器则可以用 gunicorn 或 uwsgi 来跑。5.2 运行过程中的常见问题与处理方法我在调试过程中遇到的问题非常多挑几个高频的列成表格方便直接对照排查问题现象可能原因解决办法Hadoop 启动后进程少一个格式化多次、目录冲突清空 tmp 目录并重新格式化再启动访问 HDFS 路径报错路径不存在或权限不足先用hadoop fs -mkdir -p创建确认当前用户有权限Spark 日志报 Java heap space内存分配不足提交任务时加--driver-memory 2g或把 JVM 参数调大Spark 结果中文乱码PySpark 或 CSV 编码问题统一使用 UTF-8 编码并在读取 CSV 时指定 encodingPy4J 连接失败Spark Context 未正确启动手动启动 pyspark 测试再执行脚本Django 请求接口返回 500ORM 表名或字段不匹配检查模型 Meta 的 db_table 与 MySQL 表是否一致ECharts 图不显示数据格式不对或容器高度为 0输出接口 JSON 检查结构确认 div 有明确 heightExcel 打开 CSV 中文乱码本地编码与 UTF-8 冲突用文本编辑器检查编码或转成 UTF-8-BOM 再给 Excel 用还有一个特别常见的错误在 Windows 下用本地 Spark 连接 HDFS提示找不到 winutils。解决办法是下载对应版本的 winutils.exe放到 Hadoop 的 bin 目录然后设置HADOOP_HOME环境变量。不同 Hadoop 版本对应的 winutils 版本不完全一样要尽量选择和 Hadoop 版本一致的包。5.3 演示与文档整理的方向由于这个项目本身带了“源码文档调试”这几个交付项我习惯在最终整理时把内容分成三个部分环境搭建文档、系统设计说明、操作演示脚本。环境搭建文档要细到“下载哪个版本、解压到哪个目录、修改哪个配置文件、启动哪条命令”因为这份文档不只是给老师看的还是给一个月后的自己看的。系统设计说明建议写清楚三点整体架构、数据流向、每个模块的核心代码解释。把数据从“新浪网页 - 爬虫 - HDFS - Spark - MySQL - Django 后端 - ECharts 大屏”的流动过程用一个数据流向图来描述可以放在文档引言部分。操作演示脚本则是给答辩现场用的建议只留三条命令一条启动 Hadoop一条执行 Spark 分析脚本一条启动 Django然后进入大屏页面。答辩被追问“为什么用 Spark”、“Spark 与 MapReduce 有什么区别”时可以围绕“内存计算、API 表达能力、适合迭代式分析”回答。如果被追问“数据量多大才算大数据”不要说“10TB”这种在自己项目里没法佐证的数值可以说“在集群资源足够的情况下Spark 能横向扩展本项目主要验证的是从存储到分析再到可视化的完整流程”。我自己在做这个项目时最深的一个体会是项目链路看着长真正占用时间最多的不是分析逻辑而是环境调试和数据格式问题。建议按“先本地小数据跑通、再往 HDFS/Spark 迁移数据、最后做 Django 可视化和大屏美化”的顺序去推进不要在第一天就追求全套组件都装好并启动成功。先将 Spark 分析脚本在本地跑通再逐步过渡到 HDFS 读取会让整个过程顺畅非常多。最后再分享一个小的实战技巧分析结果导入 MySQL 之后先不要急着写前端用 Django 的 admin 或者直接连数据库查一下每张表的数据量。一个表如果是空大屏上对应图表就会空白而这个问题出在“导入过程”而不是“前端代码”排查起来很浪费时间。把数据层先确认好再调展示层整个系统的开发节奏会稳很多。
返回列表