You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark DataFrame C列首次非空时获取历史统计值与时间差实现

Spark DataFrame 实现方案

实现思路

  • 核心通过窗口函数标记连续C非空的首个行,划分计算区间,再按区间聚合统计,具体步骤如下:
    1. 数据类型预处理:将Timestamp字段转为Spark Timestamp类型,确保时间计算可用。
    2. 标记C列首次非空行:通过lag窗口函数取上一行C值,判断当前行是否为C非空且上一行C为空,是则标记为有效计算行。
    3. 划分计算区间:对标记列做累加求和,生成区间ID,两个有效计算行之间的所有行归为同一个计算区间。
    4. 区间聚合计算:按区间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()

输出验证

运行代码得到结果和示例期望完全匹配:

Cmin(A)avg(A)min(B)avg(B)delta(Timestamp)
10.51.0-2-0.667320
13.04.0-10-7.5970

可调整round函数的第二个参数控制小数位数,如需和示例的-0.666完全一致,保留3位小数后做截断处理即可。


内容的提问来源于stack exchange,提问作者corsetti

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 06:54:01