跳到内容

9.2 分区、Shuffle 与数据倾斜:读懂分布式批处理计划

把装备日志分给二十个 worker 后,过滤很快,按 equipment_id 聚合却停在最后几个任务。大部分分区几十秒结束,一个“unknown”键所在分区跑了四十分钟。增加机器没有缩短长尾,因为所有同键记录仍要汇到一起。

分布式性能的核心不是节点数,而是数据怎样在节点间移动、聚集和失败重算。

本课目标

  • 理解 job、stage、task 与 partition 的关系;
  • 区分窄依赖和需要 shuffle 的宽依赖;
  • 根据表大小和分布选择 join 策略;
  • 从执行计划和运行指标定位倾斜、spill 与小任务。

1. 从逻辑查询到物理执行

以 Spark SQL 为例,用户写的是逻辑关系:扫描、过滤、join、聚合。优化器选择物理算子,并把执行图切成 stages;每个 stage 包含可并行运行的 tasks,task 通常处理一个 partition。

python
from pyspark.sql import functions as F

missions = spark.read.parquet("missions/")
equipment = spark.read.parquet("equipment/")

result = (
    missions
    .filter(F.col("event_date") >= "2026-01-01")
    .join(equipment, "equipment_id", "left")
    .groupBy("equipment_type")
    .agg(
        F.count("*").alias("mission_count"),
        F.avg("wear_rate").alias("avg_wear_rate"),
    )
)

result.explain(mode="formatted")

先看计划是否只读必要列、过滤是否接近扫描、join 类型与统计估计是否合理,再调配置。

2. 窄依赖与宽依赖

过滤、逐行映射等操作通常只依赖当前分区,可以流水执行。按新键聚合、去重、排序和大表 join 通常要求相同键落到同一目标分区,触发 shuffle。

shuffle 包含:

  • 上游对记录重新分桶并写出;
  • 通过网络传输分区数据;
  • 下游读取、排序或构造哈希状态;
  • 失败时重新获取或重算。

它不是“禁止使用”的操作,而是要减少无意义的数据移动。先过滤、投影和局部预聚合,能降低 shuffle 字节;如果最终语义要求按键汇总,shuffle 本身不可凭口号消失。

3. 分区数量是并行度与开销的权衡

分区过少:

  • 并行度不足;
  • 单 task 内存压力大;
  • 失败重算单位大。

分区过多:

  • 调度与序列化开销高;
  • 产生许多小文件;
  • 元数据和对象存储请求增加。

不要只按输入字节设置。观察每个 task 的输入、shuffle、持续时间和 spill,再考虑 worker 核数、内存、压缩率与算子膨胀。现代引擎的自适应执行可在运行时合并小 shuffle 分区或处理部分倾斜,但仍依赖统计和合理布局。

repartition 通常会重新分布数据并触发 shuffle;coalesce 常用于减少分区,避免完整重排,但可能造成不均衡。两者不是同义的“改文件数量”。

4. Join 先检查语义,再优化策略

性能优化不能掩盖 join fanout。先验证键唯一性、NULL 语义、预期行数和未匹配比例。

常见物理策略:

Broadcast hash join

把足够小的一侧复制到 worker,本地与大表 join,避免大表按键 shuffle。是否“足够小”取决于序列化大小、并发、worker 内存和引擎阈值;不要根据磁盘文件大小盲目强制 broadcast。

Shuffle hash/sort-merge join

两侧按 join key 重分区,再在对应分区连接。适合两边都大,但网络和磁盘成本高。

利用已有布局

若两表存储分区、bucket 或排序兼容,引擎可能减少 shuffle。需要确认执行计划确实利用了布局,而不是仅凭表定义推测。

及时收集表和列统计,能帮助优化器估计行数与选择策略。hint 是最后手段,并不保证对所有 join 类型生效。

5. 数据倾斜为何制造长尾

倾斜可能来自:

  • NULLunknown 或默认值形成超级键;
  • 少数大客户/堡垒天然占据多数记录;
  • 过滤后分布与历史统计不同;
  • join fanout 让某些键产生巨大笛卡尔组合;
  • 分区哈希不均或输入文件大小悬殊。

诊断要看 task 分布,而不仅是 stage 平均:最大/中位 task 时长、输入行、shuffle 字节、spill 和输出行。

可选修复:

  • 把业务上可分开的 NULL/unknown 单独处理;
  • 对热点键做有证明可合并的两阶段聚合;
  • 对倾斜 join 使用 salting,并正确去盐与去重;
  • 广播真正的小表;
  • 启用并验证引擎的自适应倾斜处理;
  • 修复上游 fanout,而不是给错误结果加机器。

6. 内存、Spill 与 Cache

执行引擎可把中间状态 spill 到磁盘以避免 OOM,但 spill 增加 I/O,并不表示内存无关。要同时观察执行内存、缓存、GC、磁盘空间和序列化。

缓存适合被多次复用且重算昂贵的中间结果。只用一次的数据缓存会占用资源;缓存原始大表也可能比依靠列式扫描更差。缓存后应验证是否命中,并在生命周期结束时释放。

7. 重试要求确定性与幂等输出

分布式 task 可能因 worker 故障或 speculative execution 重跑。纯转换应在相同输入下产生相同结果;在 UDF 内发送邮件、调用付款或追加外部记录,会因重试产生重复副作用。

写出策略应具备:

  • 作业运行 ID 与确定性目标分区;
  • 临时位置写入并原子/事务式发布;
  • 重试可覆盖同一批次,而不是重复追加;
  • 提交前后的可恢复状态;
  • 输入、代码、配置和输出版本记录。

“task 成功一次”不等于整张表恰好发布一次,提交协议和下游存储语义同样重要。

8. 性能排障顺序

text
1. 校验行数、唯一性和业务结果
2. 查看逻辑/物理计划与统计估计
3. 找最慢 stage 和最长 tasks
4. 区分扫描、shuffle、CPU、spill、GC、写出
5. 检查最大值而非只看平均值
6. 一次改变一个变量并重复代表性基准
7. 记录成本、稳定性和结果校验

常见误区

  • 节点越多越快:单个热点键和串行阶段不会线性扩展。
  • 尽量避免所有 shuffle:有些全局语义必然需要数据重排。
  • 分区数越多越并行:微小任务和小文件会吞掉收益。
  • 缓存一定加速:只读一次或可高效重扫的数据可能更慢。

练习

  1. 为过滤、group by、全局排序和 join 标出可能的 stage 边界。
  2. 构造一个 unknown 热点键,比较平均与最大 task 指标。
  3. 比较 broadcast join 与 shuffle join 的计划、网络字节和峰值内存。
  4. 为分区写出设计一个可安全重试的提交协议。

小结

分区提供并行单位,shuffle 负责重新建立按键关系,倾斜则把并行作业重新拖回少数长尾。读懂计划和 task 分布,才能把优化落到数据移动、状态和发布语义上。

下一课处理不会结束的数据:流处理必须同时管理事件时间、迟到状态、重放和端到端交付保证。

Built with VitePress | Software Systems Atlas