WRITING / 2026.10.02

大数据专题(四):Spark 批处理,一条 SQL 如何变成分布式计算

沿着播放明细关联视频维度的执行过程理解 Stage、Task、Shuffle、Join 和 AQE,用百万条热点数据观察真实计划与运行指标。

查询没有改,为什么今天慢了

一条日报 SQL 已经运行数周:播放明细关联视频维度,再按分类统计播放次数和观看时长。普通日期能够按时完成,某次活动却明显变慢。集群整体 CPU 利用率不高,增加资源后改善也不明显。观察任务列表时,大多数任务早已结束,最后几个任务迟迟没有完成。

这类问题提醒我们,SQL 的长度与实际工作量没有简单关系。同样的逻辑表达,输入分布、统计信息、文件布局和物理计划不同,可能产生完全不同的运行过程。优化之前,先确认结果口径不变,再找到耗时发生在哪个阶段。

本篇使用 Spark 3.5.7 本地四线程执行模式,构造一百万条事实记录及一万条维度记录。它实际运行 Spark 计划和任务,但不是四台机器的集群,也没有测试跨节点网络开销。实验的目的,是把执行机制和可观测证据连接起来。

Spark SQL 从计划到 Shuffle 和任务执行的过程

SQL 描述结果,计划决定执行方式

SQL 说明需要哪些关系、过滤条件、关联和聚合。引擎首先解析语句、解析字段和类型,形成逻辑计划,再经过规则优化和物理计划选择,将计算交给实际执行过程。相同的业务逻辑可以采用不同关联方式,过滤也可能被移动到更早的位置。

因此,优化建议应落实到某个计划变化。例如“先过滤”需要确认条件能被下推、没有被表达式阻止,也没有改变关联语义;“广播小表”需要确认维度实际能够放进执行资源的内存,而不仅是表名叫维度;“增加分区”需要确认问题确实来自单任务过大,而不是调度开销。

在实验 SQL 中,我们使用 EXPLAIN FORMATTED 查看物理结构。它可以帮助识别扫描、交换、排序与关联算子,但静态计划不是最终运行事实。启用自适应执行时,运行中还可能根据观测信息调整计划,因此需要结合执行后的 SQL 事件和任务指标。Spark SQL 性能文档介绍了相关策略。

读计划时可以先回答三个问题:哪些数据被扫描,哪里需要重新分布,关联的哪一侧被复制或排序。不要一开始就被所有节点名称淹没。把这三处与业务输入对应起来,往往就能解释主要成本来自哪里。

Driver、Executor 与任务粒度

Driver 负责组织应用的执行与协调,Executor 负责运行分配到的任务并管理相关数据。一个分区通常构成一次任务处理的基本数据范围;任务并行度与分区数量、可用资源以及阶段依赖共同决定。Spark 集群概览给出了这些角色的基本定义。

本地模式将执行放在同一台机器上,降低部署复杂度,但这并不会让分区、序列化、交换和任务边界消失。因此它适合学习执行计划和验证 SQL 正确性。涉及节点丢失、网络拥塞与远程存储吞吐的问题,则需要另外设计环境。

分区太少时,可用 CPU 可能没有工作;分区太多时,单个任务很短,调度和文件管理成本上升。更重要的是,分区数量相同不代表每个分区的数据量相同。一个分区拥有绝大多数热点键,其任务仍可能主导整个阶段的完成时间。

资源设置也要区分作用对象。增加 Driver 内存不会自动解决某个执行任务处理热点数据时的压力;增加任务槽也不会把一个已经形成的大分区自动切开。排障时,应先找到具体角色和阶段,再调整对应资源。

Stage 边界为什么常与 Shuffle 有关

有些计算可以在每个分区内部独立完成,例如过滤非法观看时长、选择需要的字段、计算一个派生列。还有一些计算需要让具有相同分组键的记录相遇,例如按视频 ID 聚合或执行某些关联。这时需要重新组织数据的归属。

Shuffle 将上游任务产生的数据按规则分发给下游分区。它可能涉及序列化、缓冲、磁盘写入、网络读取和排序。引擎可以在不同情形下做局部聚合与优化,但跨分区重组依然是理解执行成本的重要边界。RDD 编程指南中的 Shuffle 说明有助于建立基础模型。

观察一条查询时,不应把所有 Shuffle 都视为可以删除的坏东西。如果最终要按分类得到全局统计,某种跨分区合并往往不可避免。更有价值的问题是:能否先减少输入列、提前过滤、局部聚合,或改变关联方式,使进入交换的数据更少。

同样,看到网络或磁盘活动也不能立即断言 Shuffle 是瓶颈。扫描慢、数据倾斜、垃圾回收和输出提交都可能占据关键路径。必须用阶段时间、任务分布和输入输出指标把猜测变成可检查的解释。

关联为什么容易放大成本

事实表通常较大,维度表可能相对较小。某些关联方式需要双方按键重新分布并排序,某些方式则可以把较小一侧复制到参与计算的任务。选择取决于数据规模、统计、关联类型、配置和资源条件。

如果维度很小,广播可以避免将庞大的事实数据仅为关联而整体重新分布。但广播本身需要收集、传输并在执行端保存数据,维度膨胀后可能失去优势。只在开发环境看到一百行维度就永久强制广播,容易在真实数据增长后产生新的问题。

关联之前应先检查语义。维度键不唯一会放大事实行数,内连接会丢弃找不到维度的事实,左连接则需要决定缺失维度如何归类。性能优化不能通过无意改变连接类型、过滤位置或空值规则来获得更快结果。

本实验使用唯一的视频维度键,且事实键都存在于维度中。两种物理策略必须生成完全相同的分类播放次数和观看毫秒数;独立 Python 参考计算器逐组核对。这个前提将比较限制在执行方式,而不是让不同口径的查询参加竞速。

热点键与数据倾斜

一百万条事实中,前九十万条都属于视频零,剩余十万条按一万个视频 ID 分布。加上尾部再次命中视频零的记录,热点键共有 900,010 条。平均每键一百条这个数字,掩盖了绝大多数记录集中在一个键上的现实。

当某种计划按视频 ID 重新分区时,热点键对应的大量记录必须被送到同一个相关分区。阶段完成需要等最慢任务结束,于是少数任务成为拖尾。观察平均任务时长很可能看不出问题,应该查看分位数、最大值和任务处理的数据量。

处理倾斜有多种方式,但不能不加条件地应用。广播较小维度可以改变关联的数据移动方式;对可分解聚合进行局部预聚合,可以减少后续记录;对热点键加盐,需要补充第二阶段归并并保持关联语义。每种方法都在改变计算结构,应证明结果仍然一致。

增加分区数量可以分散不同键,却不会自然把一个普通哈希分区中的同一键平均拆开。这也是“分区加倍但最慢任务仍很慢”的常见原因。要改变热点内部的工作分布,需要引擎支持的倾斜处理或明确的业务等价改写。

AQE 的价值与观察边界

自适应查询执行使用运行时信息调整执行计划,例如处理交换后的分区和某些关联策略。它的价值在于弥补静态估计与实际数据之间的差异,但并不意味着所有输入分布、所有算子都会自动得到最优计划。

启用 AQE 后,应查看最终采用了哪些调整,而不是只记录配置为 true。如果在同一次比较中同时更换 Join 提示、广播阈值和 AQE 状态,观察到的变化属于这套整体配置,不能全部归因于其中一个开关。

本实验的基线明确关闭 AQE 和自动广播,使用 Merge 提示形成便于观察的计划;对照明确广播小维度并开启 AQE。这样可以直观看到两类计划的差异,但不是“只改变 AQE”的单变量实验。要单独评估 AQE,需要保持其他设置、输入和运行顺序一致,再做额外对照。

教学案例采用显式设置,是为了让计划稳定、便于读者查看,并不意味着生产系统应该照搬这些开关。实际选型应从默认行为和真实工作负载出发,只有观察到具体问题后才固定某项策略。

最小实验:运行与核对

在完成前两章接入和存储实验后,执行实验包中的进阶入口:

python3 advanced.py /absolute/path/to/lab-work

脚本生成百万条事实和一万条维度,在同一个 Spark 会话中缓存输入,再运行两种关联与分类聚合。输出分别落到独立目录,结果与 Python 参考答案比较。事件日志保存 SQL 执行起止时间、任务信息与 Shuffle 指标,文章使用根 SQL 执行记录,避免把嵌套执行重复计数。

核心关联可以缩写成下面的形式;完整语句及配置以下载包为准:

SELECT /*+ BROADCAST(d) */
  d.category,
  count(*) AS plays,
  sum(f.watch_ms) AS watch_ms
FROM facts f
JOIN dimensions d ON f.video_id = d.video_id
GROUP BY d.category;

正确性检查先于性能比较。两份结果都必须得到总计一百万次播放,每个分类的观看时长完全一致。只比较总数还不够,因为分错类别之后总数也可能不变;实验因此逐分类比较两个聚合值。

本次实际观察

实验在同一台 macOS arm64 机器、JDK 17、Spark local[4] 下运行。输入已在同会话中缓存,结果包含输出 JSON 的成本;下表是一次执行记录,不是多轮统计基准。

指标 Merge 基线 广播对照
根 SQL 执行耗时 896 ms 308 ms
该执行时间范围内任务数 24 9
最大任务耗时 182 ms 101 ms
任务耗时中位数 82.5 ms 36 ms
Shuffle 写出字节 5,550,118 430
磁盘溢写字节 0 0

两种方案与参考答案完全一致。计划与指标共同表明,在这份小维度、强热点样例中,广播对照减少了关联阶段的数据交换,执行时间也更短。没有发生磁盘溢写,因此不能声称本次优化解决了溢写或内存不足问题。

这里的最大任务与中位数覆盖整次 SQL 的相关任务,不是某个单独 Shuffle 阶段的倾斜比。任务数也不是机器数。复现时应先确认相同计划和正确结果,再观察自身环境的时间;不要把表格中的毫秒数作为验收阈值。

一次对照还受到 JIT、缓存温度、运行顺序、系统负载和文件系统缓存影响。若要形成生产性能结论,应预热、交替运行、收集多轮分布,并采用真实字段宽度和数据规模。本系列保留原始事件日志,就是为了让测量口径可以被检查,而不是只留下一个“快了几倍”的标题。

一个可重复的排障顺序

首先检查输入是否变化。文件数、总字节、记录数、键分布和维度大小,比代码提交记录更接近实际工作量。没有修改 SQL,并不代表没有改变执行问题。活动热点、上游重放和维度膨胀都可能让旧计划暴露新的成本。

然后检查计划是否变化。统计信息过期、配置调整或引擎版本升级可能影响关联策略。对照慢运行和正常运行的计划,定位新增交换、排序、全表扫描或广播对象变化,而不是凭印象修改十几个参数。

再看阶段和任务。确认时间花在等待资源、读取、计算、交换、垃圾回收还是输出;查看最慢任务的数据量与失败重试。资源空闲并不一定代表引擎没有使用能力,也可能是其余任务正在等待一个无法继续拆分的热点分区。

最后做最小改动的对照。保持输入和结果契约固定,每次只验证一个可解释假设;如果必须同时调整多个配置,就明确报告组合方案,避免过度归因。性能提升之外,还要检查资源峰值、失败概率和成本,防止将延迟问题转化为内存问题。

从事件日志复核测量口径

SQL 执行可能包含嵌套的子执行,直接把每一条开始、结束事件都当成独立查询,会重复统计同一工作。本实验按根执行身份选取起止时间,再收集对应时间范围内的任务指标。日志解析器保留查询名称、任务数和字节数,使表格能够追溯到实际事件。此处没有并发运行其他查询;若同一应用同时运行多个查询,应该进一步依据执行与阶段关联,而不能仅靠时间重叠归属任务。

任务耗时还包含调度与执行环境影响,不能用最大任务时长简单乘以任务数估算总时长。多个任务并行,阶段之间又存在依赖,总耗时由关键路径决定。查看时间线与阶段关系,比单独比较几个平均值更容易解释为何整体资源看起来空闲、查询却迟迟没有完成。

优化之后仍要考虑重跑

日报第一次运行很快,重复执行却把结果追加两遍,仍然是不合格的任务。输出路径、分区、表键和提交方式应支持明确的重跑语义。使用临时结果进行校验,再以稳定键替换或合并,通常比直接对生产汇总累加更容易恢复。

同样,任务的一次失败不应要求重新生成全部历史。如果原始输入和表版本能够定位,就可以缩小到受影响日期或分区。性能设计与恢复设计在这里相遇:合理的分区范围既能减少扫描,也能限制补数的影响范围。

下一篇把输入换成持续到来的事件。我们不再等一批数据全部到齐,而是维护随事件更新的状态,并观察乱序和故障如何改变结果的可见时间。Spark 的执行计划解决“这一批如何算”,Flink 的事件时间与状态则帮助回答“持续到来的数据何时算、如何恢复”。

大数据工程:从业务事件到可信指标

  1. 大数据专题(一):从业务数据库到数据平台,一条分析链路如何形成
  2. 大数据专题(二):Kafka 数据接入,分区、消费进度与可靠重放
  3. 大数据专题(三):从文件到分析表,Parquet、数仓分层与 Iceberg
  4. 大数据专题(四):Spark 批处理,一条 SQL 如何变成分布式计算 · 当前文章
  5. 大数据专题(五):Flink 实时计算,乱序事件如何得到正确结果
  6. 大数据专题(六):让指标可信,批流对账、历史补数与故障恢复

按时间浏览