如何在Spark中无循环实现类似SQL的Dataframe逐行更新?
Spark实现逐行更新逻辑的方案探讨
原始数据
| date | value |
|---|---|
| 2022-01-08 | 2 |
| 2022-01-09 | 4 |
| 2022-01-10 | 6 |
| 2022-01-11 | 8 |
原SQL逻辑及输出
以下SQL通过循环逐行更新数据,每一行更新后的值会被后续行的outer apply引用:
WHILE (@start_date <= @end_date) BEGIN update t1 set value = IIF(ISNULL(avg_value,0) < 2, 0,1) from #table t1 outer apply ( select top 1 value as avg_value FROM #table t2 WHERE value >= 2 AND t2.date < t1.date ORDER BY date DESC ) t3 where t1.date = @start_date SET @start_date = dateadd(day,1, @start_date) END
最终输出结果:
| date | value | avg_value |
|---|---|---|
| 2022-01-08 | 0 | null |
| 2022-01-09 | 0 | 0 |
| 2022-01-10 | 0 | 0 |
| 2022-01-11 | 0 | 0 |
Spark现有实现及差异
在Spark中尝试用Window函数计算辅助列avg_value,得到中间结果:
| date | value | avg_value |
|---|---|---|
| 2022-01-08 | 0 | null |
| 2022-01-09 | 4 | 2 |
| 2022-01-10 | 6 | 4 |
| 2022-01-11 | 8 | 6 |
再通过withColumn更新value列,最终输出:
| date | value |
|---|---|
| 2022-01-08 | 0 |
| 2022-01-09 | 1 |
| 2022-01-10 | 1 |
| 2022-01-11 | 1 |
两者差异的核心原因:Spark采用批量计算模式,先完成所有avg_value的计算再统一更新value,不会像SQL那样逐行迭代并即时引用更新后的值。
问题
能否在Spark中不使用循环实现上述SQL的逐行更新逻辑?原始数据量约30万行,出于性能考虑必须避免循环。
内容的提问来源于stack exchange,提问作者OdiumPura
相关产品推荐
相关产品推荐

