跳到内容

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

mapfilterselectjoin 等先构建计划;countcollect、写出等 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;
  • 重跑是否产生同一业务结果;
  • 输出发布是否原子,失败重试是否重复。

参考资料

Built with VitePress | Software Systems Atlas