把普通 Scala 集合的 map 改成 Spark RDD 的 map,函数看起来相似,执行责任却发生了变化。函数可能被序列化、送到执行端,在 action 触发时运行,还可能因再次计算或故障恢复而重复执行。捕获一个驱动端对象、修改一个外部变量、发送一次业务请求,都需要重新检查含义。

本篇使用真实 Spark 4.0.0 本地模式,运行两个有边界的实验:对同一条未缓存 RDD 执行两次 action,观察转换被重新执行;在闭包中捕获不可序列化对象,观察任务提交被拒绝。本地实验不模拟多机网络分区、执行器丢失或真实集群容错。

工程位于 examples/scala-lab/electives/E06。冻结 Spark 4.0.0、Scala 2.13.16、JDK 21,成功记录为 evidence/20261002-c/spark-local.log。主线 Scala 3.3.7 不参与这个独立工程的编译。

版本矩阵先于业务代码

Spark 4.0.0 官方文档要求 Java 17 或 21,并使用 Scala 2.13。Maven 发布的 spark-core_2.13:4.0.0 POM 进一步给出 scala-library 与 scala-reflect 为 2.13.16,因此本实验把编译器也固定为这一补丁版本。

组成部分 冻结值 核验依据
Spark Core 4.0.0 官方版本文档与发布坐标
Scala 2.13.16 Spark POM 的标准库与反射依赖
JDK 21.0.11 本机 Java 启动输出
执行环境 local[2] 程序显式配置
部署范围 单机 无集群提交步骤

Scala 语言版本相近,不代表生态库可以任意互换。_2.13 是二进制交叉发布标识,不能把 _3 产物直接改名放进 classpath。Scala 2 宏、scala-reflect、编译器插件和 TASTy 读取又有各自的兼容边界,必须按真实依赖路径验证。

Scala 3.8 起标准库与 Scala 2.13 互操作还出现新的版本分界,不能把早期 Scala 3 迁移示例概括为所有未来版本均双向兼容。本篇采用 Spark 自身要求的 2.13 组合,避免把框架实验同时变成未经验证的语言迁移实验。

JDK 启动参数也是组合的一部分。这里通过 Scala CLI 直接启动应用,没有使用 spark-submit 帮忙配置全部模块开放参数。序列化调试路径需要访问特定 JDK 内部包,runner 显式记录本次使用的 add-opens。缺少参数时出现的 IllegalAccessException,是启动配置失败,不是目标序列化反例通过。

transformation 保存计算描述

实验将三个整数分成两个分区,在 map 中把每个数字乘以 100。一个 LongAccumulator 记录转换执行次数,两个 collect 都校验最终总额为 600。

1
2
3
4
5
6
7
8
9
10
11
val executions = sc.longAccumulator("executions")
val rdd = sc.parallelize(List(1, 2, 3), 2).map { n =>
executions.add(1)
n * 100
}

assert(rdd.collect().sum == 600)
val first = executions.value
assert(rdd.collect().sum == 600)
val second = executions.value
assert(first == 3 && second == 6)

定义 rdd 时主要形成计算关系;collect 是 action,触发需要的执行。当前没有 persist 或 cache,因此第二次 collect 会再次沿计算关系处理三个元素。累计值从 3 变为 6,显示同一个 RDD 变量不代表计算结果已经被保存。

这里的两次运行来自应用主动调用两次 action,未注入 task 失败。它证明“同一转换可能执行不止一次”这一风险的一个直接场景,不能据此声称已经验证 Spark 的任务重试次数、推测执行或故障恢复策略。

Accumulator 在本例中是观测工具,不是业务数据库。Spark 对 action 与 transformation 中的累计更新有不同约定;官方文档明确提示转换重新执行时,累计更新可能重复。把它当作精确计费账本,会让执行重放改变业务金额。

业务侧副作用也面临相同问题。如果把 executions.add(1) 换成“扣款一次”,两次 collect 就可能发送两次请求。是否恰好写入一次,取决于外部系统的幂等键、提交协议和失败处理,不能从 map 看起来是纯函数的语法形式推断。

[PATTERN] 分布式转换里的副作用必须能够承受重执行。RDD 变量保存计算关系,不能当作唯一执行凭据。

闭包捕获了什么

反例对象有一个普通字段 value,没有实现 Serializable。实例方法里的 lambda 访问该字段,因此闭包需要捕获接收者对象,不能只看到源代码中的一个整数加法。

1
2
3
4
final class DriverOnly(val value: Int) {
def transform(sc: SparkContext): Array[Int] =
sc.parallelize(List(1, 2), 2).map(_ + value).collect()
}

调用 new DriverOnly(7).transform(sc) 时,Spark 清理并检查闭包,序列化失败,程序捕获包含 serializable 诊断的 SparkException。正常路径与拒绝路径在同一个进程中执行,最终还会断言拒绝确实发生。

这个例子说明字段访问可能间接捕获整个对象。真实驱动对象常常持有日志器、线程池、数据库连接或 SparkContext,即使 lambda 只用到一个配置值,也可能把更多状态带进捕获图。修复时应先理解捕获关系,再决定传递简单值或构造执行端资源。

把类简单标记为 Serializable 也未必正确。字段对象仍可能不可序列化;即使能序列化,一个打开的本地连接也不一定能在另一进程中恢复有效。序列化描述数据,不会自动复制系统资源的运行状态。

对于配置数字,可以在方法中先取出独立局部值,再让闭包捕获该值;对于外部客户端,通常需要按执行端生命周期建立连接并控制释放。这两类修复解决的问题不同,不能用“加 Serializable”统一处理。

本地模式仍有边界

local[2] 在本机运行两个工作线程,用于验证编程接口与部分任务执行过程。本次确实经过 Spark 的闭包序列化检查,因此捕获反例可以在本地复现;但这并不意味着本地线程与真实集群进程拥有相同的隔离与故障模型。

普通外部变量尤其容易造成误解。某些本地观察可能看见共享地址空间中的状态变化,而集群任务持有的是序列化后的副本。Spark 文档明确要求通过支持的分布式状态机制交互,不能依赖驱动变量被任务直接修改。

广播变量也不等于共享可变内存。它适合传播只读数据,更新策略需要单独设计。把可变对象放入广播后再尝试依靠其变化协调执行器,会混淆数据分发与一致性协议。

同时,本机文件路径在执行器上不一定存在,localhost 也指向执行器自己的主机。一个读取驱动桌面文件的函数,本地运行成功,集群运行仍可能失败。因此代码迁移到集群时,需要检查输入分发、路径、凭据和网络访问范围。

本篇显式绑定回环地址并关闭 Spark UI,避免为实验开放不必要的外部入口;这些配置不构成生产部署模板。集群认证、加密、资源配额与审计属于部署任务,应按实际集群管理器验收。

缓存改变重算成本,不提供事务

给 rdd 加上 persist,并在第一次 action 中完成物化,通常可以复用已缓存分区,减少后续重算。但缓存可能被驱逐,执行器可能丢失,分区也可能重新计算。缓存是一种执行优化,不能成为“扣款恰好一次”的正确性依据。

如果应用需要把计算结果写入外部存储,应考虑使用能够表达提交语义的输出机制,并对目标系统设计幂等键。例如按订单标识与业务版本建立唯一约束,可以让重复提交变成可识别的同一操作;仅按 task attempt 标识去重,可能无法覆盖应用重新提交整个作业。

失败响应还有结果未知问题。一次网络请求超时可能发生在服务端提交之后,直接重试会产生重复写;不重试又可能漏写。解决这类问题需要服务端可查询的操作标识与明确状态机,RDD API 不会代替它完成。

实验特意没有用缓存来“修复”计数从 3 到 6。这个变化正是要暴露的语义。若只追求让计数一直停在 3,就可能把正确性依赖错误地建立在某次缓存命中上。

观测与生产推论分开

Spark 的执行日志包含环境、driver 启动和模块诊断。程序结束前调用 sc.stop(),因此本机服务由实验自身释放。日志中的本地目录与端口是运行细节,不能拿来推断真实集群上的数据分布。

本次通过记录显示两次结果总额都为 600,转换执行累计 3 与 6,不可序列化捕获被拒绝。没有注入 executor 崩溃,也没有验证 shuffle 数据丢失、网络重试或分布式写入的端到端一致性。

若要进一步验证故障重执行,需要建立可控制的失败注入,记录 partition、stage、attempt 与外部操作键,并区分任务开始次数、成功次数和最终业务提交次数。单纯数日志行容易把重试日志、驱动打印和执行端输出混在一起。

采样任务也要避免把观测本身变成业务副作用。Accumulator 适合这个小实验,因为结果可在 action 结束后读取;对高频生产统计,还需考虑更新代价、重执行语义与监控系统的聚合方式。测试仪表应当明确它计数的是尝试还是成功。

闭包审查从依赖对象开始

闭包表达式只有一行,也可能携带很大的对象图。访问实例字段会引入接收者,接收者又可能保存连接池、线程或其他不可序列化对象。因此排查序列化错误时,应先列出表达式读取的自由变量,再沿字段依赖检查,而不是看到 lambda 短小就假定运输成本很低。把字段值复制成不可变局部数据,有时能缩小捕获范围,但仍需确认该数据本身适合序列化。

连接也不应通过闭包从 driver 运输给 executor。通常需要在执行端按适当范围创建客户端,并在相同范围关闭;具体是每分区创建还是使用连接池,要依据客户端线程安全性与资源配额决定。本实验没有连接外部数据库,因此不为某个客户端提供生命周期保证,只展示为什么捕获 driver 对象会在运输边界失败。

对外写入还需要独立于任务次数的业务标识。分区编号和某次任务尝试编号属于执行计划,数据重新分区后可能变化;业务去重键应来自稳定业务语义。即便某次本地测试只执行一次,真实任务重试、再次提交作业或上游重复数据仍然可能产生多次写入。测试计划应分别覆盖这些触发条件,不能用一个计数值概括所有交付语义。

结果与练习

最终进程退出码为 0,输出:

1
spark=4.0.0;first=3;second=6;serialization-rejected=true
场景 实测结果 解释
第一次 collect 总额 600,计数 3 三个元素执行转换
第二次 collect 总额 600,计数 6 未缓存结果重新计算
捕获 DriverOnly 被拒绝 对象不可序列化
进程退出 0 正常与拒绝路径断言通过

手算题:把第二次 collect 删除,只读取两次 executions.value,计数是否还会变成 6?不会;读取驱动端累计结果不是触发这条 RDD 的第二次 action。

修改练习:把 DriverOnly 的 value 先复制到局部变量,让闭包只使用局部值,验证结果为 8、9。随后增加缓存并记录两次 action 的计数,再主动解除缓存运行第三次 action。将“是否重算”的观测与“外部业务是否重复提交”分成不同断言。

可迁移规则 验收需要记录
版本由框架矩阵决定 编译器、反射库、运行 JDK
闭包需要序列化 捕获对象图与目标诊断
转换可能重执行 action、attempt 与业务操作键
本地结果有范围 未测的集群故障与提交协议

完整启动参数见实验说明。

参考资料

顺序导航:系列入口:00 · 上一篇:E05 · 下一篇:E07。