简介:基于Java与Spark2x技术栈的新闻网大数据实时分析可视化系统项目,面向大数据相关课程设计、毕业设计及需要掌握实时处理链路的开发者。项目围绕新闻数据采集、流式处理、结果存储与Web可视化展开,可帮助理解从Kafka接入、Spark流计算到HBase存储、前端图表展示的完整实现方案。压缩包共36个文件,约3.44MB,主要包含10个Java归档包、7个Scala源文件、6个Java源文件,以及XML配置、JavaScript脚本、Maven配置和参考步骤文档,既有可直接运行的依赖库,也有便于二次开发的源码结构。同时包含Flume对接HBase的序列化器、Spark流处理逻辑与Web图表交互相关代码,并附带项目截图与README说明,便于按模块学习、对照搭建。目前已有917人学习下载,适合希望提升Spark2x实战能力、学习大数据可视化系统搭建的初中级学习者。
1. 为什么“Java + Spark2x + 新闻网实时可视化”值得照着做一遍
看到一个“基于Java实现Spark2x新闻网大数据实时分析可视化系统”的题目,很多人的第一反应是“又是一套毕业设计模板”。但如果你愿意把这个方向按生产原型的标准重做一遍,你会发现它恰好覆盖了实时数据链路上最关键的四段:Java 生产端、Kafka 消息管道、Spark2x 流式计算、可视化大屏对接。适合正在选毕业设计题、或者从 Java 后端转大数据实时方向的开发者。这篇笔记不介绍某个现成源码包,而是讲清这条链路每一步怎么选型、怎么写、参数怎么设、以及哪些坑是绕不过去的。
2. 实时管道设计:Kafka到Spark2x的Java实现怎么选型
新闻网站最典型的实时场景是用户点击流:每一次曝光、点击、评论都是一条事件。做实时分析的第一步不是急着写 Spark 代码,而是先决定事件从哪来、以什么格式进到计算引擎,以及计算引擎用哪套 API 接数据。顺序反了,后面每一步都在打补丁。
2.1 模拟数据源与Kafka Topic设计
新闻站没有现成的埋点流量时,常见做法是写一个 Java 模拟器,按固定速率往 Kafka 里灌 JSON 事件。这个模拟器同时也是后面压测的“发压工具”,别写完就删:
import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; import java.util.UUID; public class NewsClickProducer { public static void main(String[] args) throws Exception { Properties props = new Properties(); props.put("bootstrap.servers", "node01:9092,node02:9092"); props.put("key.serializer", StringSerializer.class.getName()); props.put("value.serializer", StringSerializer.class.getName()); props.put("acks", "1"); // 允许少量丢失,换取吞吐 KafkaProducer<String, String> producer = new KafkaProducer<>(props); String[] categories = {"tech", "sports", "finance", "ent"}; while (true) { String json = String.format( "{\"newsId\":\"%s\",\"cat\":\"%s\",\"uid\":\"%s\",\"ts\":%d}", UUID.randomUUID().toString().substring(0, 8), categories[(int) (Math.random() * categories.length)], UUID.randomUUID().toString().substring(0, 6), System.currentTimeMillis() ); // key 用随机串而不是 newsId,避免某个新闻成为热点分区 producer.send(new ProducerRecord<>("news_click", UUID.randomUUID().toString(), json)); Thread.sleep(5); // 约 200 TPS,够一个演示项目用 } } }逻辑说明:每条消息包含新闻 ID、栏目、用户 ID、时间戳四个字段,覆盖后续聚合和可视化所需的最少维度。acks=1是吞吐和可靠性的折中,模拟场景没问题,生产环境看业务对丢失的容忍度。
参数说明:Topic 的partitions建议设成 8 或 16,理由在第 5 章讲并行度时会说;replication-factor至少 2,单机演练可以 1。分区数定小了,后面想扩容要动生产端和消费端,很麻烦;定大了先期成本也不高,所以给模拟数据时预留 8 个分区是常见起点。
2.2 Spark2x Java API:JavaStreamingContext与DStream的取舍
Spark2x 时代的实时计算有两个入口:老的 Spark Streaming(DStream)和新的 Structured Streaming。用 Java 做这个项目时,我建议优先盯住JavaStreamingContext+DStream这套组合,原因是它在 Spark2x 里和 Kafka 0-10 客户端的配合最成熟,网上能查到的案例也最多;Structured Streaming 的 Java 版 API 在 2.x 早期版本里抽象层绕,调试 event-time 窗口时反而容易翻车。
另外这也是 Java 面试八股里常出现的区分点:DStream 是一串连续的 RDD,本质是微批次,Spark2x 里通过Durations.seconds(n)切分;Structured Streaming 用 DataFrame/Dataset 的声明式 API,不直接暴露 RDD。对这个新闻网项目来说,用 DStream 能更直观地看到“每个批次在处理什么”,排查问题时心理负担小。
初始化代码是整个作业的地基:
import org.apache.spark.SparkConf; import org.apache.spark.streaming.Durations; import org.apache.spark.streaming.api.java.JavaStreamingContext; SparkConf conf = new SparkConf().setAppName("NewsClickStreaming"); // 本地模式仅用于功能联调,集群提交时这行要注释掉 conf.setMaster("local[2]"); JavaStreamingContext jssc = new JavaStreamingContext(conf, Durations.seconds(10)); jssc.checkpoint("hdfs://nameservice/spark-checkpoint/news");逻辑说明:local[2]里的 2 表示本地用两个线程跑,一个收数据、一个处理,联调时够用;集群模式删除这行,资源由 YARN 统一分配。checkpoint是窗口聚合的硬性要求,DStream 跨批次维护状态时会把元数据写到这里。路径必须用 HDFS 或可靠文件系统,本地/tmp一重启就没了,后面会变成“玄学丢数据”。
参数说明:批次间隔Durations.seconds(10)决定实时延迟上限 10 秒。不是越小越好,批次太密而处理不过来,反而会积压成雪崩。具体配法在第 5 章展开。
2.3 最小作业代码:窗口内热度统计的完整Java实现
下面这段是能直接抄走改的完整核心逻辑:读 Kafka、解析 JSON、按 60 秒窗口统计每个栏目的点击量。ES 输出部分放到下一章,这里先用print()验证管道通不通:
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.spark.api.java.Optional; import org.apache.spark.streaming.Durations; import org.apache.spark.streaming.api.java.*; import org.apache.spark.streaming.kafka010.*; import scala.Tuple2; import java.util.*; import java.util.function.Function; public class NewsClickStreaming { public static void main(String[] args) throws InterruptedException { String brokers = args.length > 0 ? args[0] : "node01:9092"; String groupId = args.length > 1 ? args[1] : "news_rt_group"; SparkConf conf = new SparkConf().setAppName("NewsClickStreaming"); JavaStreamingContext jssc = new JavaStreamingContext(conf, Durations.seconds(10)); jssc.checkpoint("hdfs://nameservice/spark-checkpoint/news"); Map<String, Object> kafkaParams = new HashMap<>(); kafkaParams.put("bootstrap.servers", brokers); kafkaParams.put("key.deserializer", StringDeserializer.class.getName()); kafkaParams.put("value.deserializer", StringDeserializer.class.getName()); kafkaParams.put("group.id", groupId); kafkaParams.put("auto.offset.reset", "latest"); kafkaParams.put("enable.auto.commit", false); Collection<String> topics = Collections.singletonList("news_click"); JavaInputDStream<ConsumerRecord<String, String>> stream = KafkaUtils.createDirectStream( jssc, LocationStrategies.PreferConsistent(), ConsumerStrategies.Subscribe(topics, kafkaParams) ); JavaDStream<String> jsonLines = stream.map(ConsumerRecord::value); // 这里用最简单的 split 解析,项目里换成 fastjson 或 jackson 都行 JavaPairDStream<String, Long> catCount = jsonLines .mapToPair(line -> { String cat = line.split("\"cat\":\"")[1].split("\"")[0]; return new Tuple2<>(cat, 1L); }) .reduceByKeyAndWindow( (a, b) -> a + b, (a, b) -> a - b, Durations.seconds(60), Durations.seconds(10) ); catCount.print(); jssc.start(); jssc.awaitTermination(); } }逻辑说明:mapToPair把每条 JSON 拆成“(栏目, 1)”的键值对;reduceByKeyAndWindow是 DStream 窗口聚合的核心,第一个 lambda 是窗口内累加,第二个 lambda 是滑出窗口时扣减,这种“加新减旧”的写法比每个窗口重新全量计算省资源,代价是必须开 checkpoint。print()默认打印 10 条,联调时够看,正式环境会被日志刷屏,要及时换掉。
参数说明:窗口长度 60 秒、滑动间隔 10 秒,意味着大屏上每 10 秒刷新一次最近一分钟的热度。如果业务要看“最近 5 分钟的热点新闻”,窗口改成 300 秒即可。这里有个容易忽略的细节:窗口长度必须是批次间隔的整数倍,否则 Spark2x 会直接报参数非法。
3. 可视化数据层怎么搭:ES索引、Java聚合接口与前端图表对接
计算层跑通只是第一步,实时系统的价值要在大屏上体现。可视化层常见做法是“Spark 算完写 Elasticsearch,后端查 ES 出接口,前端 ECharts 渲染”。这套组合能覆盖新闻热榜、栏目占比、时间趋势三类最常见的大屏图表。
3.1 存储选型:实时聚合结果为什么别直接写MySQL
新闻点击流的特征是事件量大、聚合维度固定、查询条件集中在时间范围和分类上。如果让 Spark 每 10 秒把结果怼进 MySQL,热点行会被频繁更新,行锁竞争很快变成瓶颈;而大屏要的“最近 N 分钟趋势”在 MySQL 里要写复杂 group by,索引怎么建都别扭。
ES 的优势在两点:一是date_histogram聚合能直接按分钟切桶,二是一次查询可以同时出总数、TopN、占比多个结果。代价是要多维护一个集群,但对这类可视化系统来说是值得的。还有一种常见折中是 Spark 算结果写 MySQL、ES 只存明细,我一般不建议入门项目这么拆,双写一致性会变成新的“坑王”。
3.2 索引设计与Java批量写入参数
索引设计直接决定查询好不好写。我通常按天滚动索引,比如news_click_rt_20250217,字段只有四个:
| 字段 | 类型 | 说明 |
|---|---|---|
| category | keyword | 栏目名,聚合用 |
| newsId | keyword | 新闻 ID,TopN 用 |
| windowStart | date | 窗口起始时间,趋势图用 |
| cnt | long | 窗口内点击数 |
主分片设 5、副本 1,对演示项目足够。写入端用BulkRequest攒批,避免每条一次 HTTP 请求:
import org.apache.http.HttpHost; import org.elasticsearch.action.bulk.BulkRequest; import org.elasticsearch.action.index.IndexRequest; import org.elasticsearch.client.RequestOptions; import org.elasticsearch.client.RestHighLevelClient; import org.elasticsearch.common.xcontent.XContentType; import scala.Tuple2; import java.util.HashMap; import java.util.Iterator; import java.util.Map; public class EsBulkWriter { private static final RestHighLevelClient CLIENT = new RestHighLevelClient( org.elasticsearch.client.RestClient.builder( new HttpHost("node03", 9200, "http") ).build() ); public static void write(Iterator<Tuple2<String, Long>> partition) throws Exception { BulkRequest bulk = new BulkRequest(); int count = 0; while (partition.hasNext()) { Tuple2<String, Long> item = partition.next(); Map<String, Object> doc = new HashMap<>(); doc.put("category", item._1); doc.put("cnt", item._2); doc.put("windowStart", System.currentTimeMillis()); bulk.add(new IndexRequest("news_click_rt") .source(doc, XContentType.JSON)); count++; if (count >= 500) { // 攒够 500 条批量提交一次 CLIENT.bulk(bulk, RequestOptions.DEFAULT); bulk = new BulkRequest(); count = 0; } } if (count > 0) { CLIENT.bulk(bulk, RequestOptions.DEFAULT); } } }逻辑说明:CLIENT定义成静态单例,避免每个分区都建连接导致 YARN 容器内存爆炸。windowStart直接用System.currentTimeMillis()取的是处理时间,如果模拟数据里的ts字段是事件时间,这里优先用ts所在的分钟桶,语义才一致。
参数说明:批量大小 500 条在这个量级够用。生产环境一般看“批字节数”,ES 官方建议攒到 5MB 左右提交一次,或者以 1000 条为下限,二选一。批量太大反而会触发 ES 的写入队列积压,返回 429 拒绝。
3.3 从ES聚合到前端图表:时间范围与热点TopN
后端接口建议按大屏组件拆,不要做一个万能查询接口。最常用的是趋势接口和 TopN 接口。趋势接口用 ES 的date_histogram:
import org.elasticsearch.action.search.SearchRequest; import org.elasticsearch.action.search.SearchResponse; import org.elasticsearch.index.query.QueryBuilders; import org.elasticsearch.search.aggregations.AggregationBuilders; import org.elasticsearch.search.aggregations.BucketOrder; import org.elasticsearch.search.aggregations.bucket.histogram.DateHistogramInterval; import org.elasticsearch.search.builder.SearchSourceBuilder; public Map<String, Object> trend(String category, long from, long to) throws Exception { SearchRequest request = new SearchRequest("news_click_rt_*"); SearchSourceBuilder src = new SearchSourceBuilder(); src.query(QueryBuilders.boolQuery() .filter(QueryBuilders.termQuery("category", category)) .filter(QueryBuilders.rangeQuery("windowStart").from(from).to(to))); src.size(0); src.aggregation(AggregationBuilders.dateHistogram("timeline") .field("windowStart") .fixedInterval(DateHistogramInterval.minutes(1)) .subAggregation(AggregationBuilders.sum("total").field("cnt"))); request.source(src); SearchResponse response = CLIENT.search(request, RequestOptions.DEFAULT); // 解析 buckets 里的 key 和 total,按时间升序返回给前端 return buildTimelineResult(response); }逻辑说明:size(0)表示只要聚合结果,不回明细;外层用filter而不是must,因为不需要算相关性分数,性能更好。dateHistogram按 1 分钟一个桶,配合sum子聚合,就能画出“每分钟点击数折线”。
前端对接时,ECharts 的配置核心就一段,难点在于后端返回的时间戳要和前端xAxis的类型对牢:
// 后端返回 [{time: 1716000000000, value: 34}, ...] option = { xAxis: { type: 'time' }, // x 轴声明为时间轴 yAxis: { type: 'value' }, series: [{ type: 'line', smooth: true, areaStyle: { opacity: 0.2 }, data: resp.map(item => [item.time, item.value]) }] };参数说明:smooth会让曲线变光滑,但也会掩盖瞬时抖动,做实时监控时建议改成false;xAxis.type = 'time'要求数据点是“毫秒时间戳 + 数值”的二元数组,如果后端返回字符串时间,ECharts 会解析失败,这是可视化联调最常见的翻车点。
4. 避坑实录:Spark2x实时任务最常见的5类翻车现场
这一章是血泪经验合集。实时任务和离线任务最大的区别是“跑起来容易,稳着跑难”,下面 5 个问题是我在多个类似项目里反复撞过的墙,每条都按“现象 → 原因 → 解决”的顺序写。
4.1 现象:Kafka消费延迟只增不减
现象:Kafka 监控面板上lag数值持续走高,大屏上的指标比真实时间滞后越来越多,最终任务整批整批地积压。
原因:最常见的是没开背压,或者单条记录处理太重。默认情况下 Spark Streaming 会按最大速率拉数据,处理不过来就把批次排队,形成雪崩。另一个隐蔽原因是每条数据都同步写 ES,网络往返把处理时间放大。
解决:先确认瓶颈。看 Spark UI 里每个 batch 的“Processing Time”是否超过批次间隔,超过就说明 CPU 或 IO 满了。然后打开两个开关:
--conf spark.streaming.backpressure.enabled=true \ --conf spark.streaming.kafka010.maxRatePerPartition=2000逻辑说明:背压开启后,Spark 会根据上一批次的处理耗时动态调整拉取速率;maxRatePerPartition是单分区每秒拉取上限的兜底,防止流量尖峰直接打垮作业。这两个参数是实时任务的“后悔药”,但不要指望它们解决所有问题,ES 写入改成批量(第 3 章的做法)才是根治。
4.2 现象:窗口统计结果对时对错
现象:大屏上同一时间点的数值多次刷新不一致,和离线数仓对账时永远差几分钟的数据。
原因:DStream 老 API 不原生支持 event-time 窗口,它默认按“处理时间”切窗口。如果模拟器里的ts是事件发生时间,而你的窗口聚合用的是 Spark 收到数据的时间,那两条时间线就对不上。还有一种是时区问题:ES 里windowStart按 UTC 存,前端按北京时间查,查出来少 8 小时。
解决:要么承认用处理时间,把模拟器里的ts改成生成时间;要么在map阶段把事件时间按分钟切桶,把桶 id 当作 key 的一部分。对新闻网这种对时间精度要求不高的场景,我建议直接用处理时间,会少踩很多坑。时区问题则在写入前统一转成东八区毫秒,前端的时区设置和后端保持一致。
4.3 现象:Executor 反复被 killed
现象:提交后任务跑几分钟就挂,YARN 日志里反复出现Container killed by YARN for exceeding memory limits,重试多少次都没用。
原因:绝大多数情况是spark.executor.memoryOverhead太小。Spark2x 默认按 executor 内存的 10% 计算堆外内存,但 ES 客户端、Kafka 消费者、网络缓冲都在堆外分配,很容易超限。YARN 不会等你慢慢涨,一旦超额直接杀容器。
解决:显式调大堆外内存,同时看 GC 日志确认不是堆内问题:
--conf spark.executor.memoryOverhead=1024 \ --conf spark.executor.extraJavaOptions="-XX:+UseG1GC -XX:MaxGCPauseMillis=200"逻辑说明:memoryOverhead设为 1GB 对多数 Spark2x 实时任务够用;如果还杀,再查是不是 RDD 缓存太多。cache()用多了会占满堆内存,改用persist(StorageLevel.MEMORY_AND_DISK())让溢写落盘。
4.4 现象:ES批量写入报429和版本冲突
现象:Spark 日志里刷EsRejectedExecutionException,重试后部分数据写进,部分丢失;或者索引里出现重复数据,总数对不上。
原因:429 是 ES 写入了,常见于批量线程过多、单批过大。版本冲突则是因为代码里给 doc 指定了固定 id,同一个窗口重复写入时主键撞车,默认版本号不一致就会被拒绝。
解决:批量线程数不要超过 ES 节点数的两倍;如果每条 doc 的内容是“窗口聚合结果”,根本不需要固定 id,让 ES 自动生成_id即可。如果业务必须覆盖旧窗口,在 IndexRequest 上配置版本,或者干脆用windowStart + category拼一个幂等 id,并用 upsert 语义。
4.5 现象:检查点恢复后数据重复或丢失
现象:任务正常跑了几天,重启后大屏数值跳变,要么比停之前多了一截,要么少了一截。
原因:checkpoint 里保存的 Spark 内部状态和 Kafka 偏移量没有对齐。enable.auto.commit=false是对的,但代码里如果从没提交过 offset,重启时 Spark 会从 checkpoint 恢复内部状态,而 Kafka 侧偏移还在旧位置,两边的进度不一致。
解决:用spark-streaming-kafka-010自带的提交机制,在处理完一个批次之后主动提交:
stream.foreachRDD((JavaRDD<ConsumerRecord<String, String>> rdd, Time time) -> { OffsetRange[] offsetRanges = ((HasOffsetRanges) rdd.rdd()).offsetRanges(); // 业务输出完成后,异步提交偏移 ((CanCommitOffsets) stream.inputDStream()).commitAsync(offsetRanges); });逻辑说明:HasOffsetRanges可以从当前批次 RDD 里取出消费到的偏移范围,commitAsync是异步提交,不阻塞批次处理。这段代码只适合放在“结果已安全输出”之后,如果先提交偏移、再写 ES 且写失败了,数据就真丢了。
5. 参数联动调优:把批次、并行度、内存和检查点设成一套
很多新手调参是一次只调一个参数,结果越调越乱。实时任务的参数是联动的:批次间隔决定窗口能取多细,Kafka 分区数决定并行度上限,内存决定窗口状态能撑多大。下面按三个维度讲我怎么配。
5.1 批次间隔与窗口长度怎么配对
规矩只有一个:窗口长度和滑动间隔都必须等于批次间隔的整数倍。推荐配置:
| 场景 | 批次间隔 | 窗口长度 | 滑动间隔 |
|---|---|---|---|
| 大屏秒级刷新 | 5s | 60s | 5s |
| 热点新闻 TopN | 10s | 300s | 10s |
| 准实时运营看板 | 30s | 1800s | 300s |
批次间隔不是越小越好。5 秒批次意味着 Spark 每 5 秒做一次任务调度和状态快照,调度开销占比上升;30 秒批次则让大屏看起来“卡顿”。我一般从 10 秒起步,压测后如果 batch 处理时间低于 5 秒,再考虑降成 5 秒。窗口越长,需要维护的状态越大,checkpoint 的压力也越大,所以不要随手设 24 小时窗口。
5.2 并行度与Kafka分区怎么匹配
Kafka Topic 的分区数决定了KafkaUtils.createDirectStream读出来的 RDD 分区数。换句话说,Topic 设 8 个分区,Spark 读端最多 8 个并行任务。常见做法是让 Topic 分区数等于目标 executor 核数,或略大一点(1.2 倍)。注意 executor 核数是指并发任务总数,不是节点数。
如果后续处理算子的并行度不够,比如某条新闻热度特别高,导致按newsId聚合时数据倾斜,一个分区的数据比其他分区大一个数量级,这属于实时任务的经典倾斜问题。解决办法是在聚合前加盐(salted key):
// 加盐打散热点:key 变成 "newsId_随机数%10",聚合后再去掉盐 JavaPairDStream<String, Long> salted = clicks .mapToPair(click -> new Tuple2<>(click.getNewsId() + "_" + (int)(Math.random() * 10), 1L)); // 先按加盐 key 局部聚合,再按原 key 汇总逻辑说明:加盐会让同一新闻的点击数分散到 10 个临时 key 上,先做一轮局部累加,再去盐做第二轮汇总,热点压力就摊开了。代价是结果会有 1 个批次间隔的额外延迟,新闻热度场景完全能接受。
5.3 内存与GC参数:Spark2x Executor的JVM设置
实时任务里,内存配置错误通常不直接报错,而是表现为频繁 Full GC、批次处理时间越来越长。一个实用的提交命令如下:
spark-submit \ --class com.news.rt.NewsClickStreaming \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 20 \ --conf spark.executor.memoryOverhead=1024 \ --conf spark.streaming.backpressure.enabled=true \ --conf spark.streaming.kafka010.maxRatePerPartition=2000 \ --conf spark.executor.extraJavaOptions="-XX:+UseG1GC -XX:MaxGCPauseMillis=200" \ news-rt-1.0.jar参数说明:executor-memory 4g是堆内内存,memoryOverhead 1024是堆外,两者加起来就是单个 executor 向 YARN 申请的总内存。executor-cores 2意味着单 executor 两个并发任务,窗口聚合任务不太吃单核,没必要开 4 核,减少 GC 干扰更重要。开启 G1GC 并限制停顿目标,比默认的 ParallelGC 对实时任务的延迟更友好。
注意,内存配置要和状态量一起算:一个 60 秒窗口、每秒 200 条、每条约 1KB,窗口状态不过 12MB,4GB 堆绰绰有余;如果你把窗口拉到 1 小时,同时cache()了多个 RDD,4GB 就不够看了。遇到瓶颈先看 Spark UI 的 Storage 页签,别盲目加内存。
6. 进阶:用滑动窗口的“斜率”给可视化系统加实时告警
可视化做到能看只是及格,能从数据里自动发现异常才算可用。一个实用的进阶做法是:不告警绝对值,告警“变化速度”。新闻热度天然有早晚高峰,某条新闻在流量平峰期突然暴增,才是运营真正关心的信号。
思路很简单:维护最近 5 个窗口的热度值,用当前窗口值减去前一个窗口值得到增量,如果增量超过历史均值的 3 倍,且绝对增量大于某个阈值(比如 500),才触发告警:
// 假设 history 是最近 5 个窗口的点击数数组,current 是当前窗口值 double baseAvg = Arrays.stream(history).average().orElse(1.0); long delta = current - history[history.length - 1]; if (delta > 500 && delta > baseAvg * 3) { // 触发告警,推送到钉钉/企微机器人 alert("新闻热度异常飙升: " + newsId + ", delta=" + delta); }这个逻辑既可以放在 Spark 端,也可以放在后端查询 ES 时算。我建议放在后端,因为它不占实时计算资源,也不影响核心链路。告警阈值不要拍脑袋定,先让系统跑几天,把“正常波动范围”打印出来再设阈值,否则凌晨低峰期随便一条新闻都能把你叫醒。
我第一次做这类项目时,恨不得把每个组件都换成最新的,结果在版本兼容上浪费了两周。做实时可视化系统,稳比新重要,Spark2x 老一点没关系,只要 Kafka 客户端版本对得上、ES 批量写入不 429、检查点和偏移不打架,这个方向就能真正落地。数据链路通了之后,换 Flink、换 ClickHouse 都只是换零件,思路是通的。
希望帮到你。
本文还有配套的精品资源,点击获取