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()
关键说明
- 有状态聚合:通过
flatMapGroupsWithState维护每个group的全局累积状态,确保从流启动开始的所有数据都被累加。 - 乱序处理:对每个批次的数据按时间戳排序,若实际流存在延迟数据,可添加
.withWatermark("timestamp", "1 hour")(根据业务调整延迟阈值)来处理。 - 按日输出:在得到全局累积和的数据流后,通过日窗口分组+
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.
相关产品推荐
相关产品推荐

