PySpark Structured Streaming窗口按时间戳取首尾记录实现方法
问题根因
F.first()和F.last()函数本身不保证按字段排序取值,分组聚合过程中shuffle分区内的记录顺序受数据到达顺序、分区分配规则影响,存在非确定性,无法保证返回分组内TimeStamp真正最大/最小的对应记录。
另外注意示例数据存在无符号整数溢出回绕(Value从65535跳变为0),不要直接用F.max(Value)/F.min(Value)取首尾值,必须以TimeStamp作为排序依据取对应Value,否则结果错误。
实现方案
以下两种方案均为Structured Streaming原生支持的确定性逻辑,不受执行顺序影响:
方案1:使用max_by/min_by函数(Spark 3.0及以上版本推荐,写法最简洁)
max_by(value_col, order_col)会返回分组内order_col取最大值时对应的value_col值,min_by逻辑反之,完全匹配按时间戳取首尾值的需求。
from pyspark.sql import functions as F Silver = (Bronze .withWatermark("TimeStamp", "1 minute") .groupBy(['sensor_id', F.window('TimeStamp', '1 minute')]) .agg( F.max_by(F.col('value'), F.col('TimeStamp')).alias('lastvalue'), F.min_by(F.col('value'), F.col('TimeStamp')).alias('firstvalue'), F.max(F.col('TimeStamp')).alias('lastTimeStamp'), F.min(F.col('TimeStamp')).alias('firstTimeStamp') ) .withColumn('value_diff', F.col('lastvalue') - F.col('firstvalue')) )
方案2:结构体聚合方案(兼容Spark 2.x版本)
Spark 2.x版本没有内置max_by/min_by函数,可以将时间戳和指标值打包为结构体——Spark对结构体的比较规则是按字段声明顺序依次比较大小,因此直接取分组内结构体的最大/最小值,即可得到对应时间戳最早/最晚的记录,再拆分结构体字段即可。
from pyspark.sql import functions as F Silver = (Bronze .withWatermark("TimeStamp", "1 minute") # 注意结构体中必须把TimeStamp放在第一个字段,保证比较时优先按时间戳排序 .withColumn("ts_value_struct", F.struct(F.col("TimeStamp"), F.col("value"))) .groupBy(['sensor_id', F.window('TimeStamp', '1 minute')]) .agg( F.max("ts_value_struct").alias("max_struct"), F.min("ts_value_struct").alias("min_struct") ) .select( "sensor_id", "window", F.col("max_struct.TimeStamp").alias("lastTimeStamp"), F.col("max_struct.value").alias("lastvalue"), F.col("min_struct.TimeStamp").alias("firstTimeStamp"), F.col("min_struct.value").alias("firstvalue") ) .withColumn('value_diff', F.col('lastvalue') - F.col('firstvalue')) )
注意:不要尝试在分组前对全量流数据执行
orderBy再聚合,Structured Streaming流查询中无聚合/窗口绑定的全局排序会直接抛出异常,即使配置触发间隔也无法保证聚合阶段的记录顺序,没有实际意义。
内容的提问来源于stack exchange,提问作者Kermit754
相关产品推荐
相关产品推荐

