Skip to content

项目6:实时计算书籍热度——Spark Streaming 实时计算框架 ​

先修基础:项目1-5(Spark概述、Scala基础、Spark Shell编程、Spark IDE编程、Spark SQL)


目录 ​


第一部分:项目背景与 Spark Streaming 概述 ​

1.1 项目背景 ​

为什么需要实时计算? ​

前面学的 Spark Core、Spark SQL 都是批处理——处理的是已经存好的历史数据。 但现实中很多场景需要实时处理:

  • 电商网站:实时统计商品热度、用户行为
  • 网站监控:实时监控访问量、异常流量
  • 金融风控:实时检测欺诈交易
  • 物联网:实时处理传感器数据

Spark Streaming 就是 Spark 的实时计算框架。

项目场景 ​

某电商平台想要实时计算书籍热度,把热度高的书推荐给用户。

书籍热度计算公式:

热度 = 用户平均评分 × 用户评分次数 × 0.3 + 书籍平均评分 × 书籍被评分次数

数据:BookRating.txt(用户对书籍的评分数据,评分范围1~5分)

本项目要做的事:

  1. 模拟实时产生的评分数据
  2. 用 Spark Streaming 实时读取和处理
  3. 实时计算书籍热度,取 Top10 存入 Hive

1.2 批处理 vs 流处理 ​

对比项批处理(Batch)流处理(Streaming)
数据特点有界、历史数据无界、实时到达
处理方式一次处理全部数据来一条处理一条(或小批量)
延迟高(分钟/小时级)低(秒级/毫秒级)
代表技术MapReduce、Spark Core、Spark SQLSpark 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)。

核心思想 ​

  1. 把实时数据流按时间片切成一小段一小段
  2. 每一段数据当成一个小批处理来算
  3. 批量输出结果

关键概念 ​

概念说明
批处理时间间隔(Batch Interval)切分数据的时间间隔,比如10秒切一次
DStream(离散流)Spark Streaming 的数据抽象,代表持续的数据流
每个时间片 → 一个RDDDStream 就是一系列 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 的步骤 ​

  1. 创建 StreamingContext:流处理的入口
  2. 创建输入源 DStream:指定从哪里读数据(Socket、HDFS、Kafka等)
  3. 操作 DStream:转换、窗口、输出等
  4. 启动流处理:调用 ssc.start()
  5. 等待停止:调用 ssc.awaitTermination()

⚠️ 注意:start() 之前只是设定执行计划,并没有真正运行;start() 之后才开始真正处理数据。


3.2 示例1:Socket 数据源(单词实时计数) ​

从网络端口读取数据,实时统计单词出现次数。

准备工作:安装 nc 工具 ​

nc(netcat)是一个网络工具,可以用来发送数据。

bash
# 检查是否安装
dnf list --installed | grep nc

# 如果没装,安装一下
dnf install -y nc

步骤1:启动 nc 发送数据(slave1 节点) ​

bash
# 在slave1上启动,监听8888端口
nc -lk 8888

步骤2:编写流处理程序(spark-shell) ​

scala
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 目录,有新文件进来就处理。

代码 ​

scala
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 目录上传文件:

bash
# 创建目录
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 常用数据源 ​

数据源方法说明
SocketsocketTextStream(host, port)从网络端口读取,学习测试用
HDFS文件textFileStream(directory)监控HDFS目录,新文件触发
KafkaKafkaUtils.createStream(...)生产环境最常用
FlumeFlumeUtils.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 的任何操作

示例 ​

scala
// 用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() 方法 ​

最基础的窗口操作,把数据按窗口重新分组。

示例 ​

scala
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 聚合。

示例:窗口内单词计数 ​

scala
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 对比 ​

对比项reduceByKeyreduceByKeyAndWindow
数据范围当前批次窗口内多个批次
时间维度无时间概念有时间窗口
参数只有聚合函数聚合函数 + 窗口长度 + 滑动步长
适用场景每批次统计近N分钟/小时统计

第七部分:DStream 输出操作 ​

7.1 常用输出操作 ​

方法说明
print()打印前10个元素(Driver端)
saveAsTextFiles(prefix, suffix)保存为文本文件
saveAsObjectFiles(prefix, suffix)保存为序列化对象文件
saveAsHadoopFiles(prefix, suffix)保存为Hadoop文件
foreachRDD(func)对每个RDD做任意输出操作(最灵活)

💡 注意:输出操作是触发计算的,没有输出操作的话,转换操作都不会执行。


7.2 print() 方法 ​

最简单的输出,在控制台打印前10个元素。

scala
wordCounts.print()

7.3 saveAsTextFiles() 方法 ​

保存为文本文件,每个批次生成一个文件夹。

scala
lines.saveAsTextFiles(
  "hdfs://master:8020/tipdm/data/saveAsTextFiles/sahf",  // 前缀
  "txt"                                                   // 后缀
)

生成的文件路径类似:sahf-时间戳.txt

💡 说明:每个批次生成一个文件夹,文件夹里是该批次的数据。


7.4 foreachRDD() 方法(重点+难点) ​

最灵活的输出操作,可以对每个 RDD 做任意操作。 常用于把结果写入 MySQL、Redis 等外部系统。

正确使用方式很重要! ​

❌ 错误写法1:在 Driver 端创建连接 ​

scala
dstream.foreachRDD { rdd =>
  val connection = createNewConnection()  // 在Driver端创建,错误!
  rdd.foreach { record =>
    connection.send(record)
  }
}

问题:连接对象在 Driver 创建,要序列化发送到 Worker,连接对象很难序列化。

❌ 错误写法2:每条记录创建一个连接 ​

scala
dstream.foreachRDD { rdd =>
  rdd.foreach { record =>
    val connection = createNewConnection()  // 每条记录创建一个,错误!
    connection.send(record)
    connection.close()
  }
}

问题:每条记录都创建连接,开销太大,严重影响性能。

✅ 正确写法:每个分区创建一个连接 ​

scala
dstream.foreachRDD { rdd =>
  rdd.foreachPartition { partitionOfRecords =>
    val connection = createNewConnection()  // 每个分区创建一个连接
    partitionOfRecords.foreach(record => connection.send(record))
    connection.close()
  }
}

原理:一个分区用一个连接,分区内多条记录复用连接。

⭐ 优化写法:用连接池 ​

scala
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 中创建表 ​

sql
create database spark;
use spark;
create table searchKeyWord(
  insert_time datetime,
  keyword varchar(30),
  search_count integer
);

步骤2:编写 Spark Streaming 程序 ​

scala
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:用户对书籍的评分数据

字段说明 ​

列索引字段说明
0UserID用户ID
1BookID书籍ID
2Rating评分(1~5分)

书籍热度公式 ​

热度 = 用户平均评分 × 用户评分次数 × 0.3 + 书籍平均评分 × 书籍被评分次数

公式中的四个变量:

  • u:用户平均评分(user_avg_rating)
  • x:用户评分次数(user_count)
  • y:书籍平均评分(book_avg_rating)
  • z:书籍被评分次数(book_count)

任务6.1:获取输入数据源 ​

需求 ​

  1. 写一个日志生成模拟器,模拟实时产生评分数据
  2. 用 Spark Streaming 实时读取新文件

步骤1:日志生成模拟器(CreateData.scala) ​

每隔60秒,随机从 BookRating.txt 中抽取100条记录,写入新文件。

scala
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) ​

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()
  }
}

运行方式 ​

  1. 先运行 CreateData 生成数据
  2. 再运行 rating 读取数据
  3. 每隔60秒处理一批新文件

💡 说明:两个程序要同时运行,一个产生数据,一个消费数据。


任务6.2:计算用户评分次数及平均评分 ​

需求 ​

统计每个用户的评分次数和平均评分。

实现代码 ​

scala
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:计算书籍被评分次数及平均评分 ​

需求 ​

统计每本书的被评分次数和平均评分,并和前面的数据合并。

实现代码 ​

scala
// 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 + 书籍平均评分 × 书籍被评分次数

实现代码 ​

scala
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()

完整程序结构 ​

scala
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 排序
scala
val sorted = dstream.transform(rdd => {
  rdd.map(x => (x._2, x._1))
     .sortByKey(false)
     .map(x => (x._2, x._1))
})

9.4 排错通用思路 ​

  1. 先本地测试:用 Socket 数据源,手动输入数据测试
  2. 加 print():每个步骤都 print 一下,看哪一步出问题
  3. 设短间隔:测试时批处理间隔设短一点,不用等太久
  4. 看日志:设置日志级别为 WARN 或 ERROR,减少干扰
  5. 确认参数:窗口长度、滑动步长、批处理间隔的关系

第十部分:实习 / 面试高频考点 ​

10.1 概念类(高频) ​

Q1:Spark Streaming 是什么?和 Spark Core 有什么区别? ​

Spark Streaming 是 Spark 的实时计算框架,用于处理流式数据。 区别:

  • Spark Core 处理的是有界的批量数据(RDD)
  • Spark Streaming 处理的是无界的实时数据流(DStream)
  • Spark Streaming 底层还是 RDD,是微批处理模式

Q2:Spark Streaming 的运行原理? ​

微批处理(Micro-batching):

  1. 把实时数据流按时间片(批处理间隔)切成一小段一小段
  2. 每一段数据对应一个 RDD
  3. 对每个 RDD 进行批处理
  4. 批量输出结果 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 的操作分几类? ​

三类:

  1. 转换操作(Transformation):map、flatMap、filter、reduceByKey、transform 等
  2. 窗口操作(Window):window、reduceByKeyAndWindow 等
  3. 输出操作(Output):print、saveAsTextFiles、foreachRDD 等

Q7:窗口操作的两个关键参数是什么?有什么规则? ​

两个参数:

  • 窗口长度(Window Length):窗口包含多长时间的数据
  • 滑动步长(Slide Interval):窗口多久滑动一次 规则:窗口长度和滑动步长都必须是批处理间隔的整数倍。

Q8:reduceByKey 和 reduceByKeyAndWindow 有什么区别? ​

  • reduceByKey:只对当前批次的数据按 Key 聚合
  • reduceByKeyAndWindow:对窗口内多个批次的数据按 Key 聚合
  • 一个是单批次,一个是跨批次(窗口)
  • reduceByKeyAndWindow 多了窗口长度和滑动步长两个参数

Q9:foreachRDD 的正确使用方式是什么? ​

正确方式是用 foreachPartition,每个分区创建一个连接:

scala
dstream.foreachRDD { rdd =>
  rdd.foreachPartition { partition =>
    val conn = createConnection()
    partition.foreach(record => conn.send(record))
    conn.close()
  }
}

错误方式:

  1. 在 Driver 端创建连接(连接对象不能序列化)
  2. 每条记录创建一个连接(开销太大) 优化方式:用连接池复用连接。

Q10:DStream 怎么排序? ​

DStream 没有 sort 方法,需要用 transform() 转成 RDD 来排:

scala
dstream.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 这种
  • 但对于大部分场景,秒级延迟已经够用了
  • 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:

scala
dstream.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 端口读取菜单数据(类别、名称、价格),实时统计所有菜品的总价格。

实现代码 ​

scala
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月

基于 Vite 强力驱动 | 纯静态轻量托管