5.2 Spark:Lineage、DAG、DataFrame 与数据倾斜
Spark 把多个 transformation 组织成 DAG,遇到 action 才生成作业。性能关键不是“数据都在内存”,而是减少扫描、Shuffle、序列化和不必要物化,并让失败分区可重算。
RDD 的三个核心属性
- partition:并行和数据局部性的单位;
- dependency/lineage:分区由哪些父分区计算而来;
- function:如何从父数据生成当前数据。
RDD 不可变。丢失分区可以沿 lineage 重算,而不必始终复制每个中间结果。
Lineage 很长或源数据不再可用时,checkpoint 把状态写到可靠存储并截断依赖。cache/persist 是性能优化,executor 丢失后仍可重算;checkpoint 更偏恢复边界,语义不同。
Transformation 与 Action
map、filter、select、join 等先构建计划;count、collect、写出等 action 触发执行。
重复 action 若没有 persist,可能重复计算整个 lineage。persist 也不是越多越好:缓存占内存、可能溢写磁盘,还需要驱逐和序列化。
collect() 把全部结果拉到 driver,只适合确定很小的数据。生产代码要有硬上限,不能凭开发样本判断。
窄依赖与宽依赖
窄依赖中,一个父分区只供少数子分区使用,例如 map/filter,可流水执行在同一 stage。
宽依赖中,多个父分区数据进入多个子分区,例如 groupByKey/repartition/join,需要 Shuffle,形成 stage 边界。
优先用 reduceByKey/aggregateByKey 进行 map-side combine,而不是先 groupByKey 把所有值跨网络搬运。但自定义聚合仍需满足正确的结合语义。
DataFrame 让引擎看见结构
RDD 中的任意函数对优化器是黑盒;DataFrame 提供 schema 和表达式,使优化器能够:
- predicate pushdown;
- column pruning;
- 常量折叠;
- join 重排与策略选择;
- 更紧凑的内存和代码生成路径。
UDF 可能阻断部分优化,优先使用内置表达式。解释执行计划时同时看逻辑计划、物理计划和运行时统计,不能只看源码的链式顺序。
Join 策略与倾斜
- broadcast join:把足够小的一侧广播到 executors,避免大表 shuffle;
- sort-merge join:双方按 key 重分区和排序,适合大数据;
- shuffle hash 等策略:适用性由引擎和统计决定。
广播阈值必须考虑序列化后大小和 executor 数量。热点 join key 可以拆分、单独处理或使用自适应倾斜优化,但空值或默认 key 往往是数据质量问题,不应仅靠加资源掩盖。
小文件与分区数量
太多小文件增加元数据、任务调度和打开成本;太少大分区降低并行度并增加失败重算范围。
输出分区数应根据数据量、下游读取和目标文件大小设计。coalesce(1) 会把全量输出压到单个 task,通常只适合很小结果。
验证作业而不只看“成功”
- 输入、过滤、输出记录数;
- 主键唯一性、空值率和范围约束;
- 每分区大小与最大/中位比;
- Shuffle、spill、GC、executor loss;
- 计划是否使用期望 pushdown/join;
- 重跑是否产生同一业务结果;
- 输出发布是否原子,失败重试是否重复。
参考资料
- Zaharia et al., Resilient Distributed Datasets
- Apache Spark, RDD Programming Guide
- Apache Spark, SQL Performance Tuning