SparkStreaming中的转换算子2--有状态的转换算子updateStateByKey
2022/9/2 23:23:11
本文主要是介绍SparkStreaming中的转换算子2--有状态的转换算子updateStateByKey,对大家解决编程问题具有一定的参考价值,需要的程序猿们随着小编来一起学习吧!
- 将之前批次的状态保存,
package SparkStreaming.trans import org.apache.spark.SparkConf import org.apache.spark.storage.StorageLevel import org.apache.spark.streaming.dstream.DStream import org.apache.spark.streaming.{Seconds, StreamingContext} object ByUpdateByKey { def main(args: Array[String]): Unit = { val conf = new SparkConf().setMaster("local[3]").setAppName("transform3") val ssc = new StreamingContext(conf, Seconds(3)) ssc.checkpoint("hdfs://node1:9000/sparkstreaming") val ds: DStream[String] = ssc.socketTextStream("node1", 44444, StorageLevel.MEMORY_ONLY) val ds1: DStream[(String, Int)] = ds.flatMap(_.split(" ")).map((_, 1)) /* (A, B) A:之前批次处理得到的结果 B:当前批次处理得到的结果 */ val ds2 = ds1.updateStateByKey((array: Seq[Int], state: Option[Int]) => { var num: Int = state.getOrElse(0) for (elem <- array) { num += elem } Option(num) }) ds2.print() ssc.start() ssc.awaitTermination() } }
这篇关于SparkStreaming中的转换算子2--有状态的转换算子updateStateByKey的文章就介绍到这儿,希望我们推荐的文章对大家有所帮助,也希望大家多多支持为之网!
- 2024-05-01为什么公共事业机构会偏爱 TiDB :TiDB 数据库在某省妇幼健康管理系统的应用
- 2024-04-26敏捷开发:想要快速交付就必须舍弃产品质量?
- 2024-04-26静态代码分析的这些好处,我竟然都不知道?
- 2024-04-26你在测试金字塔的哪一层?(下)
- 2024-04-26快刀斩乱麻,DevOps让代码评审也自动起来
- 2024-04-262024年最好用的10款ER图神器!
- 2024-04-2203-为啥大模型LLM还没能完全替代你?
- 2024-04-2101-大语言模型发展
- 2024-04-17基于SpringWeb MultipartFile文件上传、下载功能
- 2024-04-14个人开发者,Spring Boot 项目如何部署