ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Java+Spark2x构建新闻网实时看板:从Kafka到Redis全链路解析

Java+Spark2x构建新闻网实时看板:从Kafka到Redis全链路解析 简介这是一套基于Java与Spark2x构建的新闻网大数据实时分析可视化课程设计项目适合大数据专业学生、Spark入门开发者及需要完成类似毕设或课设的读者。项目覆盖从Flume日志采集、Kafka消息接入、HBase存储到Spark实时处理与前端可视化展示的完整链路源码结构清晰可直接作为二次开发或学习模板。资源包共36个文件以Java/Scala源码为主含10个jar依赖、7个scala、6个java另附2个xml配置、2个js前端脚本、3个png架构图及README、参考步骤txt与LICENSE压缩包仅3.44MB轻量易部署。目前已有916人在CSDN学习下载。通过该资源可掌握Spark2x与HBase/Flume/Kafka的整合方式理解DataFrame、Spark Streaming处理新闻数据的业务逻辑并参考其Web可视化交互图表设计对完整走通大数据实时分析项目极具参考价值。1. 为什么“JavaSpark2x”是新闻网实时看板最务实的组合一个新闻网站的运营看板编辑要的不是昨天发了多少稿而是此刻哪篇稿子正在被疯转、哪个频道的流量在往上冲。这类场景的数据链路是典型的“高吞吐、低延迟、多端展示”点击流源源不断产生计算侧需要在秒级或分钟级给出滚动统计前端大屏再通过接口把结果拉出来。用 Java 写采集、用 Spark2x 做流式聚合、用 Redis 做结果缓冲、用 Spring Boot 提供查询 API是这套体系里最成熟也最容易招聘到人的路线。选择这个组合的原因很直接Java 提供了从模拟器、后端接口到操作工具的完整工程生态Spark2x 在离线与实时两条链路上都能覆盖尤其适合要做“实时分析 可视化”的教学项目或中小型资讯平台。你不需要在架构上追求极致的 Flink 状态管理也不需要自己从 C 语言开始造轮子只要把一个真实的点击流从产生、接入、计算到图表展示完整打通这套系统就能直接回答业务上“现在发生了什么”。文章会按数据进入系统的顺序展开最后落在验证方法和调优参数上过程中所有代码都基于 Java 语言编写可直接落到工程里跑。2. 数据接入用 Java 模拟新闻网点击流喂给 Spark2x2.1 为什么接入层选 Kafka 而不是直连 SocketSpark2x 的 Streaming 模块确实支持从 TCP Socket 读数据用于功能演示很直观但真实新闻网场景下几乎没人这么干。Socket 接入有两个硬伤一是没有缓存能力计算端一旦重启或者在做 checkpoint 恢复这期间的点击事件就全部丢掉了二是没有消费位点管理无法知道“处理到哪一条”。Kafka 的 Topic 天然支持多消费者组同一个 Topic 可以被“实时计算”和“离线补算”各消费一遍这对后面做数据对账和故障重放非常方便。接入层的基本结构是前端或服务端记录用户访问行为写入统一的点击流 TopicSpark2x 的 Receiver 以 Direct 模式消费。Direct 模式不会把数据先放到 Receiver 的存储层而是由 Streaming 的调度器直接向 Kafka 拉取这样既省了一层内存拷贝也让 offset 存在 checkpoint 里而不是 ZooKeeper 中语义更可靠。2.2 点击流消息体的字段设计消息体建议直接用 JSON不要用自定义二进制协议。JSON 可读性强便于在排查问题的时候直接用 kafka-console-consumer 查看内容。针对新闻网站的场景字段应至少覆盖“谁、在什么时间、看了什么、属于哪个频道”但也不要贪多。消息字段如下字段类型说明userIdString匿名用户标识可用 UUID 或会话 IDarticleIdString文章 ID热点统计的主键channelString所属频道如 news、sports、techplatformString来源平台如 pc、app、wechatremoteIpString客户端 IP可用于地域分析tsLong事件发生时间毫秒时间戳urlString点击的完整 URL便于排查脏数据最容易被忽略的是 ts 字段。很多项目直接用 Kafka 自带的消息时间戳但这个时间表示的是“消息进入 Kafka 的时间”一旦采集端有积压或者网络抖动它和用户真实点击时间会偏差很大。在模拟器里显式生成 ts能让后续的窗口计算基于业务时间而不是到达时间。另外 url 字段不要省后面做数据清洗时过滤爬虫和静态资源请求主要靠它。2.3 用 Java 写一个可调速的模拟数据源没有真实业务流量时模拟器是整套系统的“发牌员”。下面这段代码会按固定速率随机生成用户点击事件并把 JSON 消息发送到 Kafka。import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import com.alibaba.fastjson.JSONObject; import java.util.Properties; import java.util.Random; import java.util.UUID; public class ClickStreamSimulator { public static void main(String[] args) throws InterruptedException { Properties props new Properties(); props.put(bootstrap.servers, node01:9092,node02:9092); props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); // acks1 表示 leader 写入成功即返回 // 模拟场景不要求强一致减少延迟。 props.put(acks, 1); // 攒批发送降低小消息带来的网络开销 props.put(linger.ms, 5); props.put(batch.size, 16384); KafkaProducerString, String producer new KafkaProducer(props); Random random new Random(); String[] channels {news, sports, tech, finance}; String[] platforms {pc, app, wechat}; while (true) { JSONObject event new JSONObject(); event.put(userId, UUID.randomUUID().toString().substring(0, 8)); event.put(articleId, String.valueOf(1000 random.nextInt(500))); event.put(channel, channels[random.nextInt(channels.length)]); event.put(platform, platforms[random.nextInt(platforms.length)]); event.put(remoteIp, 192.168. random.nextInt(100) . random.nextInt(255)); event.put(ts, System.currentTimeMillis()); event.put(url, /article?id event.getString(articleId)); producer.send(new ProducerRecord(news-click, event.getString(articleId), event.toJSONString())); Thread.sleep(50); // 50ms 一条约每秒 20 条 } } }代码逻辑分为三部分构造 Producer 参数、循环造数、发送消息。重点说明几个参数linger.ms是消息在缓冲区等待其他消息一起发送的时间设成 5 毫秒可以在吞吐量和延迟之间取得不错的平衡设成 0 则每条消息立即发送机器越多反而越浪费带宽batch.size是批量发送的字节数上限对点击流这种单条不到 200 字节的消息16KB 是稳妥起点。消息的 key 设为 articleId 可以让同一个文章的事件进入同一个分区保证消费端处理时不会发生同 key 乱序。2.4 模拟器与真实采集的差异点模拟器接入到 Topic 后可以用kafka-topics.sh --describe --topic news-click查看分区数和 Leader 分布kafka-console-consumer.sh --bootstrap-server ... --from-beginning --topic news-click直接观察消息内容。要特别注意真实生产环境里的点击流 URL 会带很多跟踪参数比如?utm_sourcexxxfromsinglemessage所以要在接入层做一次清洗把静态资源.jpg、.css、.js和机器请求过滤掉否则统计出来的热点会被搜索引擎爬虫带偏。Spark2x 的消费端也要处理这种情况下一章会在编码逻辑里说明过滤条件。3. Spark2x 流式计算窗口聚合是“实时看板”的发动机3.1 在 Spark2x 里选 Spark Streaming 还是 Structured Streaming这个问题在 Spark 2.3 之后变得需要认真对待。Structured Streaming 是 Spark 2.0 引入的新一代流计算 API基于 DataFrame 模型自动增量执行。但我依然推荐在 Spark2x 的工程化项目里用 Spark Streaming 的 DStream API有两个原因第一Java 版本的 DStream API 使用起来更直观map、filter、reduceByKeyAndWindow 这些算子和 RDD 一脉相承排查问题时能直接看到每个批次的数据分布第二DStream 的 foreachRDD 可以拿到每个微批的 RDD 再自由输出到外部系统对“既写 Redis 又写 HDFS”这种多路输出场景更容易控制。如果以后要迁移到 Spark 3.x再来评估是否转向 Structured Streaming 也不迟。3.2 从 Kafka 消费到窗口聚合的完整 Java 实现下面代码完成四件事从 Kafka 拉取点击事件、过滤脏数据、统计最近 60 秒内各文章的点击量、把结果写回内存表供查询。其中窗口聚合使用reduceByKeyAndWindow它是实时热门统计的核心算子。import org.apache.spark.SparkConf; import org.apache.spark.api.java.JavaPairRDD; import org.apache.spark.api.java.JavaRDD; import org.apache.spark.streaming.Duration; import org.apache.spark.streaming.api.java.JavaDStream; import org.apache.spark.streaming.api.java.JavaPairDStream; import org.apache.spark.streaming.api.java.JavaStreamingContext; import org.apache.spark.streaming.kafka010.ConsumerStrategies; import org.apache.spark.streaming.kafka010.KafkaUtils; import org.apache.spark.streaming.kafka010.LocationStrategies; import org.apache.kafka.common.serialization.StringDeserializer; import com.alibaba.fastjson.JSONObject; import scala.Tuple2; import java.util.Arrays; import java.util.Collection; import java.util.HashMap; import java.util.Map; public class NewsRealtimeAnalyzer { public static void main(String[] args) throws InterruptedException { SparkConf conf new SparkConf() .setAppName(NewsRealtimeAnalyzer) .setMaster(yarn) // 用 Kryo 替代 Java 默认序列化显著降低网络传输开销 .set(spark.serializer, org.apache.spark.serializer.KryoSerializer) .set(spark.streaming.kafka.maxRatePerPartition, 2000); JavaStreamingContext jssc new JavaStreamingContext(conf, new Duration(2000)); jssc.checkpoint(/data/spark-streaming-checkpoint/news); MapString, Object kafkaParams new HashMap(); kafkaParams.put(bootstrap.servers, node01:9092,node02:9092); kafkaParams.put(key.deserializer, StringDeserializer.class); kafkaParams.put(value.deserializer, StringDeserializer.class); kafkaParams.put(group.id, news-realtime-group); kafkaParams.put(auto.offset.reset, latest); kafkaParams.put(enable.auto.commit, false); CollectionString topics Arrays.asList(news-click); JavaDStreamString stream KafkaUtils.createDirectStream( jssc, LocationStrategies.PreferConsistent(), ConsumerStrategies.Subscribe(topics, kafkaParams) ).map(record - record.value()); // 1. 解析 JSON并过滤非文章点击的脏数据 JavaDStreamJSONObject parsed stream .map(NewsRealtimeAnalyzer::parseEvent) .filter(obj - obj ! null obj.getString(url).startsWith(/article) !obj.getString(url).contains(utm_sourcespider)); // 2. 生成 (articleId, 1) 的形式用于计数 JavaPairDStreamString, Long articleOne parsed .mapToPair(obj - new Tuple2(obj.getString(articleId), 1L)); // 3. 滑动窗口窗口 60 秒滑动 10 秒 JavaPairDStreamString, Long windowedCounts articleOne .reduceByKeyAndWindow((a, b) - a b, new Duration(60000), new Duration(10000)); // 4. 输出到 Redis详见第四章 windowedCounts.foreachRDD(rdd - { JavaPairRDDString, Long sorted rdd.sortByKey(false); sorted.foreachPartition(NewsRealtimeAnalyzer::writeToRedis); }); jssc.start(); jssc.awaitTermination(); } private static JSONObject parseEvent(String line) { try { return JSONObject.parseObject(line); } catch (Exception e) { return null; } } private static void writeToRedis(IteratorTuple2String, Long partition) { // 在 executor 端初始化 Redis 连接见 4.4 } }这段代码的链路是createDirectStream拉取消息 →map转 JSON →filter清脏 →mapToPair组合键值 →reduceByKeyAndWindow做滑窗累加 →foreachRDD输出。参数里最值得关注的是maxRatePerPartition这个值限制每个分区每秒最多消费 2000 条防止处理速度跟不上消费速度时把 executor 内存打爆。enable.auto.commit设为 false配合后续的 checkpoint 机制可以做到“至少一次”的语义宁可重复统计不能漏统计。3.3 三个时间参数的关系new Duration(2000)是批次间隔new Duration(60000)是窗口长度new Duration(10000)是滑动步长。三者的配合规则如下参数含义本项目的推荐值影响batchDurationRDD 生成的频率2 秒越小延迟越低但调度开销越大windowDuration一次统计覆盖的时间范围60 秒决定“热门”有多热slideDuration两次统计之间的间隔10 秒决定大屏数据刷新的密度必须注意windowDuration 和 slideDuration 都应该是 batchDuration 的整数倍。如果 slide 设为 7 秒而 batch 是 2 秒Spark 会直接抛出 IllegalArgumentException。从业务视角理解这三个参数热点榜反映的是“最近一分钟的文章热度”每 10 秒刷新一次大屏每 10 秒感受到一次数字跳动既不会觉得卡顿也不会因过于频繁而制造焦虑。3.4 序列化与 checkpoint 的踩坑点Spark Streaming 作业在集群上运行时闭包里的每个对象都要在 driver 和 executor 之间传输。如果自定义的统计类没有实现Serializable作业会在启动时报Task not serializable。常见的解决方式是使用 Kryo 序列化并注册需要用到的类。上面代码里set(spark.serializer, org.apache.spark.serializer.KryoSerializer)后还建议追加conf.set(spark.kryo.registrationRequired, true)并手动注册类否则遇到未注册的类会直接报错虽然后麻烦但能从一开始暴露问题。checkpoint 目录里保存了两类信息RDD 的血缘DAG和 Offset 状态。一旦你修改了消费逻辑里算子的结构例如把mapToPair和reduceByKeyAndWindow的顺序调换Spark 在恢复时会发现 operator ID 对不上报 “Aborting old checkpoint”。这是本算子在恢复机制上一个很容易踩的坑调整业务逻辑后建议先停掉作业清空 checkpoint 目录再重新提交否则恢复过程会比冷启动更慢。4. Redis 存储层让大屏查询不再直面 Spark4.1 为什么计算结果要先进 Redis而不是让前端直接查 SparkSpark 计算出的结果如果直接通过 JDBC 或 Thrift 暴露给前端会带来两个问题一是每次大屏刷新都触发一次 Spark 作业或 SQL 查询集群资源白白消耗在低频的查询上二是 Spark 的输出结果在批次之间是断续的前端的压力测试很容易打到“数据还没算完”的中间状态。引入 Redis 作为结果缓冲层后Spark 只管“算完就写”查询服务只管“读 Redis 返回”两端互不阻塞。这也是“实时分析 可视化”架构里最常见的做法。4.2 键设计与过期策略不同的统计口径需要使用不同的 Redis 数据结构。点击总量用 String热点排名用 ZSet频道分布用 Hash这样每个查询接口对应一种数据结构的原生操作性能最好。用途Key 模式类型过期时间文章当日 PV 合计news:pv:article:{yyyyMMdd}Hash48 小时全网最近 1 分钟 PVnews:trend:total:1mString90 秒频道最近 1 分钟分布news:trend:channel:1mHash90 秒文章热度 TopNnews:hot:articles:1mZSet90 秒注意过期时间的设计。因为窗口是 60 秒滑动间隔 10 秒所以键的过期时间设为 90 秒给“最后一秒的数据还有机会被刷新”留出缓冲。超过 90 秒没有新数据写入键自动消失查询端读到空值后返回空数组前端就会展示占位图而不是报错。如果用固定不变的 key还需要一个后台任务去清理过期键现在用 Redis 自身的 TTL 机制就能搞定少写一个调度器。4.3 Spark 端写 Redis 的 Java 工具类在窗口聚合结果的foreachRDD里写 Redis注意不要在每个元素上创建 Jedis 实例而是按分区批量操作。下面代码展示了在 executor 端用 pipeline 批量更新键值的方法。import redis.clients.jedis.Jedis; import redis.clients.jedis.JedisPool; import redis.clients.jedis.Pipeline; import java.io.Serializable; import java.text.SimpleDateFormat; import java.util.Date; public class RedisSink implements Serializable { private static JedisPool pool; public static void init() { if (pool null) { // 连接池参数可按集群规模调整 pool new JedisPool(node03, 6379); } } public static void writeHotRank(String articleId, long count) { init(); try (Jedis jedis pool.getResource()) { String windowKey news:hot:articles:1m; Pipeline pipe jedis.pipelined(); pipe.zadd(windowKey, count, articleId); // 只保留前 50防止 ZSet 无限增长 pipe.zremrangeByRank(windowKey, 0, -51); pipe.expire(windowKey, 90); pipe.sync(); } } }这段代码有两个细节值得讲。第一zadd的 score 使用 count而 Spark 窗口统计的 count 是“窗口内累加值”所以每次窗口滑动都会用新值覆盖同一篇文章的旧 score这是 ZSet 做 TopN 的标准写入方式。第二zremrangeByRank保留 Top50防止某些蹭热点刷量的文章长期占据内存。如果业务要求看一小时榜key 换成news:hot:articles:60m同时把窗口统计的时间范围同步调整。4.4 executor 端初始化 Redis 连接的正确姿势很多从 DStream 写外部存储的代码会在 driver 端创建一个 JedisPool 实例然后在foreachRDD的闭包里直接使用。这在本地模式没问题但到了 YARN 集群上driver 的 JedisPool 对象会被序列化并发送给每个 executor每个 executor 的反序列化过程会重新创建连接池导致连接数等于“executor 数量 × task 并发数”很容易把 Redis 的连接数打到上限。正确的做法是像上面的代码一样在foreachPartition的迭代器内部用静态方法懒加载连接池。静态变量属于 executor 上的 JVM 实例不会被反复序列化同一个 executor 上的多个 task 共享这个连接池连接数就被控制在“executor 数 × 池大小”的范围内。如果项目用的是 Spring 管理 Redis 客户端可以把 JedisPool 定义在 executor 启动时执行的初始化类里效果相同。5. Java 后端查询服务把 Redis 数据优雅地交给可视化大屏5.1 构建基于 Spring Boot 的 REST 查询接口到这一步Spark 已经把计算结果写入 Redis后端要做的就是提供一组“无状态、只读”的查询接口。无状态的意思是接口不保存任何会话数据每次请求直接读 Redis这样前端大屏可以同时部署多个实例并挂负载均衡不会出现不同节点读到不同缓存的情况。接口设计遵循“一屏一接口”的原则大屏上的每个块对应一个编辑好的 JSON 结构返回。常见接口如下接口路径功能对应的 Redis 结构GET /api/trend/total最近 1 分钟全网点击趋势StringGET /api/stats/channel各频道实时流量占比HashGET /api/hot/articles?limit20当前热门文章 TopNZSetGET /api/stats/source来源平台分布pc/app/wechatHash5.2 Controller 实现与 Redis 查询代码下面以热度榜接口为例展示 Controller 层的实现。重点是把 Redis 的 ZSet 查询结果转换成前端可渲染的 JSON 格式同时把读不到的 key 处理成空列表而不是抛异常。import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.data.redis.core.ZSetOperations; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import java.util.*; import java.util.stream.Collectors; RestController public class RealtimeStatController { private final StringRedisTemplate redisTemplate; public RealtimeStatController(StringRedisTemplate redisTemplate) { this.redisTemplate redisTemplate; } GetMapping(/api/hot/articles) public MapString, Object hotArticles(RequestParam(defaultValue 20) int limit) { String key news:hot:articles:1m; // 按 score 降序取前 limit 名 SetZSetOperations.TypedTupleString tuples redisTemplate.opsForZSet().reverseRangeWithScores(key, 0, limit - 1); ListMapString, Object items new ArrayList(); if (tuples ! null) { items tuples.stream().map(tuple - { MapString, Object item new HashMap(); item.put(articleId, tuple.getValue()); item.put(score, Objects.requireNonNull(tuple.getScore()).longValue()); return item; }).collect(Collectors.toList()); } MapString, Object result new HashMap(); result.put(code, 0); result.put(data, items); result.put(updateTime, System.currentTimeMillis()); return result; } }这个接口返回的数据结构做了两件事一是把 Redis 的分数从 Double 转成 Long避免前端拿到3.0000000001这种浮点噪声二是把updateTime放进响应体里前端可以根据这个时间戳判断数据是否还在更新超过 90 秒没变化就展示“数据链路中断”的重试提示而不是显示一个静止的假数据。真实工程里还可以在这个环节加一层 Caffeine 本地缓存例如每 2 秒刷新一次减少对 Redis 的访问压力。对于大屏 5 秒轮询的频次不加也完全够用。5.3 可视化前端的渲染与连接方式后端只负责把数据给到前端具体渲染交给 ECharts 或 AntV。在新闻网站的看板场景里需要展示的不只是“当前热度”还有一个小时内的趋势折线图。这时前端需要把 Redis 中多个时间窗口的键组合处理例如news:trend:total:1m、news:trend:total:5m、news:trend:total:30m各自保存一个值前端按时序排列这些键就能画出一条“1 分钟 / 5 分钟 / 30 分钟”粒度的时间线。注意粒度从细到粗的变化本质上是查询不同的键而不是对同一个键做聚合这样设计后端接口最简单。5.4 实时数据与离线数据的口径对齐做新闻热点看板一定会遇到一个尴尬实时榜显示的 Top1 和离线统计的 Top1 对不上。原因不在计算逻辑而在于实时链路把爬虫过滤了一部分离线链路则全量统计。建议在清洗层统一过滤规则这里把url包含utm_sourcespider的行为一律视为爬虫在两个链路里共用同一个过滤维度这样才能保证“吃瓜群众看的数据”和“运营同学复盘的数据”是一套口径。6. 用一条测试消息验证实时链路顺带调优背压参数拿到这套系统时第一步不是去调优而是把验证方法固定下来。用模拟器只向 Kafka 发送一条带有特殊标记的文章点击事件比如articleId1001, channelsports然后观察大屏的体育频道是否在 60 秒窗口的统计值中增加了 1。如果一切正常这个标记文章的分数会出现在热度榜尾部。如果 30 秒后依然看不到优先检查三处Kafka 的 Topic 是否收到了这条消息Spark 作业是否有报错日志Redis 里news:trend:channel:1m键是否已生成。用redis-cli hgetall news:trend:channel:1m可以直接确认当前缓存值比盯着大屏盲猜高效得多。验证通过后关注流处理的两个关键指标批处理时间和调度延迟。在 Spark UI 的 Streaming 页签里可以看到每一批的 Scheduling Delay 和 Processing Time。当 Processing Time 稳定大于批量间隔本例为 2 秒时说明系统处理能力不足。优先调节 Kafka 消费速率上限spark.streaming.kafka.maxRatePerPartition和背压开关推荐值如下参数默认值推荐值说明spark.streaming.backpressure.enabledfalsetrue根据处理能力动态调整消费速率spark.streaming.backpressure.initialRate无500重启后消费速率的起点spark.streaming.kafka.maxRatePerPartition无2000单分区每秒最大消费条数spark.streaming.kafka.consumer.poll.ms5121000单次拉取超时时间背压机制在 Spark 2.x 里已经稳定开启后系统会根据上一批的处理时间自动调整读取速率避免数据积压引发 OOM。如果集群内存紧张还可以把spark.streaming.kafka.maxRetries设为 3控制因网络抖动导致的拉取重试次数。调优完成后把模拟器的发送频率从每秒 20 条逐步提到每秒 200 条对比观察大屏刷新频率和 Spark UI 的 Processing Time就能画出一条本集群的“最大吞吐-延迟”曲线。本文还有配套的精品资源点击获取
返回列表