ARTICLE DETAIL

资讯详情

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

Spark 2.x实时新闻话题统计:Kafka到MySQL完整实现与避坑指南

Spark 2.x实时新闻话题统计:Kafka到MySQL完整实现与避坑指南 简介这份资源是面向计算机相关专业学生与大数据入门者的Spark 2.X新闻话题实时统计分析项目实战包可用于毕业设计、课程设计、作业或项目立项演示。项目已通过导师评审答辩成绩95分代码经测试可正常运行适合在现有基础上二次开发或直接复用。压缩包共499个文件约6.2MB以400个xml配置、47个class编译文件、18个jar依赖包为主另含9个scala源码、6个properties配置、4个js与2个html页面、2个java文件及说明文档覆盖从依赖配置到核心逻辑的完整工程结构。内容预览可见JDBCSink、StructuredStreamingKafka、StreamingKafka8/10、WeblogService、MySqlPool等模块涉及Kafka接入、结构化流处理、Web日志分析与MySQL连接池等实时统计关键环节。目前已有62人学习适合希望掌握Spark Streaming与Kafka整合、理解新闻话题实时统计流程的读者参考。1. 拆开这个 Spark 2.x 新闻话题实时统计包它到底能跑出什么结果如果你手头正压着一个大数据课程设计或者毕设题目是「新闻话题实时统计分析」要求用 Spark 做流式处理还要有可视化或者统计结果输出那这个包大概率能直接救急。它不是那种只有几页 PPT 和一堆截图的空壳项目而是把 Scala 源码、编译后的 class 文件、MySQL 连接池、Kafka 对接逻辑都塞进去了连.bak备份文件都在说明作者确实在本地反复跑过。核心链路是 Kafka 收新闻数据Spark Structured Streaming 或者 Spark Streaming 消费做话题词频统计最后通过 JDBCSink 落到 MySQL。适合谁适合已经装好 Spark 2.x、Kafka、MySQL但卡在「怎么把流式统计结果写进关系库」这一步的人。也适合想拿一个能跑通的骨架去改毕设的人因为它的类名和包结构很直白改起来不费劲。2. 环境对齐Spark 2.x 与 Kafka 版本匹配的硬约束2.1 为什么这个包锁死在 Spark 2.x 和 Scala 2.11打开压缩包你会看到StructuredStreamingKafka$.class和StreamingKafka8$.class、StreamingKafka10$.class同时存在。这不是作者乱放而是 Spark 2.x 时代对接 Kafka 的两条路Spark Streaming 走 Kafka 0.8 或 0.10 直连Structured Streaming 走 Kafka 0.10 的 source。StreamingKafka8对应的是老版KafkaUtils.createDirectStreamStreamingKafka10对应的是LocationStrategies那套新 API。如果你本地装的是 Spark 3.x这些 class 文件直接扔进去大概率报NoSuchMethodError因为 Spark 3 把 Kafka 0.8 的支持砍了Scala 也升到了 2.12。所以第一步不是急着跑而是把环境压回 Spark 2.4.x Scala 2.11 Kafka 0.10 或 0.11。常见做法是用 CDH 或者 Apache 官方二进制包别用最新版。2.2 从零把依赖版本钉死我一般会先写一个build.sbt把版本锁死避免 IDEA 自动拉最新版导致编译过不去。下面这段是我根据包内 class 文件反推出来的最小依赖集你可以直接抄// build.sbt name : NewsTopicStreaming version : 1.0 scalaVersion : 2.11.12 val sparkVersion 2.4.8 val kafkaVersion 0.10.2.2 libraryDependencies Seq( org.apache.spark %% spark-core % sparkVersion, org.apache.spark %% spark-sql % sparkVersion, org.apache.spark %% spark-streaming % sparkVersion, org.apache.spark %% spark-streaming-kafka-0-10 % sparkVersion, org.apache.spark %% spark-sql-kafka-0-10 % sparkVersion, org.apache.kafka %% kafka % kafkaVersion, mysql % mysql-connector-java % 5.1.47 )逻辑说明spark-streaming-kafka-0-10和spark-sql-kafka-0-10必须同时引入因为包里既有 DStream 写法也有 Structured Streaming 写法。mysql-connector-java用 5.1.x 而不是 8.x是因为MySqlPool.class里大概率用的是老版com.mysql.jdbc.Driver换成 8.x 会报驱动类找不到。参数上sparkVersion选 2.4.8 是因为它是 2.x 最后一个稳定版对 Kafka 0.10 兼容最好。如果你用 2.3.xspark-sql-kafka-0-10的 API 略有差异StructuredStreamingKafka里的option(kafka.bootstrap.servers)写法不变但startingOffsets的行为有区别建议直接上 2.4.8。2.3 把源码目录还原成 IDEA 能认的结构压缩包里 class 文件和 scala 源码混在一起直接导入 IDEA 会乱。我一般先按下面步骤理一遍# 假设解压到 news-spark 目录 cd news-spark mkdir -p src/main/scala/com/news/streaming mkdir -p src/main/resources # 把 .scala 和 .scala.bak 挪进 scala 目录 find . -name *.scala* -exec mv {} src/main/scala/com/news/streaming/ \; # class 文件单独放一个 lib 目录方便反编译对照 mkdir -p lib/classes find . -name *.class -exec mv {} lib/classes/ \;逻辑说明src/main/scala是 sbt 默认源码路径包名com.news.streaming是我根据WeblogService这种类名猜的你可以在源码第一行package声明里确认。.bak文件不要删它往往是作者改参数前的备份对比一下能看出哪些配置被调过。lib/classes里的 class 文件用 JD-GUI 或者javap -p反编译能快速看到JDBCSink里 MySQL 表名和字段名省得去猜。3. 核心链路拆解Kafka 进、Spark 算、MySQL 落3.1 StructuredStreamingKafka 的消费与解析逻辑StructuredStreamingKafka这个类名暗示它用的是 Spark 2.x 的 Structured Streaming。典型写法是从 Kafka 读出来是DataFramevalue 是二进制需要 cast 成 string 再解析。我根据常见新闻话题统计场景把核心代码补全成这样// StructuredStreamingKafka.scala val spark SparkSession.builder() .appName(NewsTopicStreaming) .master(local[2]) .getOrCreate() val df spark.readStream .format(kafka) .option(kafka.bootstrap.servers, localhost:9092) .option(subscribe, news-topic) .option(startingOffsets, latest) .load() import spark.implicits._ val lines df.selectExpr(CAST(value AS STRING)).as[String] // 假设新闻格式为 topic,content val topicDF lines.map(_.split(,)(0)).groupBy(value).count() val query topicDF.writeStream .outputMode(complete) .foreach(new JDBCSink()) .start() query.awaitTermination()逻辑说明startingOffsets设成latest表示只消费新数据做实时统计演示够用如果你要复现历史数据改成earliest。outputMode(complete)是因为groupBy().count()需要全量输出换成append会报错。foreach(new JDBCSink())是自定义 sink对应包里的JDBCSink.class。参数上master(local[2])本地跑至少给两个线程一个收 Kafka 一个做计算给一个线程会卡住。subscribe的 topic 名要和你在 Kafka 里创建的一致别照抄。3.2 JDBCSink 怎么写才能不丢数据JDBCSink是整条链路最容易翻车的地方。Structured Streaming 的foreach要求实现ForeachWriter里面open、process、close三个方法必须成对。我见过太多人把 MySQL 连接写在process里结果每条数据开一次连接跑几分钟就连接数爆了。正确做法是用MySqlPool做连接池在open里取连接close里归还。下面是我改过的版本// JDBCSink.scala class JDBCSink extends ForeachWriter[Row] { var conn: Connection _ var stmt: PreparedStatement _ override def open(partitionId: Long, version: Long): Boolean { conn MySqlPool.getConnection() // 从池里拿 stmt conn.prepareStatement( INSERT INTO topic_count(topic, cnt, ts) VALUES(?,?,?) ON DUPLICATE KEY UPDATE cnt?, ts? ) true } override def process(row: Row): Unit { val topic row.getString(0) val cnt row.getLong(1) stmt.setString(1, topic) stmt.setLong(2, cnt) stmt.setLong(3, System.currentTimeMillis()) stmt.setLong(4, cnt) stmt.setLong(5, System.currentTimeMillis()) stmt.executeUpdate() } override def close(errorOrNull: Throwable): Unit { if (stmt ! null) stmt.close() if (conn ! null) MySqlPool.returnConnection(conn) } }逻辑说明ON DUPLICATE KEY UPDATE要求 MySQL 表里topic字段有唯一索引否则每次都是 insert表会爆。MySqlPool是包内已有的类你需要在open之前确保池已经初始化通常放在main方法最前面。参数上partitionId和version在open里可以用来做幂等判断但这里简化处理。注意close里一定要判空因为open失败时close也会被调用不判空会抛 NPE 把整个流干掉。3.3 MySQL 建表与连接池参数建表语句不能少我一般直接跑这段CREATE TABLE topic_count ( topic VARCHAR(64) NOT NULL, cnt BIGINT DEFAULT 0, ts BIGINT DEFAULT 0, PRIMARY KEY (topic) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;逻辑说明topic做主键是为了配合上面的ON DUPLICATE KEY UPDATE。ts存毫秒时间戳方便后面做趋势图。utf8mb4是为了支持中文话题名用utf8在某些 MySQL 版本下会截断。连接池参数在MySqlPool里通常有maxActive、maxIdle、maxWait我一般把maxActive设成 10 到 20maxWait设 3000 毫秒别设太大否则流处理线程会一直等连接吞吐直接掉下来。4. 避坑与排查跑不起来先看这五条4.1 报 NoSuchMethodError 或 ClassNotFoundException现象提交任务后立刻抛java.lang.NoSuchMethodError: org.apache.spark.sql.streaming.DataStreamReader.option或者ClassNotFoundException: org.apache.kafka.common.serialization.StringDeserializer。原因Spark 版本和 Kafka 依赖版本对不上或者spark-sql-kafka-0-10没打进 fat jar。解决用sbt assembly打胖包确认build.sbt里spark-sql-kafka-0-10的版本和spark-core完全一致别一个 2.4.8 一个 2.4.0。4.2 MySQL 连接数暴涨然后任务挂掉现象跑几分钟后报Too many connections或者Communications link failure。原因JDBCSink里每条记录都DriverManager.getConnection没走连接池。解决确认MySqlPool在open里被调用并且close里归还了连接。如果MySqlPool本身没实现单例自己加一个object MySqlPool保证全局只有一个池。4.3 Kafka 数据消费到了但统计结果一直是空现象Spark UI 里能看到inputRowsPerSecond有值但 MySQL 表里没数据。原因outputMode设成了append而groupBy后的聚合结果在append模式下只有 watermark 触发后才输出没设 watermark 就永远不输出。解决改成complete模式或者加上withWatermark(timestamp, 10 minutes)并确保数据里有时间列。4.4 中文话题名入库变成问号现象MySQL 里查出来topic字段是???。原因JDBC URL 没加useUnicodetruecharacterEncodingutf8或者表字符集是latin1。解决连接串改成jdbc:mysql://localhost:3306/news?useUnicodetruecharacterEncodingutf8mb4表也确认是utf8mb4。4.5 本地跑正常提交到集群就报序列化错误现象org.apache.spark.SparkException: Task not serializable。原因JDBCSink里引用了外部不可序列化的对象比如直接持有SparkSession或者Connection作为成员变量且在process里用。解决把连接相关的东西都放在open和close里process只做纯数据操作。MySqlPool如果是 object 单例在 executor 端会重新初始化注意配置要能读到。5. 进阶技巧用反编译对照源码快速改毕设5.1 用 javap 看 class 文件里的真实字段包里的.class文件不是摆设。当你发现.scala源码缺了某个方法或者.bak和当前版本不一致时直接反编译最快。我常用这条命令javap -p -c lib/classes/JDBCSink.class jdbcsink.txt逻辑说明-p显示私有成员-c反汇编字节码。打开jdbcsink.txt搜INSERT INTO就能看到作者实际用的表名和字段搜getConnection就能看到连接池调用方式。这比猜快得多尤其适合改毕设时快速定位要改哪几个参数。5.2 把统计维度从话题扩展到时间窗口原包大概率只按话题分组计数。如果你毕设要求「每 5 分钟统计一次热门话题」需要把groupBy(value)改成带窗口的聚合import org.apache.spark.sql.functions._ val windowed lines .withColumn(ts, current_timestamp()) .groupBy(window($ts, 5 minutes), $value) .count()逻辑说明window($ts, 5 minutes)会生成一个struct类型的窗口列落库时需要拆成window.start和window.end。current_timestamp()是处理时间不是事件时间做演示够用如果新闻数据里自带时间戳换成to_timestamp($event_time)更准。参数上窗口大小和滑动间隔可以按需改5 minutes换成10 minutes就是十分钟窗口。5.3 验证统计结果是否正确的笨办法别只看 MySQL 里有没有数据要验证数字对不对。我一般开两个终端一个用kafka-console-producer手动发几条已知话题的新闻比如发三条sports,xxx、两条tech,yyy另一个终端盯 MySQL 表。如果sports的cnt是 3、tech是 2说明链路正确。如果数字对不上先查 Kafka 里是不是有重复消费再看groupBy之前有没有做distinct。这个笨办法能排掉八成逻辑错误。从那以后我每次拿到这种带 class 文件的源码包都先javap一遍再动手改省得在源码和编译产物不一致的地方来回翻车。希望帮到你。本文还有配套的精品资源点击获取
返回列表