ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Spark2.x新闻实时分析系统:Kafka流式计算与可视化实践

Spark2.x新闻实时分析系统:Kafka流式计算与可视化实践 简介面向大数据专业毕业设计与Spark初学者的完整项目资源聚焦新闻网场景下实时分析可视化系统的工程实现。资源共35个文件打包后3.43MB涵盖7个Scala和6个Java核心源码、10个依赖JAR包以及XML配置、JS前端页面、HTML展示模板、图片效果图和README说明等目录结构清晰。已有45人学习。内容包含数据采集、预处理、实时分析、结果展示四大模块的完整代码集成文本挖掘、情感分析等算法部署文档覆盖Flume、HBase、Spark环境搭建步骤附参考操作指南自带新闻数据集可直接运行对完成课程设计、毕业设计或研究可复现项目均有直接帮助。整体功能闭环代码规范清晰可在此基础上二次开发是理解Spark分布式处理与可视化落地的实用样例。1. 这个Spark2.x新闻网实时分析系统到底在做什么一个新闻网站每天产生数十万条点击、评论、分享记录运营要的是“此刻正在发生什么”——哪条新闻正在被疯转、哪个关键词热度在飙升、哪个榜单在剧烈变化。传统的T1离线统计根本追不上这个节奏所以就有了基于Spark2.x的新闻网大数据实时分析可视化系统它把Kafka里的新闻行为流数据接进来用Spark Streaming做秒级到分钟级的窗口计算把结果写进Redis供前端可视化大屏拉取。整套东西在毕设里能同时覆盖“采集-计算-存储-展示”四条链路这也是它在毕业设计里拿高分的关键原因。这篇笔记按我实际做过的方案拆开讲链路怎么选型、版本怎么配对、窗口参数怎么设、哪些地方是黑匣子。适合三类人——正在做大数据方向毕业设计的学生、准备Spark实战面试的开发者、以及接了实时看板需求但还没想清楚架构的一线工程师。2. 系统链路拆解Kafka到Spark再到可视化屏的数据流向2.1 为什么链路里要有Kafka削峰与解耦常见做法是先把新闻站点的埋点日志收集到Kafka再由Spark Streaming消费。很多第一次做的人会想“直接用Spark消费数据库不行吗”行但一遇热点新闻就翻车。新闻流量有一个明显特征突发性强某个事件爆发时同一秒内的点击量可能飙升几十倍。如果让Spark直接对接业务库流量尖峰会把数据库连接池打满Spark端也因为没有缓冲而反复失败重启。Kafka在这里起两个作用一是削峰生产者猛写时数据先堆积在Kafka里Spark端按自己的最大消费速率拉取不会被打死二是解耦新闻前端、评论系统、用户行为采集各自往里写Spark不关心上游是谁。参数上有几个值得认真调的acksall保证生产者写入不丢数据retries3配合enable.idempotencetrue处理网络抖动导致的重复topic分区数建议设成Spark消费并行度的2到3倍比如Spark分配5个executor分区数给10到15个这样即使某个executor故障退出剩下的executor还能把分区接管过来。2.2 Spark Streaming的两种接入方式Receiver与DirectSpark 2.x时代Kafka接入有Receiver-based和Direct两种方式。Receiver方式在内部把Kafka数据先存进WAL再交给Spark处理逻辑上绕了一圈而且默认的spark.streaming.receiver.writeAheadLog.enable一旦没打开executor宕机就会丢数据。Direct方式从Spark 1.3开始引入它让Spark直接连接Kafka分区拿offset不做中间缓冲。val kafkaParams Map[String, Object]( bootstrap.servers - node01:9092,node02:9092, key.deserializer - classOf[StringDeserializer], value.deserializer - classOf[StringDeserializer], group.id - news-realtime-group, auto.offset.reset - latest, enable.auto.commit - (false: java.lang.Boolean) ) val stream KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](Array(news-click), kafkaParams) )这段代码是Direct方式的典型写法。createDirectStream不会自动提交offset配合enable.auto.commitfalse由你自己管理消费进度。LocationStrategies.PreferConsistent会在所有executor上均匀分布分区数据避免数据倾斜到个别节点。ConsumerStrategies.Subscribe支持动态发现topic新增分区比Assign固定分区列表更省心。很多毕设版本里用的还是老API如果你手上的源码包出现KafkaUtils.createStream说明是Receiver方式建议改成上面这种。Direct方式在故障恢复、背压控制和exactly-once语义上都比Receiver干净这也是Spark 2.x时代的主流选择。2.3 计算结果存放为什么实时看板选Redis不选MySQL实时计算的产出一段一段写入MySQL也不是不行但可视化大屏是高频读场景前端每几秒拉一次接口MySQL每次都要走SQL解析和磁盘IO扛不住。Redis把热点数据放在内存里读写都在微秒级而且数据结构天然适合排行榜和时序曲线。我一般设计三组Redis键news:hot:rank用ZSET存储新闻热度排行member是新闻IDscore是热度值news:click:trend:{date}用HASH存储按小时粒度的点击量field是小时value是点击次数news:top:{newsId}用STRING存单条新闻的标题、摘要、当前热度前端详情页直接取。写入时的序列化建议用JSON虽然比二进制多占一点空间但对前端JavaScript来说解析零成本。Redis的过期时间要给排行榜设置expire比如保留24小时不然旧数据会越堆越多内存到头后触发淘汰策略把有用的热数据也挤掉。这是实时看板里最常见的“越跑越慢”的元凶之一。3. 用源码包搭建Spark2.x实时分析环境版本配对与最小集群3.1 版本配对JDK、Scala、Spark、Kafka的兼容矩阵拿到毕设源码包的第一件事不是打开IDE而是核对版本。Spark 2.x对版本兼容极其敏感网上大量报错其实就是版本不配对造成的。我自己搭环境时用的组合供参考组件推荐版本说明JDK1.8Spark 2.x官方要求Java 8Scala2.11.xSpark 2.2之前的版本只认Scala 2.11高版本Scala编译的代码跑不了Spark2.4.x2.x系列里最稳定的一个版本Kafka0.10.x 或 1.xSpark 2.4对应kafka-clients 0.10Redis3.x 或 4.x3.x够用4.x支持更多内存策略这里最容易踩的坑是Scala版本。Spark 2.x在发布时会同时出_2.11和_2.12两个编译版本如果你的业务代码是用Scala 2.11编译的却把Spark换成了2.12版本运行时直接报NoSuchMethodError或ClassNotFound。这个错误很迷惑人看起来像是代码问题其实是编译期版本不匹配。还有一个隐藏问题源代码包里如果用了spark-streaming-kafka-0-8这个依赖它对应Kafka 0.8/0.9而你的Kafka如果是2.x的新版连接时会出现协议不兼容。检查pom文件或build.sbt里的依赖坐标确保spark-streaming-kafka的版本和Kafka broker版本在同一代。3.2 Standalone集群最小配置内存参数决定生死Spark 2.x有三种部署模式Local、Standalone、YARN。毕设环境通常只有几台机器Standalone最合适它不需要额外部署HDFS和YARN一条start-master.sh就能拉起来。看似简单但executor的内存参数设不对集群跑起来后你会被各种OOM和Lost executor折磨。# spark-env.sh 关键配置 export JAVA_HOME/usr/local/jdk1.8.0_202 export SPARK_MASTER_HOSTnode01 export SPARK_WORKER_CORES4 export SPARK_WORKER_MEMORY8g export SPARK_EXECUTOR_MEMORY4g export SPARK_DRIVER_MEMORY2g参数含义先说清楚SPARK_WORKER_MEMORY是每台worker节点能给executor分配的总内存上限SPARK_EXECUTOR_MEMORY是单个executor的堆内存。注意executor内存不是越大越好它要减去spark.memory.overhead这部分留给JVM堆外内存的配额当executor内存设到8g时overhead默认是executor内存的10%也就是额外要留0.8g给堆外。提交作业时我习惯在spark-submit里显式指定spark-submit \ --master spark://node01:7077 \ --class com.news.realtime.HotNewsAnalysis \ --executor-memory 4g \ --executor-cores 2 \ --total-executor-cores 8 \ --conf spark.streaming.kafka.maxRatePerPartition1000 \ --conf spark.streaming.backpressure.enabledtrue \ news-realtime.jarmaxRatePerPartition是每个分区每秒最多拉取多少条backpressure.enabledtrue开启背压让Spark根据处理速度自动调整消费速率。这两个配置是防止“消费太快、处理不过来”的关键。很多集群跑一段时间后数据延迟越来越严重就是没开背压Kafka数据堆积在Spark接收端GC时间飙升。3.3 本地跑通完整链路的最小命令序列环境搭好后先别急着跑完整业务用最小链路验证每一段是否通畅。启动顺序是固定的# 1. 启动ZookeeperKafka依赖它做broker元数据管理 /usr/local/kafka/bin/zookeeper-server-start.sh config/zookeeper.properties # 2. 启动Kafka broker /usr/local/kafka/bin/kafka-server-start.sh config/server.properties # 3. 创建新闻点击topic3个分区2个副本 /usr/local/kafka/bin/kafka-topics.sh --create \ --zookeeper node01:2181 \ --replication-factor 2 \ --partitions 3 \ --topic news-click # 4. 启动Spark master和worker /usr/local/spark/sbin/start-all.sh # 5. 启动Redis redis-server /etc/redis/redis.conf这套顺序里有一个容易忽略的细节Kafka 0.10.x之后推荐用--bootstrap-server而非--zookeeper来创建topic但如果你手上的Kafka是0.9及更早版本只能写--zookeeper。源码包里如果依赖的是spark-streaming-kafka-0-8那Kafka大概率是0.9之前的老版本命令要跟着老版本文本来。每条命令起来后用jps检查进程是否存活Kafka进程、QuorumPeerMainZookeeper、Master和Worker都要在列表里。链路启动完往Kafka里手动写几条测试数据/usr/local/kafka/bin/kafka-console-producer.sh \ --broker-list node01:9092 \ --topic news-click然后在Redis里查是否出现计算结果对应的键。这条调试路径能帮你快速定位问题是在采集段、计算段还是存储段而不是等前端大屏一片空白才下手查。4. 新闻实时指标的窗口计算与状态更新核心参数与代码结构4.1 窗口计算batchInterval、windowDuration、slideDuration的关系实时分析里最核心的概念是窗口。Spark Streaming里的时间单位不是“一条数据”而是一批数据batchInterval决定了每隔多久收集一次数据组成一个RDD。窗口操作则是把多个batch的数据聚合起来计算。很多源码包里的新闻热点统计都用到reduceByKeyAndWindow它的三个时间参数必须理解透val clickStream stream.map(record (record.key(), 1L)) val hotCounts clickStream.reduceByKeyAndWindow( (a: Long, b: Long) a b, // 窗口内合并 (a: Long, b: Long) a - b, // 窗口滑出时减掉旧数据 Seconds(60), // 窗口长度统计最近60秒 Seconds(10), // 滑动间隔每10秒滑动一次 2 // 分区数 )Seconds(60)是windowDuration代表每次计算覆盖过去60秒的数据Seconds(10)是slideDuration代表每10秒产出一个新结果。.window(60s).slide(10s)意味着窗口有6个batch的数据重叠。第二个参数(a: Long, b: Long) a - b是反向函数Spark用它来减掉滑出窗口的旧batch避免每个窗口都重新计算全部数据这就是增量计算的原理。参数设置的边界往往出现在这两处一是slideDuration必须是batchInterval的整数倍如果batch是5秒slide设成7秒直接报IllegalArgumentException二是窗口长度不要太大娱乐新闻的热度统计用60秒窗口没问题但如果要统计“近一小时热点”窗口跨度过大中间状态占用的内存也会线性增长这种情况更适合用下面说的mapWithState。4.2 状态计算updateStateByKey还是mapWithState窗口计算适合短时段聚合但“从今天0点开始到现在的累计点击量”这种长周期统计窗口就不合适了——窗口滑出后数据就丢了而且窗口越长中间状态越大。这时候要按key维护跨批次的状态Spark 2.x提供了两个API老牌的updateStateByKey和新一代的mapWithState。// mapWithState方式维护每一条新闻的累计点击量 val stateSpec StateSpec.function( (newsId: String, click: Option[Long], state: State[Long]) { val current state.getOption().getOrElse(0L) click.getOrElse(0L) state.update(current) (newsId, current) } ).timeout(Minutes(30)) val stateStream clickStream.map(record (record.key(), record.value())).mapWithState(stateSpec)StateSpec.function有三个入参key、当前batch内该key的值、以及历史状态。state.update写入新状态.timeout(Minutes(30))表示如果某个新闻ID超过30分钟没有新数据进来状态自动清除防止内存被大量冷数据占满。mapWithState比updateStateByKey的优势是性能更高它只对变化的key做增量更新而updateStateByKey每次都要在所有key上扫描全量状态。在新闻这种key数量持续增长、且大量长尾新闻只有一两次点击的场景下这个差别会被放大到肉眼可见。源码包里如果用的是updateStateByKey建议改成mapWithState代码量差不多但集群CPU和内存占用会明显降下来。4.3 结果下沉Redis写入时的序列化与键设计计算完成的结果要写进Redis这里有个容易翻车的细节Spark Streaming的输出操作是异步的在foreachRDD里直接写Redis每个partition都会创建一个连接如果连接不复用数据量一大Redis连接数直接爆掉。dstream.foreachRDD(rdd - { rdd.foreachPartition(partition - { // 每个partition只创建一个Jedis连接用完关闭 Jedis jedis new Jedis(node03, 6379); partition.forEachRemaining(item - { String newsId item._1(); Long count item._2(); // 用ZADD更新热榜 jedis.zadd(news:hot:rank, count, newsId); // 用HINCRBY累加当天分时数据 String hour LocalDateTime.now().format(DateTimeFormatter.ofPattern(yyyyMMddHH)); jedis.hincrBy(news:click:trend: LocalDate.now(), hour, count); }); jedis.close(); }); });foreachPartition而不是foreach是这里的关键前者让每个partition上的所有数据共享一个Jedis连接后者每一条数据都创建连接性能差异在每秒几千条产出时就是一个数量级。zadd和hincrBy都是原子操作适合并发写同一个key的场景。另一个细节是写入频率。窗口滑动是10秒一次但某个热点新闻可能在滑动间隔内疯狂增长如果只在窗口结束才写Redis大屏上的数字会有最长10秒的“停滞感”。可以用updateStateByKey维护中间状态每两秒读一次当前状态写到Redis牺牲一点Redis写入量换来大屏数据平滑刷新这是我在实际项目里验证过值得的做法。5. Spark实时项目避坑毕业设计里最容易翻车的5个细节5.1 现象Kafka消费不到数据offset一直不提交表现就是Spark应用日志里看不到任何记录Kafka的消费组offset却一直没有变化。原因最常见的是auto.offset.reset配置错误。如果设为latest而消费者组是新建的Spark启动时会从topic的最新offset开始消费此前写入的测试数据全部跳过了。另一个原因是topic分区数小于executor并行度部分executor空转看起来像“消费不到”。解决调试阶段把auto.offset.reset设为earliest确认逻辑没问题后再改回latest。同时用kafka-consumer-groups.sh --describe --group news-realtime-group查看消费组的分区分配情况确认每个分区都有对应的消费者线程。5.2 现象窗口统计结果重复计算数据比预期大好几倍表现同一个新闻的点击量在多个窗口里被重复累加最终数值远大于实际点击量。原因Kafka的生产者重试机制导致消息重复写入或者Spark Streaming的batch在executor故障后重新执行。Direct方式下如果没有手动管理offset应用重启后会从Zookeeper或Kafka里读旧的offset重新消费一批消息。解决Kafka生产端开启enable.idempotencetrue可以消除生产者重试导致的重复。Spark端在手动提交offset前先等结果写入Redis成功后再提交保证“数据处理完成才更新消费进度”。这样即使重启最多重复处理最后一批数据不至于整个窗口全部重新计算。5.3 现象Redis连接被拒或者内存暴涨表现跑着跑着日志里出现JedisConnectionException: Unexpected end of stream或者Redis内存被占满写入开始报OOM错误。原因foreachRDD里每条记录都创建Jedis连接连接数超过了redis.conf里的maxclients限制Redis开始拒绝新连接。另一类原因是Redis键没有设置过期时间新闻热榜、分时曲线这些数据只增不减几小时就把内存吃满了。解决把foreach改成foreachPartition复用连接同时给Redis设置maxmemory 4gb和allkeys-lru淘汰策略。业务上给news:click:trend:*这类键设24小时过期时间给news:hot:rank设12小时过期时间定时任务在每天凌晨清理前一天的历史键。5.4 现象集群模式下executor频繁丢失表现Spark UI上看到一个executor刚启动就退出反复重试日志里有Container killed by YARN for exceeding memory limits或者ExecutorLostFailure。原因executor实际使用的内存超过了申请的内存。Spark执行端的内存包括堆内存和堆外内存如果spark.executor.memory4g而spark.memory.overhead没设置当数据量大或者GC频繁时堆外内存就可能超过系统限制被集群杀掉。解决显式配置spark.executor.memoryOverhead1g给堆外留足空间。同时把spark.executor.extraJavaOptions加上-XX:UseG1GC -XX:MaxGCPauseMillis200减少Full GC。这是一个你只要不主动加配置就永远不知道原因的隐藏参数踩过一次后我养成了固定写三兄弟的习惯memory、memoryOverhead、extraJavaOptions。5.5 现象大屏前端数据不刷新表现ECharts图表第一次有数据但之后一直保持不变刷新页面才更新。原因前端用setInterval定时调用接口但后端接口返回的JSON里带有HTTP缓存头浏览器直接命中了缓存没有真正请求后端。另一个原因是Redis连接池中的连接已过期后端接口查询时报错但前端捕获了异常后静默吞掉。解决后端接口在HTTP响应头里加Cache-Control: no-cache前端在ajax请求里加cache: false。后端查Redis时用带testOnBorrowtrue的连接池配置让连接池在取出连接时先做一次PING检测过期的连接自动丢弃重建。这个坑用curl测试接口时看不出来因为curl默认不缓存只有浏览器会踩到。6. 可视化数据接口的验证技巧从Redis到ECharts一屏跑通标题里写着“源码部署文档全部数据资料”拿到手后我的建议是不要直接改代码先用一条链路验证数据通不通。把Spark Streaming计算结果写进Redis之后先用命令行确认Redis里真的有数据redis-cli zrevrange news:hot:rank 0 10这是检查实时计算链路最直接的办法。如果这个命令能返回按热度排序的新闻ID列表说明数据从Kafka到Spark到Redis这一段是通的剩下的问题全部集中在接口层和前端。如果返回为空就不要去查前端代码问题一定在计算链路里。接口层验证我习惯用curl而不是直接开前端页面curl -H Cache-Control: no-cache http://node03:8080/api/news/hotRank返回JSON后再从后端往ECharts方向排查。一个实用的验证技巧是把ECharts的series数据临时换成写死的假数据如果图表能正常渲染说明前端代码没有毛病问题出在接口数据格式上。如果假数据也不显示那就要检查前端network里是否有报错、option配置是否把data字段写错位置。最后一步就是给接口加日志每5秒打印一次请求参数和返回条数用时间戳对账前端页面上看到的图表数据和Redis里的实际数值。我在自己的毕设里靠这个办法发现了前端传参小时区差8小时的问题——前端把北京时间当UTC发送后端按UTC去查Redis整整差了8个小时的数据。这类时间口径问题在实时系统里极常见只有对账才能抓出来。Spark 2.x实时分析这条路最花时间的不是Spark本身而是数据链路里各类版本冲突、时间口径、连接复用这些小问题。先把最小链路跑通再逐步替换成自己的业务逻辑是我反复验证过最稳的做法。希望这篇笔记里的参数和排错思路能帮你在毕设里少走几个来回。本文还有配套的精品资源点击获取
返回列表