Spark 2.1.1中Dataset循环Union操作性能缓慢问题咨询
其实你提到的“惰性操作”没错,但逻辑执行计划的构建与优化本身是即时的,这就是你看到耗时的核心原因。让我拆解一下具体问题:
核心原因
1. 大量Union导致逻辑计划过度膨胀
你的代码会生成720个左右的子Dataset(30天×24小时),然后通过reduce(_.union(_))把它们逐个合并。在Spark 2.1.1中,每次Union操作都会创建一个新的UnionExec节点,并且把两个子计划合并成一个更大的计划。随着Union次数增加,这个逻辑计划会变成一个深度极高的链式结构,Spark的优化器需要不断遍历、合并这些计划节点,这个过程的时间开销会随着Union次数线性增长。
2. 老版本Spark缺乏Union合并优化
Spark 2.3及以后的版本加入了Union合并优化规则(CombineUnions),它会自动把连续的多个Union操作合并成一个包含所有子节点的单一Union计划,大幅减少计划的复杂度。但Spark 2.1.1没有这个优化,所以每一次Union都会被单独处理,计划构建的成本会越来越高。
3. 重复的过滤条件解析开销
每个getDsForOneHour都会生成一组独立的过滤条件(year === x and month === y...),这些条件都会被加入到各自的子计划中。当合并所有计划时,Spark需要逐一解析、验证这些过滤表达式,720次重复操作累积起来也会增加不少耗时。
优化建议
1. 替换循环Union为单次范围过滤
这是最有效的优化方式:不需要拆分每个小时的过滤再Union,直接生成一个覆盖整个时间范围的过滤条件。比如:
// 如果有timestamp列,直接用时间范围过滤最简洁 val myDs = inputDs.where( col("timestamp").between(fromDate.toInstant.toEpochMilli, toDate.toInstant.toEpochMilli) ) // 如果只有年/月/日/小时列,构造组合范围条件 val startYear = fromDate.getYear val startMonth = fromDate.getMonthValue val startDay = fromDate.getDayOfMonth val startHour = fromDate.getHour val endYear = toDate.getYear val endMonth = toDate.getMonthValue val endDay = toDate.getDayOfMonth val endHour = toDate.getHour val myDs = inputDs.where( (col("year") > startYear) || (col("year") === startYear && col("month") > startMonth) || (col("year") === startYear && col("month") === startMonth && col("day") > startDay) || (col("year") === startYear && col("month") === startMonth && col("day") === startDay && col("hour") >= startHour) || (col("year") < endYear) || (col("year") === endYear && col("month") < endMonth) || (col("year") === endYear && col("month") === endMonth && col("day") < endDay) || (col("year") === endYear && col("month") === endMonth && col("day") === endDay && col("hour") < endHour) )
这种方式只生成一个过滤计划,完全避免了多次Union的开销。
2. 批量Union减少计划深度
如果必须按小时拆分处理(比如每个小时有额外的业务逻辑),可以把多个子Dataset分成批次,每批次先Union,再合并批次结果。比如每10个小时合并一次:
val hourIter = Iterator.iterate(fromDate)(_.plus(ofHours(1))).takeWhile(_.isBefore(toDate)) val batchSize = 10 val batchedDs = hourIter.grouped(batchSize).map { batch => batch.map(next => getDsForOneHour(inputDs, next.getYear, next.getMonthValue, next.getDayOfMonth, next.getHour)) .reduce(_.union(_)) }.reduce(_.union(_))
这样Union的总层级从720降到了约72,能显著降低计划构建的复杂度。
3. 升级Spark版本
如果业务允许,升级到Spark 2.3+版本,原生的CombineUnions优化会自动帮你合并多次Union操作,不需要修改代码就能获得性能提升。
内容的提问来源于stack exchange,提问作者Marious

