项目6:实时计算书籍热度——Spark Streaming 实时计算框架
先修基础:项目1-5(Spark概述、Scala基础、Spark Shell编程、Spark IDE编程、Spark SQL)
目录
- 第一部分:项目背景与 Spark Streaming 概述
- 第二部分:Spark Streaming 框架与运行原理
- 第三部分:初步使用 Spark Streaming
- 第四部分:DStream 编程模型
- 第五部分:DStream 转换操作
- 第六部分:DStream 窗口操作(重点+难点)
- 第七部分:DStream 输出操作
- 第八部分:项目实战——实时计算书籍热度
- 第九部分:常见问题与排错指南
- 第十部分:实习 / 面试高频考点
- 附录:习题解析
第一部分:项目背景与 Spark Streaming 概述
1.1 项目背景
为什么需要实时计算?
前面学的 Spark Core、Spark SQL 都是批处理——处理的是已经存好的历史数据。 但现实中很多场景需要实时处理:
- 电商网站:实时统计商品热度、用户行为
- 网站监控:实时监控访问量、异常流量
- 金融风控:实时检测欺诈交易
- 物联网:实时处理传感器数据
Spark Streaming 就是 Spark 的实时计算框架。
项目场景
某电商平台想要实时计算书籍热度,把热度高的书推荐给用户。
书籍热度计算公式:
热度 = 用户平均评分 × 用户评分次数 × 0.3 + 书籍平均评分 × 书籍被评分次数
数据:BookRating.txt(用户对书籍的评分数据,评分范围1~5分)
本项目要做的事:
- 模拟实时产生的评分数据
- 用 Spark Streaming 实时读取和处理
- 实时计算书籍热度,取 Top10 存入 Hive
1.2 批处理 vs 流处理
| 对比项 | 批处理(Batch) | 流处理(Streaming) |
|---|---|---|
| 数据特点 | 有界、历史数据 | 无界、实时到达 |
| 处理方式 | 一次处理全部数据 | 来一条处理一条(或小批量) |
| 延迟 | 高(分钟/小时级) | 低(秒级/毫秒级) |
| 代表技术 | MapReduce、Spark Core、Spark SQL | Spark Streaming、Flink、Storm |
| 适用场景 | 离线分析、数据仓库 | 实时监控、实时推荐、风控 |
💡 理解:
- 批处理就像"每天下班算一次账"
- 流处理就像"每收一笔钱马上记账"
第二部分:Spark Streaming 框架与运行原理
2.1 Spark Streaming 是什么?
Spark Streaming 是 Spark 的子框架,用于处理流式数据。
- 可伸缩、高吞吐量、容错能力强
- 和 Spark SQL、MLlib、GraphX 无缝集成
- 支持多种数据源:Kafka、Flume、HDFS、Socket 等
- 结果可以保存到 HDFS、MySQL、Dashboard 等
⚠️ 注意:Spark Streaming 从 Spark 3.4.0 起被官方标记为"不推荐使用"(Deprecated), 3.5.1 还支持,但后续版本可能会完全弃用。新的实时框架是 Structured Streaming。 但很多老项目还在用 Spark Streaming,所以还是要学。
2.2 运行原理(微批处理)
Spark Streaming 不是真正的"来一条处理一条",而是微批处理(Micro-batching)。
核心思想
- 把实时数据流按时间片切成一小段一小段
- 每一段数据当成一个小批处理来算
- 批量输出结果
关键概念
| 概念 | 说明 |
|---|---|
| 批处理时间间隔(Batch Interval) | 切分数据的时间间隔,比如10秒切一次 |
| DStream(离散流) | Spark Streaming 的数据抽象,代表持续的数据流 |
| 每个时间片 → 一个RDD | DStream 就是一系列 RDD 的序列 |
💡 形象比喻:
- 数据流就像一条河,源源不断
- 批处理时间间隔就像"每隔10秒舀一桶水"
- 每一桶水就是一个 RDD
- DStream 就是一桶接一桶的水(RDD 序列)
2.3 DStream 和 RDD 的关系
DStream = 时间1的RDD + 时间2的RDD + 时间3的RDD + ...- 对 DStream 的操作 → 对每个 RDD 的操作
- DStream 的转换操作 → 底层还是 RDD 的转换操作
- DStream 也有惰性求值的特点,需要输出操作才会真正执行
💡 理解:DStream 就像"RDD 的播放列表",每个时间点播放一个 RDD。
第三部分:初步使用 Spark Streaming
3.1 使用 Spark Streaming 的步骤
- 创建 StreamingContext:流处理的入口
- 创建输入源 DStream:指定从哪里读数据(Socket、HDFS、Kafka等)
- 操作 DStream:转换、窗口、输出等
- 启动流处理:调用
ssc.start() - 等待停止:调用
ssc.awaitTermination()
⚠️ 注意:start() 之前只是设定执行计划,并没有真正运行;start() 之后才开始真正处理数据。
3.2 示例1:Socket 数据源(单词实时计数)
从网络端口读取数据,实时统计单词出现次数。
准备工作:安装 nc 工具
nc(netcat)是一个网络工具,可以用来发送数据。
# 检查是否安装
dnf list --installed | grep nc
# 如果没装,安装一下
dnf install -y nc步骤1:启动 nc 发送数据(slave1 节点)
# 在slave1上启动,监听8888端口
nc -lk 8888步骤2:编写流处理程序(spark-shell)
import org.apache.spark.streaming.{Seconds, StreamingContext}
// 1. 创建StreamingContext,批处理间隔10秒
val ssc = new StreamingContext(sc, Seconds(10))
// 2. 创建Socket输入DStream
val lines = ssc.socketTextStream("slave1", 8888)
// 3. 单词计数
val words = lines.flatMap(_.split(" "))
val wordCounts = words.map(x => (x, 1)).reduceByKey(_ + _)
// 4. 输出结果
wordCounts.print()
// 5. 启动流处理
ssc.start()
ssc.awaitTermination()步骤3:在 nc 端口输入数据
在 slave1 的 nc 窗口输入:
I am learning Spark Streaming now然后在 spark-shell 中就能看到统计结果。
停止程序
按 Ctrl + C 停止。
💡 说明:
Seconds(10)表示每10秒处理一批print()默认打印前10个元素awaitTermination()等待程序终止(一直运行,直到手动停止)
3.3 示例2:HDFS 目录数据源
监控 HDFS 目录,有新文件进来就处理。
代码
import org.apache.spark.streaming.{Seconds, StreamingContext}
// 创建StreamingContext
val ssc = new StreamingContext(sc, Seconds(10))
// 监控HDFS目录
val lines = ssc.textFileStream("/tipdm/data/sparkStreaming/temp")
// 单词计数
val words = lines.flatMap(_.split(" "))
val wordCounts = words.map(x => (x, 1)).reduceByKey(_ + _)
// 输出
wordCounts.print()
// 启动
ssc.start()
ssc.awaitTermination()测试
程序运行后,另开一个终端,往 HDFS 目录上传文件:
# 创建目录
hdfs dfs -mkdir -p /tipdm/data/sparkStreaming/temp
# 上传第一个文件
hdfs dfs -put a.txt /tipdm/data/sparkStreaming/temp/
# 等10秒以上,再上传第二个文件
hdfs dfs -put b.txt /tipdm/data/sparkStreaming/temp/每次上传新文件,Spark Streaming 都会在10秒内处理并输出结果。
⚠️ 注意:
- textFileStream 只监控新文件,已经存在的文件不会处理
- 文件必须是"移动"或"上传"到目录中的,修改已有文件不会触发
- 上传间隔要大于批处理间隔,否则可能漏处理
3.4 常用数据源
| 数据源 | 方法 | 说明 |
|---|---|---|
| Socket | socketTextStream(host, port) | 从网络端口读取,学习测试用 |
| HDFS文件 | textFileStream(directory) | 监控HDFS目录,新文件触发 |
| Kafka | KafkaUtils.createStream(...) | 生产环境最常用 |
| Flume | FlumeUtils.createStream(...) | 日志采集场景 |
💡 说明:生产环境最常用的是 Kafka,本课程主要学 Socket 和 HDFS 两种基础的。
第四部分:DStream 编程模型
4.1 什么是 DStream?
DStream(Discretized Stream,离散流)是 Spark Streaming 的核心抽象。
- 代表持续的实时数据流
- 底层是一系列 RDD,每个时间片一个 RDD
- 对 DStream 的操作 = 对每个 RDD 的操作
4.2 DStream 操作分类
和 RDD 类似,DStream 的操作也分三类:
| 操作类型 | 说明 | 类似RDD的 |
|---|---|---|
| 转换操作(Transformation) | 从一个 DStream 生成新的 DStream | 转换算子 |
| 窗口操作(Window) | 基于滑动窗口的计算 | 没有对应,特有 |
| 输出操作(Output) | 触发计算,输出到外部 | 行动算子 |
💡 注意:DStream 的转换操作也是惰性的,需要输出操作才会真正执行。
第五部分:DStream 转换操作
5.1 常用转换操作
| 方法 | 说明 |
|---|---|
map(func) | 每个元素转成一个新元素 |
flatMap(func) | 每个元素转成0或多个元素 |
filter(func) | 过滤符合条件的元素 |
reduceByKey(func) | 按Key聚合(每个批次内) |
union(otherStream) | 合并两个DStream |
count() | 统计每个批次的元素个数 |
reduce(func) | 每个批次内聚合 |
join(otherStream) | 两个DStream按Key连接 |
transform(func) | 对每个RDD做任意操作 |
💡 注意:这些转换操作都是针对每个批次独立的,批次之间不累积。 如果要跨批次累积,需要用窗口操作或者有状态的操作。
5.2 transform() 方法(重要)
transform() 是一个很灵活的方法,可以直接操作底层的 RDD。
用途
- DStream API 不够用时,直接用 RDD 的方法
- 可以调用 RDD 的任何操作
示例
// 用transform实现flatMap的效果
val words = lines.transform(rdd => rdd.flatMap(_.split(" ")))
words.print()💡 理解:transform 就像"后门",DStream 做不了的事,直接转成 RDD 做。
5.3 转换操作 vs RDD 操作对比
| DStream 操作 | 对应 RDD 操作 | 区别 |
|---|---|---|
map(func) | map(func) | DStream 是每个批次的 RDD 都 map |
flatMap(func) | flatMap(func) | 同上 |
filter(func) | filter(func) | 同上 |
reduceByKey(func) | reduceByKey(func) | 只在当前批次内聚合 |
transform(func) | - | 直接操作 RDD,最灵活 |
第六部分:DStream 窗口操作(重点+难点)
6.1 什么是窗口操作?
普通的转换操作只处理当前批次的数据。 窗口操作可以处理一段时间内的数据(多个批次)。
两个关键参数
| 参数 | 说明 |
|---|---|
| 窗口长度(Window Length) | 窗口包含多长时间的数据 |
| 滑动步长(Slide Interval) | 窗口每隔多久滑动一次 |
💡 形象比喻:
- 窗口就像"滑动的窗户"
- 窗口长度 = 窗户有多宽(能看到多少数据)
- 滑动步长 = 窗户多久往前挪一次
- 每次挪一下,就计算一次窗户里的数据
重要规则
窗口长度和滑动步长都必须是批处理间隔的整数倍!
⚠️ 面试常考:窗口长度和滑动步长必须是批处理间隔的整数倍,否则会报错。
6.2 常用窗口操作
| 方法 | 说明 |
|---|---|
window(windowLength, slideInterval) | 生成窗口化的DStream |
countByWindow(windowLength, slideInterval) | 统计窗口内元素总数 |
reduceByWindow(func, windowLength, slideInterval) | 窗口内整体聚合 |
reduceByKeyAndWindow(func, windowLength, slideInterval) | 窗口内按Key聚合 |
countByValueAndWindow(windowLength, slideInterval) | 窗口内按值计数 |
6.3 window() 方法
最基础的窗口操作,把数据按窗口重新分组。
示例
import org.apache.spark.streaming.{Seconds, StreamingContext}
val ssc = new StreamingContext(sc, Seconds(1)) // 批处理间隔1秒
val lines = ssc.socketTextStream("slave1", 8888)
val words = lines.flatMap(_.split(" "))
// 窗口长度3秒,滑动步长2秒
val windowWords = words.window(Seconds(3), Seconds(2))
windowWords.print()
ssc.start()
ssc.awaitTermination()效果说明
- 批处理间隔:1秒(每秒一个批次)
- 窗口长度:3秒(每个窗口包含3个批次的数据)
- 滑动步长:2秒(每2秒计算一次)
计算时机:第2秒、第4秒、第6秒... 每次计算的数据:前3秒的数据
6.4 reduceByKeyAndWindow() 方法
最常用的窗口操作,窗口内按 Key 聚合。
示例:窗口内单词计数
import org.apache.spark.streaming.{Seconds, StreamingContext}
val ssc = new StreamingContext(sc, Seconds(1))
val lines = ssc.socketTextStream("slave1", 8888)
val words = lines.flatMap(_.split(" "))
val pairs = words.map(word => (word, 1))
// 窗口长度3秒,滑动步长2秒,按Key求和
val windowWordCounts = pairs.reduceByKeyAndWindow(
(a: Int, b: Int) => a + b, // 聚合函数
Seconds(3), // 窗口长度
Seconds(2) // 滑动步长
)
windowWordCounts.print()
ssc.start()
ssc.awaitTermination()💡 理解:reduceByKeyAndWindow 就是"窗口版的 reduceByKey", reduceByKey 只算当前批次,reduceByKeyAndWindow 算窗口内所有批次。
6.5 窗口操作示意图
假设批处理间隔=1秒,窗口长度=3秒,滑动步长=2秒:
时间: 0s 1s 2s 3s 4s 5s 6s
批次: [0] [1] [2] [3] [4] [5] [6]
第2秒计算:窗口包含 [0][1][2] → 3个批次
第4秒计算:窗口包含 [2][3][4] → 3个批次
第6秒计算:窗口包含 [4][5][6] → 3个批次6.6 reduceByKey vs reduceByKeyAndWindow 对比
| 对比项 | reduceByKey | reduceByKeyAndWindow |
|---|---|---|
| 数据范围 | 当前批次 | 窗口内多个批次 |
| 时间维度 | 无时间概念 | 有时间窗口 |
| 参数 | 只有聚合函数 | 聚合函数 + 窗口长度 + 滑动步长 |
| 适用场景 | 每批次统计 | 近N分钟/小时统计 |
第七部分:DStream 输出操作
7.1 常用输出操作
| 方法 | 说明 |
|---|---|
print() | 打印前10个元素(Driver端) |
saveAsTextFiles(prefix, suffix) | 保存为文本文件 |
saveAsObjectFiles(prefix, suffix) | 保存为序列化对象文件 |
saveAsHadoopFiles(prefix, suffix) | 保存为Hadoop文件 |
foreachRDD(func) | 对每个RDD做任意输出操作(最灵活) |
💡 注意:输出操作是触发计算的,没有输出操作的话,转换操作都不会执行。
7.2 print() 方法
最简单的输出,在控制台打印前10个元素。
wordCounts.print()7.3 saveAsTextFiles() 方法
保存为文本文件,每个批次生成一个文件夹。
lines.saveAsTextFiles(
"hdfs://master:8020/tipdm/data/saveAsTextFiles/sahf", // 前缀
"txt" // 后缀
)生成的文件路径类似:sahf-时间戳.txt
💡 说明:每个批次生成一个文件夹,文件夹里是该批次的数据。
7.4 foreachRDD() 方法(重点+难点)
最灵活的输出操作,可以对每个 RDD 做任意操作。 常用于把结果写入 MySQL、Redis 等外部系统。
正确使用方式很重要!
❌ 错误写法1:在 Driver 端创建连接
dstream.foreachRDD { rdd =>
val connection = createNewConnection() // 在Driver端创建,错误!
rdd.foreach { record =>
connection.send(record)
}
}问题:连接对象在 Driver 创建,要序列化发送到 Worker,连接对象很难序列化。
❌ 错误写法2:每条记录创建一个连接
dstream.foreachRDD { rdd =>
rdd.foreach { record =>
val connection = createNewConnection() // 每条记录创建一个,错误!
connection.send(record)
connection.close()
}
}问题:每条记录都创建连接,开销太大,严重影响性能。
✅ 正确写法:每个分区创建一个连接
dstream.foreachRDD { rdd =>
rdd.foreachPartition { partitionOfRecords =>
val connection = createNewConnection() // 每个分区创建一个连接
partitionOfRecords.foreach(record => connection.send(record))
connection.close()
}
}原理:一个分区用一个连接,分区内多条记录复用连接。
⭐ 优化写法:用连接池
dstream.foreachRDD { rdd =>
rdd.foreachPartition { partitionOfRecords =>
val connection = ConnectionPool.getConnection() // 从连接池取
partitionOfRecords.foreach(record => connection.send(record))
ConnectionPool.returnConnection(connection) // 归还连接池
}
}原理:用连接池复用连接,减少创建销毁的开销。
💡 面试常考:foreachRDD 的正确使用方式,三种写法的区别和问题。
7.5 示例:写入 MySQL 数据库
步骤1:在 MySQL 中创建表
create database spark;
use spark;
create table searchKeyWord(
insert_time datetime,
keyword varchar(30),
search_count integer
);步骤2:编写 Spark Streaming 程序
import org.apache.spark.sql.SparkSession
import org.apache.spark.streaming.{Seconds, StreamingContext}
import java.sql.DriverManager
object WriteDataToMySQL {
def main(args: Array[String]): Unit = {
val spark = SparkSession.builder()
.appName("WriteDataToMySQL")
.master("local[*]")
.getOrCreate()
val ssc = new StreamingContext(spark.sparkContext, Seconds(30))
ssc.sparkContext.setLogLevel("ERROR")
// 从Socket读取数据
val itemsStream = ssc.socketTextStream("slave1", 8888)
val itemPairs = itemsStream.map(x => (x.split(",")(0), 1))
// 窗口计算:窗口60秒,滑动60秒
val itemCounts = itemPairs.reduceByKeyAndWindow(
(x: Int, y: Int) => x + y,
Seconds(60),
Seconds(60)
)
// 取热度Top3
val hottestWord = itemCounts.transform(rdd => {
val top3 = rdd.map(x => (x._2, x._1))
.sortByKey(false)
.map(x => (x._2, x._1))
.take(3)
ssc.sparkContext.makeRDD(top3)
})
hottestWord.print()
// 写入MySQL
hottestWord.foreachRDD(rdd => {
rdd.foreachPartition(partitionOfRecords => {
val url = "jdbc:mysql://master:3306/spark"
val user = "root"
val password = "123456"
Class.forName("com.mysql.cj.jdbc.Driver")
val conn = DriverManager.getConnection(url, user, password)
// 清理30秒前的数据
conn.prepareStatement(
"delete from searchKeyWord where insert_time < now() - interval 30 second"
).executeUpdate()
conn.setAutoCommit(false)
val stmt = conn.createStatement()
// 批量插入
partitionOfRecords.foreach(record => {
stmt.addBatch(
s"insert into searchKeyWord(insert_time, keyword, search_count) " +
s"values(now(), '${record._1}', '${record._2}')"
)
})
stmt.executeBatch()
conn.commit()
conn.close()
})
})
ssc.start()
ssc.awaitTermination()
ssc.stop()
}
}代码说明
- 用
foreachPartition每个分区创建一个 MySQL 连接 - 用
addBatch批量插入,提高性能 - 用
transform配合 RDD 的sortByKey实现排序(DStream 没有 sort 方法)
💡 技巧:DStream 没有 sort 方法,需要排序时用 transform() 转成 RDD 来排。
第八部分:项目实战——实时计算书籍热度
8.0 数据说明
数据文件
BookRating.txt:用户对书籍的评分数据
字段说明
| 列索引 | 字段 | 说明 |
|---|---|---|
| 0 | UserID | 用户ID |
| 1 | BookID | 书籍ID |
| 2 | Rating | 评分(1~5分) |
书籍热度公式
热度 = 用户平均评分 × 用户评分次数 × 0.3 + 书籍平均评分 × 书籍被评分次数公式中的四个变量:
- u:用户平均评分(user_avg_rating)
- x:用户评分次数(user_count)
- y:书籍平均评分(book_avg_rating)
- z:书籍被评分次数(book_count)
任务6.1:获取输入数据源
需求
- 写一个日志生成模拟器,模拟实时产生评分数据
- 用 Spark Streaming 实时读取新文件
步骤1:日志生成模拟器(CreateData.scala)
每隔60秒,随机从 BookRating.txt 中抽取100条记录,写入新文件。
import java.io.{File, PrintWriter}
import java.text.SimpleDateFormat
import java.util.Date
import scala.io.Source
import scala.util.Random
object CreateData {
def main(args: Array[String]): Unit = {
var i = 0
while (true) {
// 读取源文件
val filename = "E:\\data\\BookRating.txt"
val lines = Source.fromFile(filename).getLines().toList
val firerow = lines.length
// 创建新文件
val writer = new PrintWriter(
new File(s"E:\\data\\StreamingData\\StreamingData-$i.txt")
)
i = i + 1
var j = 0
while (j < 100) {
// 随机选一行
val time = new Date().getTime
val format = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
println("=" * 10 + format.format(time) + "当前时间点" + "=" * 10)
writer.write(lines(index(firerow)) + "\n")
println(lines(index(firerow)))
j = j + 1
}
writer.close()
Thread.sleep(60000) // 等60秒
}
}
// 生成随机索引
def index(length: Int) = {
val rdm = new Random
rdm.nextInt(length)
}
}步骤2:实时读取数据(rating.scala)
import org.apache.spark.sql.SparkSession
import org.apache.spark.streaming.{Seconds, StreamingContext}
object rating {
def main(args: Array[String]): Unit = {
// 1. 创建SparkSession
val spark = SparkSession.builder()
.appName("rating")
.master("local[*]")
.enableHiveSupport()
.getOrCreate()
// 2. 创建StreamingContext,批处理间隔60秒
val ssc = new StreamingContext(spark.sparkContext, Seconds(60))
ssc.sparkContext.setLogLevel("ERROR")
// 3. 监控目录,读取新文件
val stream = ssc.textFileStream("E:\\data\\StreamingData")
// 4. 按制表符切分,转成三元组
val splitData = stream.map { x =>
val y = x.split("\t")
(y(0), y(1), y(2).toInt)
}
// 5. 打印测试
splitData.print()
// 6. 启动
ssc.start()
ssc.awaitTermination()
ssc.stop()
}
}运行方式
- 先运行 CreateData 生成数据
- 再运行 rating 读取数据
- 每隔60秒处理一批新文件
💡 说明:两个程序要同时运行,一个产生数据,一个消费数据。
任务6.2:计算用户评分次数及平均评分
需求
统计每个用户的评分次数和平均评分。
实现代码
import spark.implicits._
import java.text.SimpleDateFormat
import java.util.Date
splitData.foreachRDD { line =>
// 1. RDD转DataFrame
val dataFrame = line.toDF("UserID", "BookID", "Ratings")
dataFrame.show()
// 2. 计算用户平均评分
val user_rating = dataFrame
.groupBy("UserID")
.avg("Ratings")
.withColumnRenamed("avg(Ratings)", "user_avg_rating")
// 3. 计算用户评分次数
val user_count = dataFrame
.groupBy("UserID")
.count()
.withColumnRenamed("count", "user_count")
// 4. 连接:原始数据 + 用户平均评分
val data_user_rating = dataFrame
.join(user_rating, user_rating("UserID") === dataFrame("UserID"))
.drop(user_rating("UserID"))
// 5. 再连接:+ 用户评分次数
val data_user = data_user_rating
.join(user_count, user_count("UserID") === data_user_rating("UserID"))
.drop(user_count("UserID"))
// 打印结果
val time = new Date().getTime
val format = new SimpleDateFormat("HH:mm:ss")
println("=" * 20 + format.format(time) + "=" * 20)
data_user.show()
}代码说明
| 代码 | 说明 |
|---|---|
line.toDF(...) | RDD 转 DataFrame,指定列名 |
groupBy("UserID").avg("Ratings") | 按用户分组,求平均评分 |
withColumnRenamed(...) | 重命名列,默认列名是 avg(Ratings) |
join(...) | 连接两个 DataFrame |
.drop(...) | 删除重复的 UserID 列 |
💡 为什么要 join 回去? 因为后面计算热度需要每条记录都带上用户的统计信息,所以要把统计结果 join 回原始数据。
任务6.3:计算书籍被评分次数及平均评分
需求
统计每本书的被评分次数和平均评分,并和前面的数据合并。
实现代码
// 1. 计算书籍平均评分
val book_rating = dataFrame
.groupBy("BookID")
.avg("Ratings")
.withColumnRenamed("avg(Ratings)", "book_avg_rating")
// 2. 计算书籍被评分次数
val book_count = dataFrame
.groupBy("BookID")
.count()
.withColumnRenamed("count", "book_count")
// 3. 连接:用户数据 + 书籍平均评分
val data_user_book_rating = data_user
.join(book_rating, book_rating("BookID") === data_user("BookID"))
.drop(book_rating("BookID"))
// 4. 再连接:+ 书籍评分次数
val total_data = data_user_book_rating
.join(book_count, book_count("BookID") === data_user_book_rating("BookID"))
.drop(book_count("BookID"))
// 打印
val time = new Date().getTime
val format = new SimpleDateFormat("HH:mm:ss")
println("=" * 20 + format.format(time) + "=" * 20)
total_data.show(5)最终表结构
连接后,每条记录都包含:
- UserID、BookID、Ratings(原始数据)
- user_avg_rating、user_count(用户统计)
- book_avg_rating、book_count(书籍统计)
💡 理解:就像给每条记录"贴上"用户的统计信息和书籍的统计信息, 这样每条记录都能直接计算热度了。
任务6.4:实时计算书籍热度
需求
根据公式计算每本书的热度,取 Top10 存入 Hive。
热度公式
热度 = 用户平均评分 × 用户评分次数 × 0.3 + 书籍平均评分 × 书籍被评分次数实现代码
import org.apache.spark.sql.functions.{col, desc}
// 1. 计算热度,新增hot列
val BookHot = total_data.withColumn(
"hot",
col("user_avg_rating") * col("user_count") * 0.3 +
col("book_avg_rating") * col("book_count")
)
// 2. 按热度降序,取前10,保存到Hive
BookHot.sort(desc("hot"))
.limit(10)
.write.mode("overwrite")
.saveAsTable("book.topBookHot")
// 3. 打印Top10
val time = new Date().getTime
val format = new SimpleDateFormat("HH:mm:ss")
println("=" * 20 + format.format(time) + "=" * 20)
BookHot.sort(desc("hot")).limit(10).show()完整程序结构
object rating {
def main(args: Array[String]): Unit = {
// 创建SparkSession和StreamingContext
// ...
// 读取数据流
val stream = ssc.textFileStream("E:\\data\\StreamingData")
val splitData = stream.map { x =>
val y = x.split("\t")
(y(0), y(1), y(2).toInt)
}
import spark.implicits._
// 对每个批次的RDD做处理
splitData.foreachRDD { line =>
val dataFrame = line.toDF("UserID", "BookID", "Ratings")
// 步骤1:用户统计
// ...
// 步骤2:书籍统计
// ...
// 步骤3:计算热度
// ...
// 步骤4:保存到Hive并打印
// ...
}
ssc.start()
ssc.awaitTermination()
ssc.stop()
}
}代码说明
| 代码 | 说明 |
|---|---|
withColumn("hot", ...) | 新增 hot 列,值是热度计算公式 |
sort(desc("hot")) | 按热度降序排列 |
limit(10) | 取前10名 |
.write.mode("overwrite").saveAsTable(...) | 覆盖写入 Hive 表 |
foreachRDD { line => ... } | 对每个批次的 RDD 做处理 |
💡 说明:
- 因为是模拟场景,时间间隔只有60秒,大部分用户只评了一次
- 真实场景中,时间间隔建议设为24小时以上
- 用
mode("overwrite")每次覆盖,表中始终是最新的热度排行
第九部分:常见问题与排错指南
9.1 启动和配置类问题
问题1:启动后没有数据输出
现象:程序启动了,但一直没有输出 原因:
- 数据源没数据(nc 没输入、HDFS 目录没新文件)
- 批处理间隔太长,还没到计算时间
- 数据源地址/端口写错了 解决:
- 确认数据源有数据产生
- 把批处理间隔设短一点测试
- 检查地址、端口、路径是否正确
问题2:Socket 连接不上
现象:报 Connection refused 原因:
- nc 没启动
- 主机名/端口写错了
- 防火墙没开端口 解决:
- 确认 nc -lk 端口 已经启动
- 确认主机名和端口正确
- 检查防火墙设置
问题3:textFileStream 不处理文件
现象:文件上传了但没处理 原因:
- 文件是程序启动前就存在的(只处理新文件)
- 文件不是"移动"进去的,而是在目录里修改的
- 文件上传太慢,还没传完就被检测到了 解决:
- 程序启动后再上传文件
- 先传到别的目录,再 mv 到监控目录
- 确保文件完整写入后再出现在监控目录
9.2 窗口操作类问题
问题4:窗口操作报错
现象:报窗口长度或滑动步长相关错误 原因:窗口长度或滑动步长不是批处理间隔的整数倍 解决:
- 确保窗口长度 = N × 批处理间隔
- 确保滑动步长 = M × 批处理间隔
- 比如批处理1秒,窗口可以是3秒、5秒,不能是2.5秒
⚠️ 重要:窗口长度和滑动步长必须是批处理间隔的整数倍!
问题5:窗口结果不符合预期
现象:统计结果和预期不一样 原因:
- 窗口长度和滑动步长搞混了
- 滑动步长比窗口小,有重叠;比窗口大,有遗漏 解决:
- 确认窗口长度和滑动步长的含义
- 画图理解窗口滑动的过程
9.3 输出操作类问题
问题6:foreachRDD 写入数据库报错
现象:报 NotSerializableException 或连接错误 原因:
- 在 Driver 端创建了连接对象(连接对象不能序列化)
- 每条记录创建一个连接,太多了 解决:
- 用 foreachPartition,每个分区创建一个连接
- 或者用连接池
- 参考7.4节的正确写法
问题7:DStream 怎么排序?
现象:DStream 没有 sort 方法 原因:DStream API 确实没有 sort 解决:
- 用 transform() 转成 RDD,用 RDD 的 sortByKey 排序
val sorted = dstream.transform(rdd => {
rdd.map(x => (x._2, x._1))
.sortByKey(false)
.map(x => (x._2, x._1))
})9.4 排错通用思路
- 先本地测试:用 Socket 数据源,手动输入数据测试
- 加 print():每个步骤都 print 一下,看哪一步出问题
- 设短间隔:测试时批处理间隔设短一点,不用等太久
- 看日志:设置日志级别为 WARN 或 ERROR,减少干扰
- 确认参数:窗口长度、滑动步长、批处理间隔的关系
第十部分:实习 / 面试高频考点
10.1 概念类(高频)
Q1:Spark Streaming 是什么?和 Spark Core 有什么区别?
Spark Streaming 是 Spark 的实时计算框架,用于处理流式数据。 区别:
- Spark Core 处理的是有界的批量数据(RDD)
- Spark Streaming 处理的是无界的实时数据流(DStream)
- Spark Streaming 底层还是 RDD,是微批处理模式
Q2:Spark Streaming 的运行原理?
微批处理(Micro-batching):
- 把实时数据流按时间片(批处理间隔)切成一小段一小段
- 每一段数据对应一个 RDD
- 对每个 RDD 进行批处理
- 批量输出结果 DStream 就是一系列 RDD 的序列。
Q3:什么是 DStream?和 RDD 有什么关系?
DStream(离散流)是 Spark Streaming 的数据抽象,代表持续的数据流。 关系:
- DStream 底层是一系列 RDD,每个时间片一个 RDD
- 对 DStream 的操作会转化为对每个 RDD 的操作
- DStream 也有惰性求值的特点
Q4:DStream 和 DataFrame 有什么区别?
- DStream 是 Spark Streaming 的数据抽象,是 RDD 的序列,用于流处理
- DataFrame 是 Spark SQL 的数据抽象,是结构化表格,用于批处理
- 可以用 foreachRDD + toDF() 把 DStream 转成 DataFrame 来用 Spark SQL
Q5:Spark Streaming 的数据源有哪些?
基础来源:
- Socket(socketTextStream):测试用
- HDFS 文件(textFileStream):监控目录新文件 高级来源:
- Kafka:生产环境最常用
- Flume:日志采集
- Kinesis:AWS 的消息队列
Q6:DStream 的操作分几类?
三类:
- 转换操作(Transformation):map、flatMap、filter、reduceByKey、transform 等
- 窗口操作(Window):window、reduceByKeyAndWindow 等
- 输出操作(Output):print、saveAsTextFiles、foreachRDD 等
Q7:窗口操作的两个关键参数是什么?有什么规则?
两个参数:
- 窗口长度(Window Length):窗口包含多长时间的数据
- 滑动步长(Slide Interval):窗口多久滑动一次 规则:窗口长度和滑动步长都必须是批处理间隔的整数倍。
Q8:reduceByKey 和 reduceByKeyAndWindow 有什么区别?
- reduceByKey:只对当前批次的数据按 Key 聚合
- reduceByKeyAndWindow:对窗口内多个批次的数据按 Key 聚合
- 一个是单批次,一个是跨批次(窗口)
- reduceByKeyAndWindow 多了窗口长度和滑动步长两个参数
Q9:foreachRDD 的正确使用方式是什么?
正确方式是用 foreachPartition,每个分区创建一个连接:
scaladstream.foreachRDD { rdd => rdd.foreachPartition { partition => val conn = createConnection() partition.foreach(record => conn.send(record)) conn.close() } }错误方式:
- 在 Driver 端创建连接(连接对象不能序列化)
- 每条记录创建一个连接(开销太大) 优化方式:用连接池复用连接。
Q10:DStream 怎么排序?
DStream 没有 sort 方法,需要用 transform() 转成 RDD 来排:
scaladstream.transform(rdd => { rdd.map(x => (x._2, x._1)).sortByKey(false).map(x => (x._2, x._1)) })利用 transform 可以调用 RDD 的任何操作。
10.2 原理类(中频)
Q11:Spark Streaming 是真正的实时吗?
不是严格的实时,是微批处理。
- 延迟是秒级(取决于批处理间隔)
- 真正的实时(毫秒级)是 Flink、Storm 这种
- 但对于大部分场景,秒级延迟已经够用了
Q12:Spark Streaming 和 Flink 有什么区别?
- Spark Streaming:微批处理,延迟较高(秒级),吞吐量大,和 Spark 生态集成好
- Flink:真正的流处理,延迟低(毫秒级),有状态计算更完善
- 现在新的项目更多用 Flink 或 Structured Streaming
- 但老项目很多还在用 Spark Streaming
Q13:Spark Streaming 为什么被标记为 Deprecated?
从 Spark 3.4.0 开始标记为不推荐使用。 原因:
- 新的实时框架是 Structured Streaming(基于 Spark SQL)
- Structured Streaming 更统一、更容易用
- 但 3.5.1 还支持,老项目还能用
Q14:DStream 的转换操作是惰性的吗?
是的,和 RDD 类似。 转换操作只是生成执行计划,不会真正执行。 需要输出操作(print、save、foreachRDD 等)才会触发计算。
10.3 实操类(高频)
Q15:怎么创建 StreamingContext?
scala// 方式1:从SparkConf创建 val ssc = new StreamingContext(conf, Seconds(10)) // 方式2:从SparkContext创建 val ssc = new StreamingContext(sc, Seconds(10)) // 方式3:从SparkSession创建(有SparkSession时) val ssc = new StreamingContext(spark.sparkContext, Seconds(10))参数是批处理时间间隔。
Q16:启动和停止 StreamingContext 的方法?
- 启动:
ssc.start()- 等待终止:
ssc.awaitTermination()- 停止:
ssc.stop()注意:start 之后才真正开始处理数据。
Q17:textFileStream 监控目录有什么注意事项?
- 只处理新文件,已存在的文件不处理
- 文件必须是"移动"或"重命名"到目录中的
- 修改已有文件不会触发处理
- 程序启动后再上传文件才会被处理
Q18:怎么把 DStream 和 Spark SQL 结合用?
用 foreachRDD,在里面把 RDD 转成 DataFrame:
scaladstream.foreachRDD { rdd => val df = rdd.toDF("col1", "col2") df.createOrReplaceTempView("table") spark.sql("select * from table").show() }这样就可以用 Spark SQL 的 API 来处理流数据了。
附录:习题解析
选择题解析
1、答案:C 解析:在 Spark Streaming 中,DStream 的输出操作可以将数据写入文件系统、数据库或其他应用中。
2、答案:B 解析:DStream 的转换操作中(如 map、filter、reduceByKey 等)每个批次的处理都是独立于其他批次的,批次之间不累积。
3、答案:A 解析:DStream 窗口操作中,window() 方法返回一个基于源 DStream 的窗口批次计算后得到的新 DStream;countByWindow() 方法返回基于滑动窗口的 DStream 中的元素的数量;reduceByWindow() 方法基于滑动窗口对源 DStream 中的元素进行聚合操作,得到一个新的 DStream;reduceByKeyAndWindow() 方法基于滑动窗口对元素为 (K,V) 键值对的 DStream 中的值,按 K 使用 func 函数进行聚合操作,得到一个新的 DStream。
4、答案:D 解析:DStream 的窗口操作中,窗口长度和滑动步长两个参数值必须是批处理间隔的整数倍。
5、答案:D 解析:DStream 的转换操作中,transform() 方法通过对源 DStream 的每个 RDD 应用 func 函数返回一个新的 DStream,用于在 DStream 上进行 RDD 的任意操作。
6、答案:B 解析:reduceByKey() 方法用于对单个 RDD 中的键值对按 Key 进行聚合操作;reduceByKeyAndWindow() 方法用于对滑动窗口内的多个 RDD 中的键值对按 Key 进行聚合,两者操作的数据源不同。
7、答案:A 解析:print() 方法用于在 Driver 中输出 DStream 中数据的前 10 个元素。
8、答案:D 解析:union() 方法用于合并两个 DStream,生成一个包含两个 DStream 中所有元素的新的 DStream。
9、答案:C 解析:rdd.foreachPartition() 方法正确使用方法为:为 RDD 的每一个分区创建一个单独的连接对象,使用该连接对象输出该 RDD 分区的所有数据至外部系统中。 A 选项直接在 Driver 端创建连接对象,该方法需要对连接对象进行序列化,并从 Driver 端发送到 Worker 上,而连接对象很难在不同机器间进行序列化操作; B 选项为每一个记录创建一个连接对象,该方法会产生非常大的开销,且可能会显著降低系统的整体吞吐量。
10、答案:B 解析:socketTextStream() 方法是 Spark Streaming 提供的专门用于通过 TCP Socket 连接读取数据的 API,参数为 (hostname, port)。
操作题解析
题目
从 Socket 端口读取菜单数据(类别、名称、价格),实时统计所有菜品的总价格。
实现代码
import org.apache.spark.sql.SparkSession
import org.apache.spark.streaming.{Seconds, StreamingContext}
object menu {
def main(args: Array[String]): Unit = {
// 初始化SparkSession和StreamingContext
val spark = SparkSession.builder()
.master("local[*]")
.appName("Streaming")
.getOrCreate()
spark.sparkContext.setLogLevel("ERROR")
val ssc = new StreamingContext(spark.sparkContext, Seconds(60))
// 获取监听端口数据
val input = ssc.socketTextStream("master", 8888)
// 数据处理:切分 → 取价格 → 求和
val sum = input.map(line => {
val data = line.split(" ")
((data(0), data(1)), data(2).toDouble)
})
.map(_._2)
.reduce(_ + _)
// 输出结果
sum.print()
// 启动程序
ssc.start()
ssc.awaitTermination()
}
}代码说明
| 代码 | 说明 |
|---|---|
socketTextStream("master", 8888) | 从 master 的 8888 端口读取数据 |
line.split(" ") | 按空格切分每行数据 |
((data(0), data(1)), data(2).toDouble) | 转成 ((类别, 名称), 价格) 的格式 |
.map(_._2) | 只取价格部分 |
.reduce(_ + _) | 求和(所有菜品价格相加) |
sum.print() | 打印结果 |
💡 说明:这里的 reduce 是对每个批次内的数据求和,不是累积的。 如果要累积求和,需要用 updateStateByKey 或窗口操作。
笔记版本:V1.0 对应教材:《Spark大数据技术与应用(第3版)》人民邮电出版社 对应项目:项目6 实时计算书籍热度——Spark Streaming实时计算框架 最后更新:2026年8月