小菜鸟

java菜鸟号正在起航

Spark RDD 编程:创建 RDD → 转换(惰性)→ 行动(触发),三步走

RDD 编程的核心流程就三步:

步骤 做什么 关键方法
1. 创建 RDD 从数据源生成 RDD parallelize / textFile
2. 转换 RDD 数据处理,惰性执行 map / filter / reduceByKey
3. 行动触发 触发计算,返回结果或写存储 collect / count / saveAsTextFile

记住:转换(Transformation)懒,行动(Action)勤。转换只是记下”要做什么”,行动才真正开始跑。

创建 RDD:从集合或文件来

1.1 从内存集合创建(测试用)

val rdd1 = sc.parallelize(List(1, 2, 3, 4, 5))
val rdd2 = sc.makeRDD(List(1, 2, 3, 4, 5))   // 推荐,简洁

parallelizemakeRDD 功能一样。numSlices 参数控制分区数,默认是 CPU 核心数。

1.2 从外部文件创建(生产用)

// 读本地文件(加 file:// 前缀)
val lines = sc.textFile("file:///path/to/file.txt")

// 读 HDFS 文件
val lines = sc.textFile("hdfs:///user/data/input.txt")

textFile 按行读取,每行作为 RDD 的一个元素。分区数默认跟 HDFS 块数一致(128MB/块)。

转换算子:数据处理的”流水线”

转换算子返回新 RDD,惰性执行。分三类:单值型、双值型、键值型。

2.1 单值型:一个 RDD → 另一个 RDD

算子 做什么 什么时候用
map 每个元素一对一转换 对每个元素做变换,比如 _ * 2
flatMap 每个元素 → 多个元素(压平) 拆分字符串、展开嵌套列表
filter 保留满足条件的元素 过滤数据,比如去掉空行
mapPartitions 按分区批量处理 每个分区创建一次连接(比 map 高效)
distinct 去重 去掉重复元素
coalesce 减少分区(不 Shuffle) 合并小分区,减少 Task 数
repartition 调整分区(会 Shuffle) 增加分区或均衡数据分布

示例:

阅读全文 »

Scala 样例类(Case Class):模式匹配的完美搭档

样例类(Case Class)是 Scala 中一种特殊的类,专门为模式匹配不可变数据建模设计。它自动生成了一系列常用方法,大幅简化了数据封装和匹配的代码。本文将详细解析样例类的特性、用法及适用场景。

样例类的定义与基本特性

样例类通过 case class 关键字定义,与普通类相比,它具有以下默认特性:

// 定义样例类(无需 new 关键字即可创建实例)
case class Student(id: Int, name: String, age: Int)

// 创建实例(自动生成 apply 方法,无需 new)
val stu1 = Student(1, "Alice", 20)
val stu2 = Student(2, "Bob", 21)

自动生成的方法

样例类会自动生成以下方法,无需手动实现:

  1. apply 方法:允许直接通过类名创建实例(如 Student(1, "Alice", 20)),无需 new 关键字。
  2. unapply 方法:支持模式匹配(核心特性),可从实例中提取构造参数。
  3. toString 方法:返回格式化的字符串(如 Student(1, Alice, 20)),便于调试。
  4. equalshashCode 方法:基于构造参数实现,支持值比较(而非引用比较)。
  5. copy 方法:用于创建实例的副本,可修改部分参数(适合不可变数据)。

样例类的核心特性详解

不可变的构造参数

样例类的构造参数默认被 val 修饰(不可变),确保实例创建后无法修改:

case class Book(isbn: String, title: String)

val book = Book("978-0134685991", "Scala Programming")

// 编译错误:无法修改 val 变量
// book.title = "New Title"
阅读全文 »

Scala 视图(View):懒加载的集合操作机制

在 Scala 中,视图(View)是一种特殊的集合转换机制,它通过懒加载(Lazy Evaluation) 延迟执行集合操作,直到真正需要结果时才计算。这种特性可以显著提升处理大型集合或复杂操作时的性能,避免不必要的中间计算。本文将详细解析视图的工作原理、使用场景及优势。

视图的基本概念

视图本质上是对集合操作的延迟封装。当对集合应用 view 方法后,后续的转换操作(如 mapfilterflatMap 等)不会立即执行,而是被记录下来,直到调用触发计算的方法(如 toListsizeforeach 等)时,才会一次性执行所有操作。

核心特性:

  • 延迟执行:转换操作仅在需要结果时才执行,而非立即计算。
  • 避免中间集合:普通集合操作会产生多个中间集合(如 list.filter(...).map(...) 会先生成过滤后的集合,再生成映射后的集合),而视图不会创建中间集合,直接在最终计算时一次性完成所有操作。
  • 适用于大型集合:对于数据量巨大或计算成本高的场景,视图能减少内存占用和计算开销。

视图的使用方法

创建视图

通过集合的 view 方法创建视图,后续操作将变为懒加载:

val numbers: List[Int] = List(1, 2, 3, 4, 5, 6)

// 创建视图,后续操作(filter)将延迟执行
val evenView = numbers.view.filter(_ % 2 == 0)

// 此时 filter 尚未执行,evenView 仅记录了操作
println(evenView)  // 输出:View(<not computed>)

触发计算

当调用需要实际结果的方法时,视图会执行所有延迟的操作:

阅读全文 »

RDD 深度解析:Spark 最核心的数据结构,为什么它这么设计

RDD 是 Spark 最核心的抽象——它是分布式数据集合,也是 Spark 实现”内存计算 + 容错 + 并行”的基础。

RDD 是什么?”Resilient Distributed Dataset”:

  • Resilient(弹性):数据丢了能重算,不用备份副本
  • Distributed(分布式):数据分散在多台机器上,并行处理
  • Dataset(数据集):一个不可变的、分区的元素集合

为什么 Spark 要设计 RDD? 分布式计算面临三个难题:

  1. 数据怎么分散到多台机器?→ 分区(Partitions)
  2. 某台机器挂了数据丢了怎么办?→ 血缘关系(Dependencies)
  3. 怎么知道数据在哪台机器上?→ 首选位置(Preferred Locations)

RDD 的五大属性就是回答这三个问题的

RDD 的五大属性(源码定义)

// Spark 源码注释:每个 RDD 由五个主要属性描述
// - A list of partitions(分区列表)
// - A function for computing each split(分区计算函数)
// - A list of dependencies on other RDDs(依赖关系)
// - Optionally, a Partitioner for key-value RDDs(分区器)
// - Optionally, a list of preferred locations to compute each split on(首选位置)

属性1:分区列表(Partitions)——数据怎么分散

protected def getPartitions: Array[Partition]

RDD 的数据分散在多个分区里,每个分区是数据的一个子集,分布在某个 Executor 上。

问题 RDD 的答案
数据怎么分散到多台机器? 切成多个分区,每个分区可以放到不同的节点
并行度是多少? 分区数 = 并行度,一个分区对应一个 Task
怎么控制并行度? repartition(n) 增加分区,coalesce(n) 减少分区

举例: 读取 HDFS 上一个 1GB 的文件(块大小 128MB),RDD 有 8 个分区,每个分区对应一个 HDFS 块。

属性2:分区计算函数(Compute)——数据怎么算

def compute(split: Partition, context: TaskContext): Iterator[T]
阅读全文 »

Spark 核心组件与调度:提交一个任务,Driver、Executor、Stage、Task 都经历了什么

Spark 跑一个任务,涉及四个核心角色:

角色 比喻 做什么
Driver 项目经理 解析代码、拆任务、调度
Executor 一线工人 真正执行 Task
Cluster Manager 人事部 分配资源(CPU、内存)
Task 最小工作单元 一个分区的数据处理

一个 Spark 任务从提交到结束,就是这四个角色协同完成一件事:把用户代码拆成 Task,分配到 Executor 上并行执行

四个角色各管什么

spark架构

Driver:任务的”大脑”

  • 解析用户代码,生成 DAG(有向无环图)
  • 把 DAG 切成 Stage,Stage 再切成 Task
  • 向 Cluster Manager 申请 Executor
  • 把 Task 分配到 Executor 上跑
  • 监控任务状态,失败了重试

Executor:任务的”手脚”

  • 跑 Task,执行具体的 mapfilterreduce 操作
  • 缓存 RDD 数据(内存/磁盘)
  • 向 Driver 汇报执行进度

Cluster Manager:资源的”管家”

  • 管理集群资源(CPU、内存)
  • 启动和停止 Executor
  • 支持 Standalone / YARN / K8s / Mesos

Task:最小的执行单元

  • 一个 Task 处理一个 RDD 分区
  • 所有 Task 加起来 = 一个 Stage 的全部计算

完整调度流程:提交一个 WordCount

spark调度分析

第一步:提交应用,启动 Driver

用户执行 spark-submit,Cluster Manager 启动 Driver 进程,Driver 创建 SparkContext。

spark-submit --master yarn --class WordCount myapp.jar
阅读全文 »
0%