尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

SparkStreaming 之 updateStateByKey 算子详解及代码实现

SparkStreaming 之 updateStateByKey 算子详解及代码实现 摘要前面的算子都是无状态的——每个 batch 算完就丢。但累计词频“累计销售额”维护会话状态这类需求必须跨 batch 记住之前的结果这就是 updateStateByKey 的用武之地。这篇拆解它最核心的 updateFunc 签名讲清必须开 checkpoint 的原因以及两个容易踩的性能坑全量 key 遍历、状态无限增长最后说说为什么新项目更该用 mapWithState。关键词Spark Streaming, updateStateByKey, 有状态转换, updateFunc, checkpoint, mapWithState一、为什么需要 updateStateByKey先分清一件事前面讲的 map、filter、foreachRDD、transform都是无状态的——每个 batch 独立算算完结果就扔了。但很多实时需求不是这样实时累计词频要统计从程序启动到现在每个词出现多少次。实时累计销售额每个商户的累计成交额。会话状态维护记录用户当前的登录状态。这些需求的共同点是要跨 batch 记住之前的结果。updateStateByKey就是干这个的它维护一张key → 状态的表每个 batch 用新数据更新这张表表本身跨 batch 存活。二、先看懂 updateFunc 签名理解 updateStateByKey关键是把它的核心参数 updateFunc 搞明白defupdateFunc(newValues:Seq[V],runningCount:Option[S]):Option[S]三个参数newValues: Seq[V]当前这个 batch 里某个 key 出现的所有 value一个 key 在一个 batch 里可能多次出现所以是 Seq。runningCount: Option[S]这个 key 的历史累计状态。返回Option[S]新的累计状态。Option是重点runningCount为None表示这个 key 是第一次出现还没有历史状态返回None表示删除这个 key 的状态。三、完整代码累计词频importorg.apache.spark.streaming.{Seconds,StreamingContext}valsscnewStreamingContext(conf,Seconds(5))// 有状态操作必须开 checkpoint否则抛异常ssc.checkpoint(hdfs://namenode:8020/checkpoint/wordcount)vallinesssc.socketTextStream(localhost,9999)defupdateFunc(newValues:Seq[Int],runningCount:Option[Int]):Option[Int]{valnewCountrunningCount.getOrElse(0)newValues.sum Some(newCount)}valstateDStreamlines.flatMap(_.split( )).map(word(word,1)).updateStateByKey(updateFunc)stateDStream.print()ssc.start();ssc.awaitTermination()两个关键点runningCount.getOrElse(0)首次出现的 keyrunningCount是None直接用None参与加法会报错必须getOrElse(0)兜底。这是新手最容易漏的一步。ssc.checkpoint(...)必须开状态存在内存里Driver 或 Executor 一挂累积的状态就全丢了。checkpoint 会把状态定期快照到 HDFS挂了能从快照恢复。不开的话updateStateByKey直接抛requirement failed: The checkpoint directory has not been set。四、两个必须知道的性能坑updateStateByKey 用起来简单但有两个坑不注意跑久了会出问题。坑一全量 key 遍历updateStateByKey 的实现是每个 batch对所有历史出现过的 key都调用一次 updateFunc——哪怕这个 key 在当前 batch 根本没新数据。也就是说key 的数量决定了每个 batch 的计算量。随着 key 越来越多用户数、商品数在涨每个 batch 的遍历开销线性增长处理会越来越慢。坑二状态无限增长状态表只增不减长期运行的话内存和 checkpoint 会无限膨胀最终 OOM 或 checkpoint 写爆。解决思路是在状态里带上最后更新时间updateFunc 里判断超时就返回None删除该 keycaseclassState(count:Int,lastUpdateTs:Long)defupdateFunc(newValues:Seq[Int],state:Option[State]):Option[State]{valnowSystem.currentTimeMillis()valoldstate.getOrElse(State(0,now))valnewCountold.countnewValues.sum// 30 天没更新就删除状态if(now-old.lastUpdateTs30L*24*3600*1000)NoneelseSome(State(newCount,now))}五、为什么新项目更该用 mapWithStateSpark 1.6 引入了mapWithState它解决了 updateStateByKey 上面两个坑只处理有数据的 key不像 updateStateByKey 那样全量遍历mapWithState 只对当前 batch 实际出现的 key 更新状态。内置超时清理可以给状态设置超时时间超时的 key 自动被移除不用自己写时间戳判断逻辑。实测上key 越多mapWithState 的优势越明显性能比 updateStateByKey 快数倍到数十倍。所以给个明确判断新项目有状态需求直接用 mapWithStateupdateStateByKey 主要用来理解有状态流式计算的原理以及维护旧代码。这也是为什么本系列用 updateStateByKey 来讲有状态转换——它是理解这套机制的最佳入口。六、总结updateStateByKey 是有状态转换维护跨 batch 的key → 状态表解决累计词频、累计销售额这类需求。updateFunc 签名(Seq[V], Option[S]) Option[S]重点理解None的两种含义历史无状态、删除状态。必须开 checkpointgetOrElse兜底首次出现的 None这两点是新手最容易踩的。两个性能坑全量 key 遍历、状态无限增长后者靠时间戳 返回 None 清理。新项目优先 mapWithStateupdateStateByKey 用于理解原理和维护旧代码。作者大数据技术实践者博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
返回列表