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

Spark流处理中如何从流起始点按组计算日维度累积和?

Spark流中实现从起始点按Group累积和并按日输出的解决方案

需求说明

按group对value列计算从时间序列起始点开始的累积和,并按日输出每日的最终累积值。

批处理场景实现(参考)

批处理下可通过窗口函数实现全局累积,再按日分组取当日最后累积值:

import org.apache.spark.sql._
import org.apache.spark.sql.functions._
import java.time.Instant

val columns = Seq("timestamp", "group", "value")
val data = List(
  (Instant.parse("2020-01-01T00:00:00Z"), "Group1",  0),
  (Instant.parse("2020-01-01T00:00:00Z"), "Group2",  0),
  (Instant.parse("2020-01-01T12:00:00Z"), "Group1",  1),
  (Instant.parse("2020-01-01T12:00:00Z"), "Group2", -1),
  (Instant.parse("2020-01-02T00:00:00Z"), "Group1",  2),
  (Instant.parse("2020-01-02T00:00:00Z"), "Group2", -2),
  (Instant.parse("2020-01-02T12:00:00Z"), "Group1",  3),
  (Instant.parse("2020-01-02T12:00:00Z"), "Group2", -3),
)

val df = spark
  .createDataFrame(data)
  .toDF(columns: _*)

// 定义按group分组、从起始到当前行的窗口
val event_window = Window
  .partitionBy(col("group"))
  .orderBy(col("timestamp"))
  .rowsBetween(Window.unboundedPreceding, Window.currentRow)

val computed_df = df
  .withColumn("cumsum", sum('value).over(event_window))
  .groupBy(window($"timestamp", "1 day"), $"group")
  .agg(last("cumsum").as("cumsum_by_day"))

computed_df.show(truncate = false)

批处理输出结果:

+------------------------------------------+------+-------------+
|window                                    |group |cumsum_by_day|
+------------------------------------------+------+-------------+
|{2020-01-01 01:00:00, 2020-01-02 01:00:00}|Group1| 1           |
|{2020-01-02 01:00:00, 2020-01-03 01:00:00}|Group1| 6           |
|{2020-01-01 01:00:00, 2020-01-02 01:00:00}|Group2|-1           |
|{2020-01-02 01:00:00, 2020-01-03 01:00:00}|Group2|-6           |
+------------------------------------------+------+-------------+

流处理遇到的问题

直接复用批处理的窗口函数在Spark流中会运行失败,普通的按日窗口聚合只能计算当日的value总和,无法得到从流起始点开始的全局累积和。

流处理解决方案

Spark流中需使用有状态聚合维护每个group的全局累积状态,再按日提取最终值:

import org.apache.spark.sql._
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming._
import java.time.Instant

val columns = Seq("timestamp", "group", "value")
val data = List(
  (Instant.parse("2020-01-01T00:00:00Z"), "Group1",  0),
  (Instant.parse("2020-01-01T00:00:00Z"), "Group2",  0),
  (Instant.parse("2020-01-01T12:00:00Z"), "Group1",  1),
  (Instant.parse("2020-01-01T12:00:00Z"), "Group2", -1),
  (Instant.parse("2020-01-02T00:00:00Z"), "Group1",  2),
  (Instant.parse("2020-01-02T00:00:00Z"), "Group2", -2),
  (Instant.parse("2020-01-02T12:00:00Z"), "Group1",  3),
  (Instant.parse("2020-01-02T12:00:00Z"), "Group2", -3),
)

// 模拟流数据源
implicit val sqlCtx: SQLContext = spark.sqlContext
val memoryStream = MemoryStream[(Instant, String, Int)]
memoryStream.addData(data)
val df = memoryStream.toDF().toDF(columns: _*)

// 定义状态类:存储每个group的当前累积和
case class CumulativeState(sum: Int)

// 状态更新函数:维护全局累积和,并输出每条数据的累积结果
val updateStateFunc = (group: String, values: Iterator[(Instant, Int)], state: GroupState[CumulativeState]) => {
  // 获取当前状态的累积和,初始为0
  var currentSum = state.getOption.map(_.sum).getOrElse(0)
  // 按时间戳排序(处理流中可能的乱序)
  val sortedValues = values.toList.sortBy(_._1)
  
  // 遍历数据更新累积和,并输出带时间戳的累积值
  sortedValues.map { case (ts, value) =>
    currentSum += value
    (ts, group, currentSum)
  }.iterator
}

// 应用有状态聚合,得到每条数据的全局累积和
val cumulativeDf = df
  .select($"group", struct($"timestamp", $"value").as("data"))
  .groupByKey(row => row.getString(0))
  .flatMapGroupsWithState(
    outputMode = OutputMode.Append(),
    timeoutConf = GroupStateTimeout.NoTimeout()
  )(updateStateFunc)
  .toDF("timestamp", "group", "cumsum")

// 按日窗口分组,取当日最后一条累积值作为当日结果
val resultDf = cumulativeDf
  .groupBy(window($"timestamp", "1 day"), $"group")
  .agg(last("cumsum").as("cumsum_by_day"))

// 启动流查询
resultDf.writeStream
  .option("truncate", false)
  .format("console")
  .outputMode("complete")
  .start()
  .processAllAvailable()

关键说明

  1. 有状态聚合:通过flatMapGroupsWithState维护每个group的全局累积状态,确保从流启动开始的所有数据都被累加。
  2. 乱序处理:对每个批次的数据按时间戳排序,若实际流存在延迟数据,可添加.withWatermark("timestamp", "1 hour")(根据业务调整延迟阈值)来处理。
  3. 按日输出:在得到全局累积和的数据流后,通过日窗口分组+last函数提取当日最终累积值,和批处理结果完全一致。

流处理输出结果

+------------------------------------------+------+-------------+
|window                                    |group |cumsum_by_day|
+------------------------------------------+------+-------------+
|{2020-01-01 01:00:00, 2020-01-02 01:00:00}|Group1| 1           |
|{2020-01-02 01:00:00, 2020-01-03 01:00:00}|Group1| 6           |
|{2020-01-01 01:00:00, 2020-01-02 01:00:00}|Group2|-1           |
|{2020-01-02 01:00:00, 2020-01-03 01:00:00}|Group2|-6           |
+------------------------------------------+------+-------------+

内容的提问来源于stack exchange,提问作者Pascal H.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 14:50:44