尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

SparkStreaming 之 reduceByKeyAndWindow 详解及代码实现

SparkStreaming 之 reduceByKeyAndWindow 详解及代码实现 摘要reduceByKey 统计的是当前这个 batch但最近 10 分钟的热门词最近 1 小时的 PV这类需求要看一个滚动的时间窗口。reduceByKeyAndWindow 就是干这个的。这篇讲清 windowDuration 和 slideDuration 两个参数怎么理解窗口重叠为什么会造成重复计算以及两个重载版本普通版 vs 带逆函数的增量版的差别和选择。关键词Spark Streaming, reduceByKeyAndWindow, 滑动窗口, windowDuration, slideDuration, 逆函数一、从 reduceByKey 到 reduceByKeyAndWindow先回顾一个前提reduceByKey是每个 batch 独立聚合只统计当前 2 秒batchInterval里的数据。但很多实时统计需求不是这样最近 10 分钟的热门搜索词。最近 1 小时的页面 PV。滑动窗口告警比如过去 5 分钟内错误数超过阈值就告警。这些需求要看的是一段滚动的时间范围reduceByKeyAndWindow就是在一个滑动窗口上做 reduceByKey 聚合。二、两个参数windowDuration 和 slideDuration理解窗口操作先把这两个参数搞清楚pairs.reduceByKeyAndWindow(func,windowDuration,slideDuration)windowDuration窗口长度一个窗口覆盖的时间范围比如 10 秒。一个窗口 多个 batch 的并集。slideDuration滑动间隔窗口多久滑动一次比如 2 秒。也就是多久输出一次窗口结果。举个例子batchInterval 2swindowDuration 10sslideDuration 4s那么一个窗口覆盖 5 个 batch10/2每滑一次前进 2 个 batch4/2。最容易踩的坑windowDuration 和 slideDuration 都必须是 batchInterval 的整数倍。不满足这个约束Spark 会在运行时报错。很多人配置窗口时随手写了个不是整数倍的秒数调试半天才发现是这里的问题。三、窗口重叠数据被重复计算窗口滑动的过程中相邻窗口是重叠的。回到上面的例子窗口 1 覆盖 batch0~batch4窗口 2 覆盖 batch2batch6——batch2batch4 被两个窗口共享。这意味着如果用最简单的实现每个窗口都从零重新聚合重叠部分的数据会被重复计算。重叠得越多浪费越大。这正是后面增量版要解决的问题。四、两个重载版本reduceByKeyAndWindow 有两个重载差别就在怎么处理重叠。版本一普通版重复计算valresultpairs.reduceByKeyAndWindow((a:Int,b:Int)ab,Seconds(10),Seconds(4))每个窗口从零重新聚合重叠部分重复算。实现简单但窗口大、滑动小时浪费严重。小窗口、性能不敏感的场景够用。版本二增量版带逆函数valresultpairs.reduceByKeyAndWindow((a:Int,b:Int)ab,// 累加函数加新滑入的数据(a:Int,b:Int)a-b,// 逆函数减滑出的数据Seconds(10),Seconds(4))增量版的思路是新窗口 旧窗口 新滑入的 batch − 滑出的 batch。只算新增和移除不重算重叠部分性能大幅提升。代价有两条必须提供逆函数。逆函数要能撤销累加函数的作用——加法配减法、乘法配除法。像 max/min 这种没有简单逆运算的聚合就没法用增量版只能用普通版。必须开 checkpoint。增量版要维护跨窗口的中间状态状态要能持久化、故障恢复。五、完整代码最近 10 分钟热门词importorg.apache.spark.streaming.{Seconds,Minutes,StreamingContext}valsscnewStreamingContext(conf,Seconds(2))// 增量版必须开 checkpointssc.checkpoint(hdfs://namenode:8020/checkpoint/hotword)vallinesssc.socketTextStream(localhost,9999)valhotWordslines.flatMap(_.split( )).map(word(word,1)).reduceByKeyAndWindow((a:Int,b:Int)ab,// __ 加新(a:Int,b:Int)a-b,// _-_ 减旧Minutes(10),// 10 分钟窗口Seconds(2)// 2 秒滑动)hotWords.print()ssc.start();ssc.awaitTermination()这段代码每 2 秒输出一次最近 10 分钟的词频统计。__处理滑入窗口的新 batch_-_处理滑出窗口的旧 batch窗口里的重叠部分不重复计算。六、总结reduceByKeyAndWindow 是在滑动窗口上做聚合解决最近 N 分钟这类滚动统计需求。windowDuration 是窗口长度、slideDuration 是滑动间隔两者都必须是 batchInterval 的整数倍。窗口滑动会重叠普通版会重复计算重叠数据增量版用逆函数只算增量性能更好。增量版的两个前提聚合函数有逆函数加法↔减法、乘法↔除法max/min 不行、必须开 checkpoint。窗口大、滑动小、性能敏感的场景优先增量版小窗口或没有逆函数的聚合用普通版。作者大数据技术实践者博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
返回列表