Apache Spark中last函数结合Window orderBy未按预期工作,是否为BUG?
这是Spark的BUG吗?last函数未返回分区内最后一个值
环境:Apache Spark 3.3.2,Google Dataproc 2.1
先看测试代码及结果:
val orderDF = Seq(("1234","created","2023-09-17 00:00:00"),("1234","updated","2023-09-17 01:00:00")).toDF("orderID","status","eventTimestamp") orderDF.show(false) // 输出: +-------+-------+-------------------+ |orderID|status |eventTimestamp | +-------+-------+-------------------+ |1234 |created|2023-09-17 00:00:00| |1234 |updated|2023-09-17 01:00:00| +-------+-------+-------------------+ orderDF.select($"orderID",$"status",$"eventTimestamp",last($"status").over(Window.partitionBy($"orderID").orderBy($"eventTimestamp")).alias("latestStatus")).show(false) // 输出: +-------+-------+-------------------+------------+ |orderID|status |eventTimestamp |latestStatus| +-------+-------+-------------------+------------+ |1234 |created|2023-09-17 00:00:00|created | |1234 |updated|2023-09-17 01:00:00|updated | +-------+-------+-------------------+------------+ orderDF.select($"orderID",$"status",$"eventTimestamp",first($"status").over(Window.partitionBy($"orderID").orderBy($"eventTimestamp".desc)).alias("latestStatus")).show(false) // 输出: +-------+-------+-------------------+------------+ |orderID|status |eventTimestamp |latestStatus| +-------+-------+-------------------+------------+ |1234 |updated|2023-09-17 01:00:00|updated | |1234 |created|2023-09-17 00:00:00|updated | +-------+-------+-------------------+------------+
这不是BUG,是Spark窗口函数的默认行为导致的:
- 当只给
last()指定orderBy但未定义窗口范围(frame)时,Spark默认的窗口范围是RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW——也就是从分区第一行到当前行。所以第一行的last()只能取到当前行的created,第二行取到当前行的updated,结果自然是每行对应自身状态。 - 而
first()配合倒序排序时,默认窗口范围同样是到当前行,但数据按时间倒序排列后,第一行就是最新的updated,后续行的窗口范围包含前面的行(即原数据的最新行),所以first()会一直取到这个最新值。
如果要让last()返回分区内真正的最后一个值,需要显式指定窗口范围覆盖整个分区:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ val fullPartitionWindow = Window.partitionBy($"orderID") .orderBy($"eventTimestamp") .rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) orderDF.select( $"orderID", $"status", $"eventTimestamp", last($"status").over(fullPartitionWindow).alias("latestStatus") ).show(false)
运行后两行的latestStatus都会是updated,符合预期。
内容的提问来源于stack exchange,提问作者Vikas
相关产品推荐
相关产品推荐

