简介:这是一套基于Spark2.2的新闻网大数据实时分析系统毕业设计/课程设计项目,面向大数据相关专业学生与初学者,用于解决新闻网站访问日志实时采集、流式计算和趋势统计等场景,并涉及智能推荐相关实现。压缩包共403个文件,约262KB,其中364个xml承担项目配置与构建描述,14个scala及5个java构成核心计算逻辑,4个sh脚本辅助环境启动,另有md/txt文档便于查阅,整体目录结构清晰。项目源码均经本地编译可运行,按附带文档配置环境即可启动运行;内容由助教老师审定,难度适中,可作为毕业设计或课程设计的完整参考模板,也能帮助理解Spark Streaming与大数据分析流程。已有242人浏览学习,下载后可获得可运行源码、环境配置文档、辅助脚本与完整的工程结构范例,适合需要快速落地Spark实时分析项目的读者,使用中遇到问题也可直接咨询作者。
1. 毕设中的Spark2.2实时分析系统:别只会跑通Demo
很多做这块毕设的人,答辩前几天才发现自己只能演示一个wordcount,稍微问一句"集群上有多少个Worker、Redis怎么支撑高并发"就卡壳。这套资源好在不是那种只有一个Demo壳子的工程,它把新闻网站从点击日志到实时热点、再到个性化推荐,整条链路都串起来了。你拿到的不是一个孤立的Spark程序,而是一套可以讲清楚"数据从哪来、算完存哪、推荐怎么用"的完整系统。这篇文章我按自己拆项目的方式,把架构、核心代码、调参、踩坑和答辩演示一遍拆给你看,让你拿到资源后能直接照着改,而不是对着压缩包发呆。
2. 系统架构与数据流:Kafka + Spark Streaming + Redis,每个组件选出性价比
2.1 数据从哪来:Flume日志采集与Kafka消息队列
新闻网站的访问日志跟普通日志不一样,它要记"谁在什么时候点了哪条新闻",还会带上来源渠道、设备ID、停留时长等字段。日志文件写在N台Web服务器上,如果直接让Spark Streaming去读文件,会有三个问题:日志文件按天滚动,路径要动态感知;文件正在写入时读容易切到半行;多台服务器不好统一管理。常规做法是每台服务器装一个Flume agent,监听日志目录,把新写入的行通过Kafka sink发到Topic里,Spark Streaming作为消费者再拉数据。
Flume的配置不算复杂,但要注意source的spooldir与taildir区别。常见做法是taildir,因为它支持断点续传位置,不会在服务重启后把昨天的日志又发一遍。
agent.sources = src agent.channels = ch agent.sinks = k1 agent.sources.src.type = TAILDIR agent.sources.src.filegroups = fg1 agent.sources.src.filegroups.fg1 = /data/logs/news/.*\.log agent.sources.src.positionFile = /opt/flume/taildir_position.json agent.channels.ch.type = memory agent.channels.ch.capacity = 50000 agent.channels.ch.transactionCapacity = 10000 agent.sinks.k1.type = org.apache.flume.sink.kafka.KafkaSink agent.sinks.k1.kafka.topic = news-click-log agent.sinks.k1.kafka.bootstrap.servers = kafka01:9092,kafka02:9092 agent.sinks.k1.kafka.flumeBatchSize = 500 agent.channels.ch.transactionCapacity = 10000 agent.sinks.k1.channel = ch这段配置的要点在positionFile。尾目录模式下,Flume会记住每个文件的偏移量,程序重飘时不会从头读,也不会漏数据。Kafka Sink的batchSize决定批量发送的条数,设太大会增加延迟,设太小会让Kafka吞吐上不去。我一般会控制在500到1000之间。
Kafka的Topic建议按日志类型分,news-click-log存曝光和点击,如果要存每篇文章的正文抓取结果,建议另开一个news-article-raw,不要混在一个Topic里。因为曝光日志有大量无效点击,拿它去做内容推荐会偏。另外Kafka的分区数不是越多越好,分区数要和下游Spark Streaming的并行度对上。你在Kafka里建分区时,先想好后面Streaming会开多少个并行接收任务,分区数取它不少。
2.2 计算在哪做:Spark Streaming的DStream与窗口
Spark Streaming在2.2时代有两种姿势:老牌的DStream和刚稳定一点的Structured Streaming。说实话,如果这是毕设,我建议你选DStream。原因很简单,网上能搜到的、能在2.2不踩坑的例程,95%都是DStream。Structured Streaming的Event Time、水印、输出模式在2.2里还不够顺手,你答辩讲原理时不占便宜,Debug时也容易绕进去。
DStream的本质是"离散流",它把源源不断的数据按固定时间片切成一批批RDD。你在DStream上写的每个算子,最终都会作用到这批RDD上。这点必须跟答辩老师讲清楚:Spark Streaming不是一条一条处理,而是微批处理,实时性来自"批间隔够短",而不是"事件触发"。
窗口操作是这套系统里最值得深挖的地方。新闻热点不是"这秒发生了啥",而是"最近一段时间里哪些新闻被点得多"。所以我们需要把当前批跟之前几批攒在一起算,这就是窗口函数。
val ssc = new StreamingContext(sparkConf, Seconds(5)) val lines = KafkaUtils.createStream[String, String, StringDecoder, StringDecoder]( ssc, kafkaParams, topics, StorageLevel.MEMORY_AND_DISK_SER) val clickCounts = lines .map(_._2.split(",")) .filter(_.length >= 4) .map(fields => (fields(3), 1L)) .reduceByKeyAndWindow((v1: Long, v2: Long) => v1 + v2, Seconds(300), Seconds(60))这里Seconds(5)是批间隔,Seconds(300)是窗口长度,Seconds(60)是滑动步长。意思是每60秒算一次"从当前时刻倒推过去5分钟内每篇新闻的点击量"。注意窗口长度和滑动步长都必须是批间隔的整数倍,否则Spark会直接报IllegalArgument异常。
有同学会问,为什么窗口长度不是批间隔的5倍?因为窗口两端可能有边界效应。比如你设窗口5秒、批间隔5秒,那说白了就是每批独立统计,没有重叠,只有相邻批拼起来才能表达"一段时间",所以保持倍数关系并不够,更重要的是理解窗口计算时会有状态积累。reduceByKeyAndWindow有两个版本:一个带filterFunc,一个是invFunc(逆函数)。如果不带逆函数,每个窗口都要从过去300秒的内存里重新算,CPU开销大;带逆函数可以只做增量计算。毕设场景下数据量不大,不带逆函数也能跑,但我建议你还是加上逆函数,因为答辩老师可能会问"窗口计算性能为什么慢"。
2.3 结果存到哪:Redis和MySQL的分工
计算结果分两类。一类用于大屏实时展示,比如"当前热点前十名",要求毫秒级读取,用Redis最合适。另一类是历史分析,比如"某天每小时的点击趋势",要长时间保存,用MySQL。
Redis的数据结构别只会用String,这里强烈推荐Sorted Set。我们把新闻ID塞进score里,点击量作为score值,然后用ZREVRANGE取前N个就是榜单。
MySQL表设计要简单实用,至少两张表:news(新闻ID、标题、发布时间、栏目)和click_stats(统计时间、新闻ID、点击量、PV、UV)。注意click_stats不要按天建表,按天建表后期SQL会很难看,直接用date字段分区就够用。
3. 实时热点统计:从Kafka消费到窗口聚合的完整代码
3.1 用Scala写一个稳定的Kafka消费器
这部分是最容易"看起来懂,做起飞"的地方。直接用KafkaUtils是的createStream是基于Receiver的,它会把数据先放到Executor的内存里,如果这个节点挂掉,内存中未处理完的数据会丢,除非开启WAL(Write Ahead Log)。开启WAL的方式是ssc.getConf.set("spark.streaming.receiver.writeAheadLog.enable", "true")。但WAL会写入HDFS,导致性能下降,而且Receiver模式在Executor重启后要重新rebalance,经常出现重复消费。
这里我建议用DirectStream(也叫No Receiver模式),它是从Kafka的offset范围直接读取,不通过Receiver,这样不会有数据先驻留内存,也不用WAL。但DirectStream在Spark 2.2里对应的Kafka版本比较讲究,我很少让它新到0.10,因为Spark2.2配Kafka0.8.2.1是最稳的。下面这段代码在老项目里能直接跑:
import kafka.serializer.StringDecoder import org.apache.spark.streaming.kafka.KafkaUtils val kafkaParams = Map( "metadata.broker.list" -> "localhost:9092", "group.id" -> "news-realtime-group", "auto.offset.reset" -> "largest", "enable.auto.commit" -> "false" ) val topics = Set("news-click-log") val lines = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder]( ssc, kafkaParams, topics )auto.offset.reset设为largest表示从最新offset开始消费;如果调试时想从头看,就改成smallest。enable.auto.commit设为false是让你把offset管理权握在自己手里,后面可以用rdd.foreachRDD处理完后手动提交,否则程序一重启,可能会从旧offset开始重复消费整段数据。
这段代码其实不依赖Receiver,所以资源消耗小。但注意这里返回的lines不是DStream吗?对,DirectStream仍以DStream形式存在,只是内部实现不同。你需要告诉老师"DirectStream根据你传入的Topic和分区数,自己维护一个偏移量范围,然后直接向Kafka broker请求对应区间的数据,而不是通过Receiver缓存"。
3.2 窗口统计与Top-N排序的正确写法
拿到点击流之后,先解析日志。日志一行大概是20250321102300,123.45.67.89,user001,news_9527,click,5,我们关心的是第4个字段newsId和第5个字段type。
窗口统计可以用前面的reduceByKeyAndWindow,但Top-N怎么取?很多同学的错误是在transform里做take,这会阻塞Streaming运行。正确的是在foreachRDD里取top,或者先用transform把RDD转换成新的DStream,再在输出端操作。
我一般这样写:
val windowedCounts = clickCounts .map(_._2.split(",")) .filter(f => f.length >= 5 && f(4).equals("click")) .map(f => (f(3), 1L)) .reduceByKeyAndWindow((a: Long, b: Long) => a + b, Seconds(300), Seconds(60)) windowedCounts.foreachRDD { rdd => val sorted = rdd.coalesce(1, true).sortBy(_._2, false).take(10) // 写入Redis sorted.foreach { case (newsId, cnt) => jedis.zadd("hot_rank", cnt.toDouble, newsId) } }这里coalesce(1)是为了让全量数据在一个分区排序。如果不做coalesce,sortBy会触发shuffle,无法得到全局有序的top10。每次foreachRDD都执行取top,相当于一个action算子,在流式任务里是允许的,因为它已经落在输出阶段。不要在transform内部做take,transform是懒转换,take会触发作业执行,两者混在一起容易造成作业重复提交或阻塞。
3.3 参数怎么调:batchDuration、windowDuration、slideDuration
下面这个表是我常用的初始值,你可以根据集群规模改:
| 参数名 | 建议值 | 说明 |
|---|---|---|
| spark.streaming.batch.duration | 2秒 | 指的是每批多少秒,实际在创建StreamingContext时传入Seconds(2) |
| windowDuration | 300秒 | 统计一个滚动窗口,比如最近5分钟 |
| slideDuration | 30秒 | 榜单刷新频率,每30秒刷新一次 |
| spark.streaming.kafka.maxRatePerPartition | 1000 | 每分区最大消费速率,防止压垮下游 |
| spark.streaming.backpressure.enabled | true | 背压机制,自动调节消费速率 |
这里有一个核心原则:windowDuration和slideDuration都必须是batchDuration的整数倍。比如batch 2秒,slide 30秒,window 300秒,这三个数都能被2整除,没问题。如果不满足,Spark会在启动时直接抛"Multiple of batch duration"异常。
参数不是拍脑袋设的。如果batchDuration过小,比如0.5秒,批次数太多,任务调度开销大;如果过大,比如10秒,数据延迟会明显。好在你这是毕设,数据量不大,按上表设置能跑得很顺。
4. 智能推荐模块:用户行为实时打分与ALS离线修正
4.1 基于内容的实时推荐:TF-IDF + 余弦相似度
推荐不能光靠热度,还需要看"这个用户点了什么,下一步给他推什么"。最简单有效的内容推荐是拿新闻的标题和正文做TF-IDF,转成向量,然后计算余弦相似度。
Spark MLlib有现成的TF-IDF实现,不过2025年回头用2.2,你会发现它接口有点老。建议先对新闻内容做分词,再构建HashingTF和IDF。
import org.apache.spark.mllib.feature.{HashingTF, IDF} import org.apache.spark.mllib.linalg.{Vector, Vectors} val tf = new HashingTF(100000) val newsFeatures = newsRDD.map { doc => val terms = doc.title.split("\\s+") (doc.id, tf.transform(terms.toSeq)) } val idf = new IDF().fit(newsFeatures.map(_._2)) val newsFeaturesVec = newsFeatures.map { case(id, vec) => (id, idf.transform(vec)) }TF-IDF会为每篇新闻生成一个稀疏向量。当用户点击了某篇新闻A,就遍历新闻库中的其他新闻,计算A的向量与其他新闻向量的余弦相似度,取TopN。
注意做大循环前,先把向量广播出去,否则每计算一次就要把全量新闻向量拉回来,十分慢。我一般用sc.broadcast(newsVectorMap)。
4.2 用ALS离线训练做个性化召回
内容相似度只能推荐"和A长得很像的文章",但用户可能想看不同类型但符合偏好的内容。所以需要训练一个协同过滤模型。ALS在Spark MLlib里很容易上手:
import org.apache.spark.mllib.recommendation.{ALS, Rating} val ratings = userClickRDD.map { case(user, newsId) => Rating(user.hashCode, newsId.hashCode, 1.0) } val model = ALS.train(ratings, rank = 20, iterations = 10, lambda = 0.01) val userId = "user001".hashCode val topRecs = model.recommendProducts(userId, 20)这里rank表示隐向量的维度,维度越高表达越精细,但维度太高容易过拟合且计算量大。iterations别超过15,太多了未必收敛且非常耗时。lambda是正则化系数,越大防止过拟合越强,我建议从0.01起调。
ALS生成的是离线候选集,你可以把它存入Redis的Sorted Set,key为user:rec:离线,value为新闻ID,score为模型预测评分。
4.3 推荐服务如何融合实时和离线结果
实际点开网站推荐栏时,要同时考虑"刚刚大家都在看的热点"和"这个用户长期喜欢的内容"。一个简单有效的融合方式是加权:
final_score = 0.6 * 离线als_score + 0.3 * 内容相似度_score + 0.1 * 实时热度_normalized把每个候选新闻的最终分数算好后,写入Redis的user:rec:final,再用ZREVRANGE取前N条返回给前端。
实时部分要从之前的榜单拿,那是一个DStream数据流;离线部分是从HDFS或MySQL读出来的。融合逻辑可以写在一个独立线程池里,每5秒从Redis里拉一次实时榜单,再合并离线候选,不用把离线模型塞进流处理中。这样流任务保持稳定,推荐服务也能独立调试。
5. Spark2.2实时项目常见踩坑:现象、原因、解决
5.1 任务不跑但也不报错,日志里全是WARN?
现象:StreamingContext启动后,控制台只打Recieved block ...和WARN spark.scheduler.TaskSchedulerImpl: Initial job has not accepted any resources,程序也不退出,但就是不出统计结果。
原因:你的Spark Streaming里启动了executor,但分配的核数不够。因为Spark Streaming至少需要两个核,一个用于调度,一个用于执行。我经常看到有人setMaster("local[1]"),那自然只有调度的核,没有执行核。
解决:本地调式用local[2],集群提交时--executor-cores 2以上。如果在YARN上,还要检查资源队列内存是否足够。
5.2 Kafka消费者组offset不生效,重启丢数据
现象:每次重启程序,都会重复消费之前已经算完的数据,甚至从最早开始刷。
原因:你用的是Receiver模式,offset由Kafka自己管理,但Receiver在Spark端通过WAL保存offset,如果WAL没开启或路径不对,重启offset自然复位。DirectStream则需要你手动提交offset,有些示例代码忘了提交。
解决:我推荐DirectStream配合手动提交。在foreachRDD处理后,把当前批的offset范围保存到外部(比如Redis或MySQL),重启时从保存位置恢复。不要在每次处理完立即提交,而是要在"处理成功且已保存结果"后提交,否则处理一半失败,offset会被标记为已消费。
5.3 窗口计算重复消费同一批数据
现象:你看到TopN里的点击量翻倍,但日志里没有那么多点击。
原因:窗口长度是批间隔的整数倍,但如果你用reduceByKeyAndWindow的invFunc写法不对,旧值没有被正确减去,就会造成累计虚高。另外,如果spark.streaming.receiver.writeAheadLog.enable为true,Receiver模式可能重复读。
解决:用带逆函数的版本:reduceByKeyAndWindow(_ + _, _ - _, window, slide, filterFunc)。我调试时会往Redis里打一份带时间戳的结果,如果同一条新闻在两次相邻窗口中重复出现但单个窗口内计数正常,说明逆函数失效,改成不带逆函数的版本试试。
5.4 Redis连接打满,实时统计变成阻塞
现象:系统运行几分钟后,Spark任务开始变慢,部分批次处理时间超过批间隔。
原因:每条数据或每个batch都新建Jedis连接,而Redis服务器默认最大连接数是10000,连接池配置不合理,导致创建连接的等待时间超过批次延迟。
解决:使用连接池,比如JedisPool,并将setMaster的SPARK广播变量里只放连接池对象,不要每次创建。还要把redis.pool.maxTotal设到500-1000,maxIdle设到100左右。另外,绝对不要在foreachRDD的循环外创建Jedis,否则会卡在driver端,造成内存泄漏。
5.5 源码能跑,但打包后ClassNotFound
现象:本地IDEA点Run没事,用spark-submit提交jar后,报NoClassDefFoundError或ClassNotFoundException,指向Kafka相关的类。
原因:你的jar包是瘦包,没把依赖的Kafka、Spark Streaming Kafka相关后缀带上。
解决:打包时用Maven的shade插件,把所有依赖合并进去,并注意排除签名文件。常见做法是在pom.xml里加maven-shade-plugin,生成一个-all.jar去提交。同时注意Kafka的版本跟spark-streaming-kafka-0-8_2.11的版本要一致,别一个0.8一个0.10。
6. 验证与优化:从零复现到答辩被倒问也扛得住
6.1 本机三件套自测:Kafka手动生产,Spark Streaming消费断言
在把整个系统跑起来之前,先做一个最小闭环,避免最后连报错都分不清是谁的锅。我习惯开三个终端:
终端一启动ZooKeeper和Kafka,创建测试Topic:
kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic news-click-log-test终端二启动生产者,手动发两条格式正常的行:
kafka-console-producer.sh --broker-list localhost:9092 --topic news-click-log-test > 20250321093000,1.2.3.4,user001,news_9527,click,5 > 20250321093005,1.2.3.5,user001,news_9528,click,3终端三启动Spark应用,看是否打印出这两个newsId的计数。这里有一个验证技巧:你在foreachRDD里打印rdd.count,如果count大于0,说明消费成功;然后把Redis里的榜单打印出来,看看是不是(news_9527, 1)和(news_9528, 1)。如果榜单是空,优先检查Kafka的Topic名称是否写错、日志中的分隔符对不对。
6.2 性能优化的两个抓手:并行度与序列化
Spark Streaming性能瓶颈通常不在CPU,而在IO和序列化。新闻日志字段多,如果你用默认的Java序列化,会非常耗内存。建议在SparkConf里写:
conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")Kryo序列化比Java快很多,但要注意给需要复用的类注册。另外,Kafka的消费分区数、Spark的并行度、Redis的写入分片要成比例。比如Kafka Topic分了3个分区,Spark设置后端并行度为3,每个分片写一个Redis槽位,这样不会出现某个Redis分片被写爆。
6.3 答辩演示技巧:用真实请求展示实时更新
答辩时不要只展示一个静态页面。你可以准备一个小脚本,每5秒向本地Kafka发送一条带当前时间戳的模拟日志,然后在大屏Dashboard上刷新,让排名实时变化。评委看到数字自己跳,会立刻觉得系统是活的。
这里有一个小坑:模拟数据的时间戳要跟系统时间接近,否则窗口计算里Old数据会被视为延迟数据忽略。你可以在生产者里加一个time.sync,确保机器时间准确。
演示时被问"如果Kafka挂了怎么办",你可以说自己系统里加了spark.streaming.receiver.writeAheadLog.enable=true,或者指向ZooKeeper的session超时设置。如果你用的是DirectStream,即使Kafka短暂不可用,消费者也会等到broker恢复后再继续,不会丢offset,前提是你手动提交offset的Redis还在。
从那以后我每次做流式项目,都强制自己在启动前先跑一遍最小闭环验证,把Kafka、Redis、Spark三方日志打出来,确认没有ClassNotFound和端口冲突才开始调业务。这次拆的这套资源,只要你按第2章到第6章的顺序走,再踩几个坑,答辩基本稳了。希望帮到你。
本文还有配套的精品资源,点击获取