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
相关产品推荐
相关产品推荐

