☰
Spark 2.X Structured Streaming 实时新闻话题统计实战
2026/10/3 2:46:40 网站建设 项目流程

简介:这份资源是面向计算机相关专业学生与大数据入门者的Spark 2.X实战项目资料,围绕新闻话题的实时统计分析场景展开,可用于毕业设计、课程设计、项目立项演示或自学进阶。压缩包共499个文件,约6.2MB,以400个xml配置、47个class编译文件、18个jar依赖包为主,另含9个scala源码、6个properties配置、4个js与2个html页面、2个java文件及说明文档,覆盖从依赖配置到核心逻辑的完整工程结构。项目已通过测试运行,功能正常,并获导师指导认可,答辩评审达95分。内容涉及Structured Streaming与Kafka对接、JDBCSink数据落地、Weblog服务与MySQL连接池等实时统计关键环节,读者可据此理解新闻话题的流式采集、处理与存储链路,并在此基础上修改扩展功能。目前已有62人学习,适合需要完整项目参考与排错思路的读者下载使用。

1. 新闻话题实时统计:Spark 2.x 项目里最值得先跑通的一条链路

新闻编辑部的需求往往来得很急:某个突发事件刚冒头,运营希望十分钟内看到全网相关报道的话题聚类、热度排序和来源分布。这个场景用离线跑批根本来不及,用 Spark 2.X 的 Structured Streaming 做微批处理却是刚好够用的方案。标题里的「新闻话题的实时统计分析」,本质是把持续流入的新闻文本流,按时间窗口切分,做分词、话题归并、词频统计和热度排序,最后落到可查询的存储里。它解决的是「数据还在流、结论就要出」的问题,适合已经会写 Spark SQL、但没搭过完整流式链路的工程师,也适合拿它当大数据毕业设计里那条能跑通、能演示、能讲清楚的主线。源码和资料再多,真正决定项目能不能立住的,是这条从数据源到结果表的链路有没有被完整跑通一次。

2. 拆解实时统计链路:从数据源到话题热度表

2.1 为什么选 Spark 2.X 而不是别的流处理框架

Spark 2.X 在 2016 到 2019 年间是国内大数据教学和项目实战的主力版本,Structured Streaming 在 2.0 引入、2.2 之后逐渐稳定,正好卡在「API 够新、资料够多、集群够好搭」的窗口上。选它做新闻话题实时统计,核心理由有三条。

第一,批流统一。新闻数据白天是流、晚上可能要回补历史,Structured Streaming 的 DataFrame/Dataset API 和 Spark SQL 共用一套抽象,同一段聚合逻辑既能跑流也能跑批,不用维护两套代码。第二,微批模型对新闻场景足够。新闻话题统计的时效要求通常是秒级到分钟级,不是毫秒级,微批的延迟完全能接受,而它带来的容错和 exactly-once 语义比纯流式方案更容易讲清楚。第三,生态成熟。Kafka 作为数据源、HDFS 或 MySQL 作为结果落地、Spark SQL 做聚合,这套组合在 2.X 时代有大量可参考的配置和踩坑记录。

对比 Flink,Flink 在事件时间和状态管理上更纯粹,但 Spark 2.X 对已经熟悉批处理的团队上手成本更低。如果你的团队本来就在用 Spark 跑离线报表,引入 Structured Streaming 几乎是零迁移成本。这也是为什么很多新闻类、舆情类项目在 2.X 时期选了 Spark 而不是另起炉灶。

2.2 数据流的四个环节与各自职责

一条完整的新闻话题实时统计链路,拆开看是四段:接入、清洗、聚合、落地。

接入层负责把新闻文本从消息队列读进来。常见做法是 Kafka 里一个 topic 存原始新闻 JSON,每条包含标题、正文、来源、发布时间戳。Spark 用readStream.format("kafka")消费,拿到的是 key/value 字节数组,需要先转成字符串再解析 JSON。

清洗层做三件事:去掉 HTML 标签和广告尾巴、按发布时间过滤掉过旧的数据、把正文切成词。中文分词在 Spark 里通常用 Ansj 或 HanLP 的 UDF 封装,把一篇文章变成一个词数组,再 explode 成多行,每行一个词。

聚合层是核心。按「时间窗口 + 话题关键词」分组,统计每个窗口内每个话题的出现次数、来源数、去重文章数。窗口一般用滑动窗口,比如窗口 10 分钟、滑动 1 分钟,这样热度曲线更新得比较平滑。

落地层把聚合结果写到外部存储。调试阶段写 console 最方便,生产上一般写 Kafka 供下游消费,或者写 MySQL 供前端查询。写 MySQL 需要用foreachBatch,因为 Structured Streaming 原生不直接支持 JDBC sink。

这四段里,清洗和聚合是最容易出问题的地方,后面会专门讲。

2.3 最小可运行代码:从 Kafka 读到话题词频

下面这段代码是整条链路的最小骨架,跑通它比读十篇原理文章都有用。假设 Kafka 里已经有 JSON 格式的新闻流,字段是title、content、source、ts。

from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col, explode, window, count, split from pyspark.sql.types import StructType, StringType, LongType # 1. 建 session,注意 2.X 里 master 和 shuffle 分区数要按集群调 spark = SparkSession.builder \ .appName("NewsTopicStreaming") \ .master("local[4]") \ .config("spark.sql.shuffle.partitions", "8") \ .getOrCreate() # 2. 定义新闻 JSON 的 schema,字段类型必须和上游一致 schema = StructType() \ .add("title", StringType()) \ .add("content", StringType()) \ .add("source", StringType()) \ .add("ts", LongType()) # 3. 从 Kafka 读流, earliest 保证首次启动从头消费 raw = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("subscribe", "news_raw") \ .option("startingOffsets", "earliest") \ .load() # 4. 字节转字符串再解析 JSON,坏数据用 PERMISSIVE 模式丢弃 parsed = raw.select( from_json(col("value").cast("string"), schema).alias("n") ).select("n.*").filter(col("content").isNotNull()) # 5. 简单分词:按空白和标点切,真实项目换成 Ansj UDF words = parsed.select( col("source"), col("ts"), explode(split(col("content"), "[\\s,,。!?]+")).alias("word") ).filter(col("word") != "") # 6. 滑动窗口聚合:窗口 10 分钟,每 1 分钟滑动一次 result = words.groupBy( window(col("ts").cast("timestamp"), "10 minutes", "1 minute"), col("word") ).agg(count("*").alias("freq")) # 7. 输出到 console,生产环境换成 foreachBatch 写 MySQL query = result.writeStream \ .outputMode("complete") \ .format("console") \ .option("truncate", "false") \ .start() query.awaitTermination()

逻辑说明:第 3 步的startingOffsets设成earliest只在首次调试时用,生产上要改成latest,否则每次重启都会重放历史数据。第 4 步的from_json遇到格式不对的行会返回 null,后面用filter丢掉,这是最省事的容错方式。第 5 步的分词是占位实现,中文场景必须换成真正的分词器,否则「北京」和「北京市」会被当成两个词。第 6 步的窗口参数是新闻场景的关键:窗口太短热度抖动大,窗口太长反应迟钝,10 分钟窗口配 1 分钟滑动是常见起点。

参数说明:spark.sql.shuffle.partitions默认 200,本地调试设成 8 能明显减少小数据量下的调度开销。outputMode用complete是因为聚合结果需要全量刷新,如果只想要增量用update,但要注意下游能不能处理更新。truncate设 false 是为了在 console 里看到完整词,不然长词会被截断。

2.4 结果落地:写 MySQL 的 foreachBatch 写法

console 只能看不能查,真正要演示或上线,得把结果写进 MySQL。Structured Streaming 没有原生 JDBC sink,标准做法是foreachBatch。

def write_to_mysql(batch_df, batch_id): # batch_df 是一个静态 DataFrame,可以当普通表操作 batch_df.select( col("window.start").alias("win_start"), col("window.end").alias("win_end"), col("word"), col("freq") ).write \ .format("jdbc") \ .option("url", "jdbc:mysql://localhost:3306/news") \ .option("dbtable", "topic_freq") \ .option("user", "root") \ .option("password", "yourpass") \ .mode("append") \ .save() query = result.writeStream \ .foreachBatch(write_to_mysql) \ .outputMode("update") \ .option("checkpointLocation", "/tmp/ckpt/news") \ .start()

这里有两个必须注意的点。一是checkpointLocation一定要设,否则重启后无法恢复状态,窗口聚合会从头算。二是mode用append还是overwrite取决于业务:如果下游按窗口查询,append 更安全;如果每次只要最新快照,可以在 batch 里先 delete 再 insert,但那样就不是幂等的了,需要自己保证。

3. 中文分词与话题归并:实时统计里最容易翻车的一环

3.1 分词 UDF 的写法与性能陷阱

中文新闻正文直接按空格切,结果基本没法看。真实项目里必须挂分词器。Spark 2.X 里常见做法是注册一个 UDF,内部调用 Ansj 或 HanLP。

from pyspark.sql.functions import udf from pyspark.sql.types import ArrayType, StringType from org.ansj.splitWord.analysis import ToAnalysis # 通过 py4j 调 JVM 侧 def seg(text): if not text: return [] # 过滤掉单字和停用词,减少后续聚合压力 return [w for w in ToAnalysis.parse(text).getTerms() if len(w.getName()) > 1 and w.getName() not in STOP_WORDS] seg_udf = udf(seg, ArrayType(StringType())) words = parsed.select( col("source"), col("ts"), explode(seg_udf(col("content"))).alias("word") )

这个 UDF 有两个性能陷阱。第一,Python UDF 在 Spark 2.X 里是逐行序列化到 Python 进程执行的,吞吐比 JVM 原生算子低一个数量级。如果数据量大,应该用 Scala 写 UDF 打成 jar,或者用pandas_udf(2.3 之后支持)。第二,分词器对象不要每次调用都 new,应该做成单例或线程安全的静态对象,否则 GC 压力会很大。

停用词表要自己维护,新闻里「记者」「报道」「通讯员」这类词出现频率极高但没有话题价值,不滤掉的话热度榜前几名永远是它们。

3.2 话题归并:从词频到「话题」还差一步

词频统计不等于话题统计。新闻里「台风」和「台风登陆」是两个词,但属于同一个话题。常见做法有两种。

一种是关键词映射表。维护一张「词 → 话题」的字典,比如「台风」「台风登陆」「强台风」都映射到「台风灾害」。分词后先查表,命中就替换成话题名,没命中就保留原词。这张表可以人工维护,也可以从历史数据里用 TF-IDF 或 TextRank 自动抽取候选词再人工确认。

另一种是共现聚类。在滑动窗口内统计词与词的共现次数,用简单的连通分量或 LDA 做粗聚类。这种做法更自动,但实时场景下计算量大,一般只在窗口结果上做二次处理,不放在主链路里。

我一般会先用映射表跑通,因为可控、可解释、出问题好排查。等链路稳定了再考虑加自动聚类。映射表本身可以放在广播变量里,避免每个 batch 都去查数据库。

# 广播话题映射表,避免每个 batch 重复拉取 topic_map = {"台风": "台风灾害", "台风登陆": "台风灾害", "强台风": "台风灾害"} broadcast_map = spark.sparkContext.broadcast(topic_map) def map_topic(word): return broadcast_map.value.get(word, word) map_udf = udf(map_topic, StringType()) topics = words.withColumn("topic", map_udf(col("word")))

广播变量在流式任务里只加载一次,后续 batch 复用,这是省 IO 的关键。注意映射表更新后需要重启流任务才能生效,如果要求热更新,得换成查 Redis 或广播 + 定时刷新。

3.3 窗口参数怎么定:三个必须一起调的旋钮

窗口大小、滑动步长、watermark 这三个参数是绑在一起的,单独调一个往往没效果。

窗口大小决定统计的平滑程度。新闻话题热度变化快,10 分钟窗口能覆盖一个事件的发酵期,再长就会把不同阶段混在一起。滑动步长决定结果更新频率,1 分钟是常见值,再短对新闻场景没意义,再长前端刷新会显得卡。watermark 决定迟到数据能等多久,新闻源的时间戳经常有延迟,设成 2 到 5 分钟比较稳妥。

result = words \ .withWatermark("event_time", "5 minutes") \ .groupBy( window(col("event_time"), "10 minutes", "1 minute"), col("topic") ).agg(count("*").alias("freq"))

watermark 设太小,迟到数据会被丢弃,统计偏低;设太大,状态会一直保留,内存涨得快。5 分钟是个经验值,具体要看数据源的时间戳分布。可以在清洗阶段先统计一下当前时间 - 事件时间的分布,再决定 watermark。

4. 避坑与排查:实时统计任务跑不稳的五个真实原因

4.1 任务跑几分钟就 OOM,日志里全是 shuffle spill

现象:Spark UI 上看到大量 spill to disk,Executor 内存持续上涨,最后 OOM。

原因:窗口聚合的状态没有及时清理,或者 shuffle 分区数太少导致单个分区数据量过大。Structured Streaming 的窗口状态默认会保留到 watermark 之后,如果 watermark 没设或设得过大,状态会无限增长。

解决:先确认 watermark 有没有设,设了的话看值是不是太大。然后把spark.sql.shuffle.partitions调大,本地调试 8 到 16,集群上按数据量给到 100 以上。还可以在groupBy之前先做一次repartition,把数据打散。

4.2 重启后结果翻倍,同一个窗口出现两次

现象:任务重启后,MySQL 里同一个窗口的数据出现重复行。

原因:checkpoint 没设或设在了临时目录,重启后无法恢复状态,窗口从头计算,而输出模式是 append,导致重复写入。

解决:checkpointLocation 必须设在一个持久化路径上,HDFS 或可靠的本地盘。如果下游不能容忍重复,写入时要按窗口主键做 upsert,或者用mode("overwrite")配合 batch 内先删后插。

4.3 中文分词结果全是单字,话题榜没法看

现象:统计结果里排名靠前的全是「的」「了」「在」这种单字。

原因:分词器没生效,或者停用词表没加载,或者分词后没过滤单字。

解决:先确认 UDF 真的被调用了,可以在 UDF 里打日志。然后检查停用词表路径对不对,Spark 任务里读本地文件要用--files分发或者放 HDFS。最后在分词后加filter(length(word) > 1),这一步能挡掉大部分噪音。

4.4 Kafka 消费延迟越来越高,lag 持续上涨

现象:Kafka 的 consumer lag 监控曲线一直往上走,任务处理速度跟不上生产速度。

原因:常见的是 Python UDF 太慢,或者每个 batch 的处理时间超过了触发间隔。也可能是maxOffsetsPerTrigger没设,一次拉太多数据导致 batch 处理时间过长。

解决:先看 Spark UI 里每个 batch 的 processing time,如果接近或超过 trigger 间隔,就要优化。把 Python UDF 换成 Scala UDF,或者设maxOffsetsPerTrigger限制每批数据量,让处理更平稳。还可以增加 Executor 数量或提高并行度。

4.5 时间戳解析失败,窗口全是 null

现象:结果里 window 字段是 null,或者所有数据都落在一个窗口里。

原因:ts字段是字符串或毫秒数,直接 cast 成 timestamp 可能失败或时区不对。Spark 2.X 默认时区是 UTC,如果数据是北京时间,窗口边界会偏 8 小时。

解决:在 SparkSession 里设spark.sql.session.timeZone=Asia/Shanghai。如果ts是毫秒数,用(col("ts")/1000).cast("timestamp")。解析失败的行会被 watermark 过滤掉,所以要先确认时间戳字段没有脏数据。

5. 把话题热度做成可查询的表:一个实用的验证技巧

链路跑通之后,怎么验证统计结果是对的?我的习惯是拿一批已知答案的历史数据回放。具体做法是:从 Kafka 里导出某一天的新闻数据,人工标注出当天真实的热门话题和大致热度排名,然后把这份数据按原始时间戳回放进流任务,对比输出结果和人工标注的偏差。偏差在可接受范围内,才说明窗口参数、分词、话题映射这一套是合理的。

这个验证过程里有个技巧值得单独说:用foreachBatch把每个 batch 的原始数据和聚合结果同时落一份到 HDFS,按 batch_id 分目录存。这样出问题时可以精确回放到某个 batch,看是数据问题还是逻辑问题。很多人只存聚合结果,排查时没有原始数据对照,只能靠猜。

def debug_batch(batch_df, batch_id): # 原始数据落一份,方便回放 batch_df.write.mode("overwrite") \ .parquet("/debug/batch_{}/raw".format(batch_id)) # 聚合结果落一份 agg = batch_df.groupBy("topic").count() agg.write.mode("overwrite") \ .parquet("/debug/batch_{}/agg".format(batch_id)) query = topics.writeStream \ .foreachBatch(debug_batch) \ .outputMode("update") \ .option("checkpointLocation", "/tmp/ckpt/debug") \ .start()

这个 debug 模式只在验证阶段开,生产上要关掉,否则 HDFS 会被小文件撑爆。验证通过后,把话题热度表接到前端,按窗口时间倒序查最近 N 个窗口的 top 话题,就是一个能演示的实时看板。

最后说个我自己的教训:早期做这类项目时,我总想一步到位把分词、聚类、情感分析全塞进主链路,结果每个环节都不稳,排查时根本不知道是哪一段出的问题。后来改成先跑通「Kafka → 分词 → 窗口词频 → MySQL」这条最短链路,确认端到端没问题,再逐个环节加功能,稳定性好了很多。实时统计这件事,链路短比功能全重要。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询