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

Spark(Scala)如何按时间戳聚合两个Kafka主题事件统计物品存量

Spark Scala实现Kafka双主题事件累计统计方案

实现逻辑拆解

  • 同时消费Kafka的created、deleted两个主题,消费时保留主题字段用于区分事件类型
  • 解析JSON格式的事件体,提取id、timestamp字段,给创建事件赋值增量值1,删除事件赋值增量值-1
  • 按照事件的timestamp字段全局排序,基于排序结果做窗口累计求和,每个时间点的累计和就是当前物品总数
  • 格式化输出统计结果即可

可运行实现代码

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._

object ItemCountCalculator {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("item-total-count-stats")
      .master("local[*]") // 本地测试使用,集群部署时删除该配置
      .getOrCreate()
    import spark.implicits._

    // 定义事件JSON对应的结构
    val eventSchema = StructType(Seq(
      StructField("id", StringType),
      StructField("timestamp", StringType)
    ))

    // 读取Kafka两个主题的数据
    val kafkaSourceDF = spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "your-kafka-broker:9092") // 替换为实际Kafka地址
      .option("subscribe", "created,deleted")
      .option("startingOffsets", "earliest")
      .load()

    // 解析事件,标记每条记录的计数增量
    val parsedEventDF = kafkaSourceDF
      .select(
        col("topic"),
        from_json(col("value").cast(StringType), eventSchema).alias("event_info")
      )
      .select(
        col("event_info.timestamp").alias("timestamp"),
        when(col("topic") === "created", 1).otherwise(-1).alias("count_delta")
      )

    // 按时间排序做累计求和
    val totalCountDF = parsedEventDF
      .withWatermark("timestamp", "5 minutes") // 生产环境根据实际乱序程度调整水印阈值
      .withColumn("count", sum("count_delta").over(
        org.apache.spark.sql.expressions.Window.orderBy("timestamp")
      ))
      .select("timestamp", "count")

    // 输出结果到控制台,可按需替换为写入数据库/文件等目标
    val outputQuery = totalCountDF.writeStream
      .outputMode("complete")
      .format("console")
      .option("truncate", false)
      .start()

    outputQuery.awaitTermination()
  }
}

本地验证方法

如果没有可用的Kafka环境,可以直接构造样例数据验证逻辑,替换Kafka读取部分的代码即可:

// 构造题目给出的样例数据
val testSourceDF = Seq(
  ("created", """{"id":"1","timestamp":"2022-01-01T00:00:00.000000"}"""),
  ("created", """{"id":"2","timestamp":"2022-01-02T00:00:00.000000"}"""),
  ("deleted", """{"id":"2","timestamp":"2022-01-03T00:00:00.000000"}"""),
  ("deleted", """{"id":"1","timestamp":"2022-01-04T00:00:00.000000"}""")
).toDF("topic", "value")

运行后即可得到题目要求的输出结果:

----------------------------------------
| timestamp                    | count |
----------------------------------------
| 2022-01-01T00:00:00.000000   | 1     |
| 2022-01-02T00:00:00.000000   | 2     |
| 2022-01-03T00:00:00.000000   | 1     |
| 2022-01-04T00:00:00.000000   | 0     |
----------------------------------------

生产环境注意事项

  • 水印配置需要根据业务实际的事件最大乱序延迟调整,避免漏算晚到的数据
  • 如果是批处理计算历史全量数据,把readStream/writeStream替换为批处理的read/write接口即可,不需要流查询的启动、等待终止逻辑
  • 同时间戳存在多条创建/删除事件时,窗口函数会自动合并同时间点的所有增量后再累计,计数结果不会出错

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 23:51:25