一间仓库雇了许多工人,货却迟迟发不出去。问题可能不在工人干得慢,而在调度员逐箱开单、逐个派活,大家只能站着等。Spark 集群也会遇到类似情况:机器很多,不等于工作真的被并行分了出去。
Reddit 用户 al_coper 分享了一个生产案例:一条每周运行、处理过去 24 个月历史数据的 Databricks 管线,原本一次约需 126 小时。重构后,运行时间降到 30 分钟以内。不过,这些数字和诊断均来自作者自述,未附截图、配置、代码或第三方复现,下文应视为一份值得参考的案例记录,而非普遍结论。
集群在等谁派活?
Spark 是一套分布式数据处理系统。它把角色大致分成两类:Driver 是“总调度室”,负责生成任务和协调执行;Executor 则是分布在集群各台机器上的“工人”,真正并行处理数据。
原管线使用 Python for 循环、ThreadPoolExecutor 和 replaceWhere。作者称,Spark 几乎没有把工作分布到集群中。换句话说,表面上已经用了并发工具和一批计算资源,实际流程仍由 Driver 逐步组织和发起。
团队通过 Spark UI 定位问题。Spark UI 是作业的监控界面,可以查看任务耗时、并行度、数据读写和资源使用。作者据此判断,几乎所有编排工作都发生在 Driver,主要瓶颈不在 Executor 的计算速度,而在任务如何被组织和下发。
这里需要谨慎区分“同时出现”和“就是根因”。现有材料不足以证明 ThreadPoolExecutor 或 replaceWhere 本身必然导致性能问题。更稳妥的理解是:在这套具体实现里,它们与 Python 循环共同构成了一种偏向 Driver 端的执行方式。
不再逐月发单
重构并不是只删掉一段 Driver 代码。作者同时调整了三处执行机制。
第一处是 Dynamic partition overwrite,即动态分区覆盖。分区可以理解为按月份、日期等字段把大表拆成若干块。动态覆盖只替换本次实际写到的分区,因此可以避免在 Python 中按月循环、反复发起写入。
第二处是 Adaptive Query Execution(AQE,自适应查询执行)。它允许 Spark 在查询运行期间根据真实数据量调整计划,例如合并过小的分区,或改变表连接策略。集群不必完全依赖运行前的估计。
第三处是改用 Spark 原生的分布式执行,让任务更多地交给集群统一拆分和并行处理,而不是由 Driver 承担细碎编排。
作者称,重构后管线从约 126 小时降至 30 分钟以内,并且业务团队验证了结果的正确性。按上限 30 分钟计算,运行时间至少缩短到原来的约 1/252。
真正值得看的是诊断顺序
这个案例最有价值的地方,不只是一个醒目的提速数字,而是它提醒我们:管线慢时,先别急着增加机器。要先确认机器是否真的拿到了足够的工作。
Spark UI 在这里承担了类似交通监控的作用。它帮助团队区分两类问题:是 Executor 正在满负荷计算,还是 Driver 忙于循环、提交和协调,导致集群空转。两者表面上都是“作业很慢”,处理办法却完全不同。
案例也说明,分布式系统的关键不只是拥有多少计算资源,还包括任务以什么粒度、通过什么路径送到这些资源手中。若中央调度成为单点瓶颈,增加更多 Executor,未必能换来相应速度。
局限与未知
- 帖子没有提供 Spark UI 截图、数据规模、集群规格、代码、测试口径和各项改动的单独效果,外部无法核验或复现。
- 提速发生在同时引入动态分区覆盖、AQE 和原生分布式执行之后,不能把全部收益单独归因于“移除 Driver 端编排”。
- 业务正确性由作者称已获业务团队验证,但材料没有披露验证范围和方法。这个案例适合作为排查思路,不足以推出 Databricks 或 Spark 管线的一般性能规律。