小菜鸟

java菜鸟号正在起航

RDD 持久化:缓存(快,作业内复用)vs 检查点(稳,跨作业复用)

Spark 默认不存数据,每次 Action 都从头算一遍。如果同一个 RDD 要被用两次,默认行为就是算了两次。

持久化就是”把中间结果存下来,下次直接用”。

两种方式:

缓存(Cache / Persist) 检查点(Checkpoint)
存哪 内存 / 本地磁盘 HDFS(分布式存储)
多快 快(内存级) 慢(要写 HDFS)
血缘关系 保留 切断
什么时候用 同一个 Job 里多次使用 跨 Job 复用,或者要切断长血缘
最常用方式 .cache() sc.setCheckpointDir(...) + .checkpoint()

记住:缓存是”快但临时”,检查点是”慢但持久”。两个一起用是黄金组合。

缓存(Cache / Persist):内存里存一份,下次直接用

val rdd = sc.textFile("data.txt")
  .flatMap(_.split(" "))
  .map((_, 1))
  .reduceByKey(_ + _)

rdd.cache()       // 标记:需要缓存
rdd.count()       // 第一次 Action:算 + 缓存
rdd.collect()     // 第二次 Action:直接从缓存读,不重算

.cache() 默认存内存(MEMORY_ONLY)。如果内存不够,数据会丢失,然后重算。

.persist(StorageLevel.XXX) 可以选存储级别:

存储级别 说明 什么时候用
MEMORY_ONLY 纯内存,不序列化 默认,数据量小用
MEMORY_ONLY_SER 内存,序列化 数据量大,用序列化省空间
MEMORY_AND_DISK 内存优先,不够写磁盘 数据不确定能不能装下内存
MEMORY_AND_DISK_SER 内存优先,序列化,不够写磁盘 数据量大、内存紧张
DISK_ONLY 纯磁盘 内存完全不够,或数据不用频繁读

缓存的时机: .cache() 只是”标记”,真正缓存发生在第一次 Action 的时候。

缓存的删除:

  • 内存不够时 LRU 自动淘汰
  • 手动:.unpersist()

检查点(Checkpoint):写入 HDFS,切断血缘

检查点把 RDD 数据写到 HDFS,同时切断血缘关系——后续 RDD 直接从检查点读数据,不追溯前面的依赖链。

// 先设检查点目录
sc.setCheckpointDir("hdfs://namenode:8020/checkpoint")

val rdd = sc.textFile("data.txt")
  .flatMap(_.split(" "))
  .map((_, 1))
  .reduceByKey(_ + _)

rdd.checkpoint()   // 标记检查点
rdd.count()        // 触发:写检查点

检查点的特点:

阅读全文 »

RDD 依赖关系:窄依赖快(无 Shuffle)、宽依赖慢(要 Shuffle),Stage 按宽依赖切

RDD 的依赖关系记录了一个 RDD 是怎么从父 RDD 算出来的。它干三件事:

作用 说明
容错 分区丢了,沿着依赖关系找到父 RDD,只重算丢失的分区
Stage 划分 宽依赖是”切 Stage”的边界,窄依赖可以合并到一个 Stage
性能优化 宽依赖 = Shuffle,能避免就避免

一句话:依赖关系 = 血缘 + 调度边界 + 容错依据

窄依赖 vs 宽依赖

窄依赖 宽依赖
子分区依赖几个父分区 1 个(或少数几个) 多个
有没有 Shuffle 没有
能不能在同一 Stage 不能(必须切 Stage)
容错代价 低(只重算一个父分区) 高(重算多个父分区)
典型算子 mapfilterflatMap reduceByKeygroupByKeyjoin

窄依赖的例子:map

父 RDD 有 2 个分区,map 之后子 RDD 也是 2 个分区,一一对应。

父分区0 → 子分区0
父分区1 → 子分区1

宽依赖的例子:reduceByKey

父 RDD 的多个分区的数据,按 key 重新分布后,可能落到同一个子分区。

父分区0 (a, b) ──┐
                 ├── 子分区0 (a)
父分区1 (a, c) ──┘
                 ├── 子分区1 (b, c)

宽依赖意味着数据要跨节点重新分布 → Shuffle。

怎么看依赖关系

val rdd = sc.textFile("data.txt")
  .flatMap(_.split(" "))
  .map((_, 1))
  .reduceByKey(_ + _)

// 看血缘(树形结构)
println(rdd.toDebugString)
阅读全文 »

Spark 序列化:什么时候需要序列化?Java vs Kryo 怎么选?NotSerializableException 怎么破?

序列化在 Spark 里发生在三个地方:

  1. Driver → Executor 传数据:广播变量、闭包(函数)要序列化后才能发给 Executor
  2. Shuffle:Map 输出的数据要序列化后通过网络传给 Reduce
  3. 缓存:RDD 存到内存/磁盘时,如果用了序列化存储级别,需要序列化

说白了:任何跨节点或跨存储的数据传递,都需要序列化。

Java 序列化 vs Kryo 序列化

Spark 默认用 Java 序列化。它不用配,啥都能序列化,但慢、体积大。

Kryo 序列化是备选方案:快(10倍)、体积小(1/3~1/5),但要手动配置。

对比项 Java 序列化 Kryo 序列化
配置 不用配 要配
速度 快 10 倍
体积 小 3-5 倍
兼容性 高(任何 Serializable) 中(需注册类)

什么时候用 Kryo? 任何时候。Spark 官方文档也建议生产环境用 Kryo。

Kryo 配置三步走

第一步:在 SparkConf 里启用 Kryo

val conf = new SparkConf()
  .setAppName("MyApp")
  .set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")

第二步:注册自定义类

conf.registerKryoClasses(Array(
  classOf[Person],
  classOf[Order],
  classOf[Array[String]]
))

注册的目的是让 Kryo 提前知道这些类,不用运行时动态发现,省时间。

第三步(推荐):强制注册

conf.set("spark.kryo.registrationRequired", "true")
阅读全文 »

Spring 上下文构建源码深度解析:从 ClassPathXmlApplicationContext 到 IOC 容器就绪

Spring 上下文(ApplicationContext)是 IOC 容器的核心载体,负责配置加载、BeanDefinition 管理、Bean 实例化与初始化的全流程。以 ClassPathXmlApplicationContext 为例,其初始化过程围绕 refresh() 方法展开,这是 Spring 最核心的源码链路之一。从 “构造函数初始化→refresh() 全景流程→关键子流程拆解→核心设计思想” 四个维度,彻底讲透上下文构建的每一步。

上下文初始化入口:ClassPathXmlApplicationContext 构造函数

创建 ClassPathXmlApplicationContext 实例时,仅需一行代码,但背后触发了完整的初始化流程:

ApplicationContext context = new ClassPathXmlApplicationContext("spring-lifecycle.xml");

构造函数核心逻辑

构造函数的本质是 “初始化配置路径 + 触发上下文刷新”,源码如下(已简化关键逻辑):

// ClassPathXmlApplicationContext 构造函数
public ClassPathXmlApplicationContext(String[] configLocations, boolean refresh, ApplicationContext parent) {
    super(parent); // 初始化父上下文(若有)
    setConfigLocations(configLocations); // 1. 保存配置文件路径(如 "spring-lifecycle.xml")
    if (refresh) {
        refresh(); // 2. 核心:触发上下文刷新(IOC 容器初始化的入口)
    }
}
阅读全文 »

Spark 数据读写:文本、JSON、CSV、SequenceFile,什么场景用什么格式

Spark 读写数据,选择什么格式取决于两个因素:

  1. 数据结构:是纯文本、表格、还是键值对?
  2. 性能要求:要快(二进制)还是要通用(文本)?
格式 适用场景 性能 跨语言
文本文件 日志、简单输出 一般 可以
JSON API 数据、半结构化 一般 可以
CSV 表格数据、Excel 导出 一般 可以
SequenceFile Hadoop 生态内部交换 不可以
对象序列化 Spark 内部中间结果 不可以

核心原则:要通用用文本/JSON/CSV,要性能用二进制格式(SequenceFile/Parquet)

文本文件:最通用,但也最基础

什么时候用: 日志文件、简单的文本输出、不需要结构化的数据。

读:

val rdd = sc.textFile("hdfs:///data/logs/2026-07-29.log")

// 指定最小分区数
val rdd = sc.textFile("hdfs:///data/logs/", minPartitions=10)

写:

rdd.saveAsTextFile("hdfs:///output/result")

写入后是一个目录,里面每个分区一个 part-* 文件。

JSON:半结构化的”万能格式”

什么时候用: API 响应、日志中的结构化数据、跨语言数据交换。

有两种方式,推荐用 DataFrame 方式:

方式1:DataFrame 方式(推荐)

val spark = SparkSession.builder().appName("JSON").getOrCreate()

// 读 JSON
val df = spark.read.json("hdfs:///data/people.json")

// 自动推断 schema,可以做 SQL 查询
df.createOrReplaceTempView("people")
spark.sql("SELECT name, age FROM people WHERE age > 18").show()

// 写 JSON
df.write.json("hdfs:///output/json_result")

方式2:RDD 方式(手动解析)

import com.fasterxml.jackson.databind.ObjectMapper

case class Person(name: String, age: Int)
val mapper = new ObjectMapper()

val rdd = sc.textFile("hdfs:///data/people.json")
  .map(line => mapper.readValue(line, classOf[Person]))
  .filter(_ != null)

rdd.map(mapper.writeValueAsString).saveAsTextFile("hdfs:///output/json_rdd")

CSV:表格数据的”标准格式”

什么时候用: Excel 导出、数据库导出的表格数据、数据分析常用的格式。

DataFrame 方式(推荐):

val spark = SparkSession.builder().appName("CSV").getOrCreate()

// 读 CSV
val df = spark.read
  .option("header", "true")        // 第一行是表头
  .option("inferSchema", "true")   // 自动推断数据类型
  .csv("hdfs:///data/users.csv")

// 查
df.show()
df.filter("age > 18").show()

// 写 CSV
df.write
  .option("header", "true")
  .csv("hdfs:///output/csv_result")

常用配置:

配置 说明 示例
header 第一行是不是表头 true / false
inferSchema 自动推断列类型 true / false
delimiter 分隔符(默认逗号) `delimiter=” “`
quote 引号字符 quote="\""
nullValue 空值表示 nullValue="NULL"

SequenceFile:Hadoop 的二进制键值对格式

什么时候用: Hadoop 生态内部数据交换、需要高效存储键值对数据。

import org.apache.hadoop.io.{IntWritable, Text}

// 读 SequenceFile(必须指定 key 和 value 的 Hadoop 类型)
val seqRDD = sc.sequenceFile[Text, IntWritable]("hdfs:///data/seqfile")

// 转成普通类型
val rdd = seqRDD.map { case (k, v) => (k.toString, v.get()) }

// 写 SequenceFile(要转回 Hadoop 类型)
rdd.map { case (k, v) => (new Text(k), new IntWritable(v)) }
  .saveAsSequenceFile("hdfs:///output/seq_result")

注意: SequenceFile 是 Hadoop 专属格式,跨语言兼容性差。如果你的数据不需要被 Hadoop 外的系统读,用 SequenceFile 性能很好。

对象序列化:Spark 内部的”快速缓存”

什么时候用: 在 Spark 作业之间传递中间结果,不需要跨语言读。

// 定义可序列化的类
case class User(id: Long, name: String) extends Serializable

val users = sc.parallelize(Seq(User(1, "Alice"), User(2, "Bob")))

// 保存为对象文件
users.saveAsObjectFile("hdfs:///output/objects")

// 读取
val loaded = sc.objectFile[User]("hdfs:///output/objects")

性能优化: 默认用 Java 序列化,速度慢。换成 Kryo 序列化:

val conf = new SparkConf()
  .set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
  .registerKryoClasses(Array(classOf[User]))

高级优化技巧

1. 输出压缩

// 压缩 CSV 输出
df.write
  .option("header", "true")
  .option("compression", "snappy")
  .csv("hdfs:///output/compressed")

支持的压缩算法:snappygziplzobzip2

2. 分区写入

// 按日期分区写入
df.write
  .partitionBy("date")
  .parquet("hdfs:///output/partitioned")

3. 控制输出文件数量

// 减少分区数,避免生成太多小文件
df.coalesce(10).write.csv("hdfs:///output/fewer_files")

// 增加分区数,提高并行度
df.repartition(50).write.csv("hdfs:///output/more_parallel")
0%