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

Spark DataFrame/SQL:计算每行timestamp前3小时的value累加和

问题描述

现有一张包含timestamp(时间戳)、value(数值)、id(标识)字段的Spark DataFrame/数据表,样本数据如下:

+--------------------+-----+---+
|           timestamp|value| id|
+--------------------+-----+---+
|2024-10-05 20:38:...|   67|  0|
|2024-10-05 19:38:...|   14|  1|
|2024-10-05 18:38:...|   80|  2|
|2024-10-05 17:38:...|    6|  3|
+--------------------+-----+---+

需求:对每一行数据,计算该行timestamp往前3小时内(包含本行)所有value的数值之和,数据量可达数亿行,需要Spark DataFrame或标准SQL的实现方案。期望输出示例:

+--------------------+-----+---+
|           timestamp|sum  | id|
+--------------------+-----+---+
|2024-10-05 20:38:...|  167|  0| /* 本行value加上前3小时内所有数据的value之和 */
|2024-10-05 19:38:...|  100|  1| /* 完整数据集下结果会大于100,以此类推 */
|2024-10-05 18:38:...|   86|  2|
+--------------------+-----+---+
解决方案

针对亿级数据量的场景,时间范围滑动窗口是最优方案,可避免全表关联的性能损耗。以下是两种实现方式:

1. Spark DataFrame API 实现

Scala版本

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

// 假设你的DataFrame名为df
val resultDF = df
  .orderBy("timestamp") // 按时间戳排序,窗口计算依赖有序数据
  .withColumn(
    "sum",
    sum("value").over(
      Window
        .orderBy(col("timestamp").cast("long")) // 转成long类型计算时间差
        .rangeBetween(-3 * 60 * 60, 0) // 范围:当前时间往前3小时(秒数)到当前行
    )
  )
  .select("timestamp", "sum", "id") // 调整字段顺序

Python版本

from pyspark.sql import Window
from pyspark.sql.functions import sum, col

result_df = df \
    .orderBy("timestamp") \
    .withColumn(
        "sum",
        sum("value").over(
            Window
                .orderBy(col("timestamp").cast("long"))
                .rangeBetween(-3 * 60 * 60, 0)
        )
    ) \
    .select("timestamp", "sum", "id")

2. 标准SQL 实现

先将DataFrame注册为临时视图,再执行SQL计算:

-- 注册临时视图(Spark环境下)
CREATE OR REPLACE TEMP VIEW data_table AS SELECT * FROM df;

-- 计算3小时滑动窗口的value和
SELECT
  timestamp,
  SUM(value) OVER (
    ORDER BY UNIX_TIMESTAMP(timestamp)
    RANGE BETWEEN 3*3600 PRECEDING AND CURRENT ROW
  ) AS sum,
  id
FROM data_table
ORDER BY timestamp DESC;

性能优化提示

  • 分区策略:如果数据按时间分区(如按天),先过滤目标时间范围减少计算量;可按timestamp的小时段分区,提升窗口计算并行度。
  • 数据排序:确保数据在窗口计算前已按timestamp有序,避免重复排序开销。
  • 资源配置:针对亿级数据,调整Spark的executor内存、核心数等参数,避免内存溢出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 13:16:07