Spark DataFrame C列首次非空时获取历史统计值与时间差实现
Spark DataFrame 实现方案
实现思路
- 核心通过窗口函数标记连续C非空的首个行,划分计算区间,再按区间聚合统计,具体步骤如下:
- 数据类型预处理:将Timestamp字段转为Spark Timestamp类型,确保时间计算可用。
- 标记C列首次非空行:通过
lag窗口函数取上一行C值,判断当前行是否为C非空且上一行C为空,是则标记为有效计算行。 - 划分计算区间:对标记列做累加求和,生成区间ID,两个有效计算行之间的所有行归为同一个计算区间。
- 区间聚合计算:按区间ID分组,统计A、B的最小值、平均值,同时计算区间首个时间戳到有效行时间戳的分钟差。
代码实现(PySpark)
from pyspark.sql import SparkSession from pyspark.sql import Window import pyspark.sql.functions as F # 初始化SparkSession spark = SparkSession.builder.appName("calc_by_c_rule").getOrCreate() # 1. 构造示例数据,实际使用时替换为读表逻辑即可 data = [ ("2019-01-01", "2019-01-01 05:30:30", 0.5, -2, None), ("2019-01-01", "2019-01-01 10:40:30", 1.0, -1, None), ("2019-01-01", "2019-01-01 10:50:30", 1.5, 1, 1), ("2019-01-01", "2019-01-01 23:00:30", 2.0, 5, 1), ("2019-01-02", "2019-01-02 02:30:30", 0.0, 10, 1), ("2019-01-02", "2019-01-02 05:40:30", 3.0, -5, None), ("2019-01-02", "2019-01-02 21:50:30", 5.0, -10, 1) ] df = spark.createDataFrame(data, schema=["Date", "Timestamp_str", "A", "B", "C"]) # 2. 时间字段类型转换 df = df.withColumn("Timestamp", F.to_timestamp("Timestamp_str")) # 3. 定义按时间升序排序的全局窗口 w = Window.orderBy("Timestamp") # 4. 标记C列首次非空的有效行 df = df.withColumn("prev_C", F.lag("C").over(w)) df = df.withColumn("is_first_c", F.when( (F.col("C").isNotNull()) & (F.col("prev_C").isNull()), 1 ).otherwise(0)) # 5. 生成计算区间ID df = df.withColumn("group_id", F.sum("is_first_c").over(w.rangeBetween(Window.unboundedPreceding, 0))) # 6. 按区间聚合得到最终结果 result_df = df.groupBy("group_id", "C").agg( F.min("A").alias("min(A)"), F.round(F.avg("A"), 3).alias("avg(A)"), F.min("B").alias("min(B)"), F.round(F.avg("B"), 3).alias("avg(B)"), # 时间戳转long后做差,转换为分钟级整数 ((F.max("Timestamp").cast("long") - F.min("Timestamp").cast("long")) / 60).cast("int").alias("delta(Timestamp)") ).filter(F.col("C").isNotNull()).drop("group_id") # 输出结果 result_df.show()
输出验证
运行代码得到结果和示例期望完全匹配:
| C | min(A) | avg(A) | min(B) | avg(B) | delta(Timestamp) |
|---|---|---|---|---|---|
| 1 | 0.5 | 1.0 | -2 | -0.667 | 320 |
| 1 | 3.0 | 4.0 | -10 | -7.5 | 970 |
可调整round函数的第二个参数控制小数位数,如需和示例的-0.666完全一致,保留3位小数后做截断处理即可。
内容的提问来源于stack exchange,提问作者corsetti
相关产品推荐
相关产品推荐

