大数据(8q)流计算updateStateByKey

2023-11-11 02:40

本文主要是介绍大数据(8q)流计算updateStateByKey,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

文章目录

  • 前言
  • updateStateByKey示例
  • updateStateByKey源码
  • Option知识补充
    • getOrElse
    • isEmpty

前言

  • 本文属于Spark Streaming分支章节
  • 流式处理中,分为有状态冇状态
  • 有状态:记录之前数据流处理的信息
  • updateStateByKey有状态Transformation

updateStateByKey示例

import org.apache.spark.{SparkConf, SparkContext}
import org.apache.spark.rdd.RDD
import org.apache.spark.streaming.dstream.{DStream, InputDStream}
import org.apache.spark.streaming.{Seconds, StreamingContext}
import scala.collection.mutable// 创建SparkContext和SparkStreamingContext
val c0: SparkConf = new SparkConf().setAppName("a0").setMaster("local[2]")
val sc: SparkContext = new SparkContext(c0)
val ssc: StreamingContext = new StreamingContext(sc, Seconds(9))
// 创建RDD队列,并放入QueueInputDStream
val rddQueue: mutable.Queue[RDD[String]] = new mutable.Queue[RDD[String]]()
val iDS: InputDStream[String] = ssc.queueStream(rddQueue, oneAtATime = false)
//===========================================================================
// 数据预处理
val dS: DStream[(String, Int)] = iDS.map((_, 1))
// 无状态
dS.reduceByKey(_ + _).print()
//设置检查点路径,用于保存状态
ssc.checkpoint("checkpoint")
// 根据 Key 来更新状态
dS.updateStateByKey(// seq是一个DStream内所有RDD相同Key连成的Value队列(seq: Seq[Int], state: Option[Int]) => {Option(seq.sum + state.getOrElse(0))}
).print()
//===========================================================================
// 启动任务:循环输入文本,按空格切分
ssc.start()
while (true) {rddQueue += sc.makeRDD(scala.io.StdIn.readLine.split(" "))
}
ssc.awaitTermination()

结果打印

updateStateByKey源码

def updateStateByKey[S: ClassTag](updateFunc: (Seq[V], Option[S]) => Option[S],partitioner: Partitioner): DStream[(K, S)] = ssc.withScope {val cleanedUpdateF = sparkContext.clean(updateFunc)val newUpdateFunc = (iterator: Iterator[(K, Seq[V], Option[S])]) => {iterator.flatMap(t => cleanedUpdateF(t._2, t._3).map(s => (t._1, s)))}updateStateByKey(newUpdateFunc, partitioner, true)
}
def updateStateByKey[S: ClassTag](updateFunc: (Iterator[(K, Seq[V], Option[S])]) => Iterator[(K, S)],partitioner: Partitioner,rememberPartitioner: Boolean): DStream[(K, S)] = ssc.withScope {val cleanedFunc = ssc.sc.clean(updateFunc)val newUpdateFunc = (_: Time, it: Iterator[(K, Seq[V], Option[S])]) => {cleanedFunc(it)}new StateDStream(self, newUpdateFunc, partitioner, rememberPartitioner, None)
}

Option知识补充

  • Option译作选项,用来表示一个值是可选的(有值或无值)
  • Option[T]是一个类型为T的可选值的容器:
    若值存在,Option[T]就是一个Some[T]
    若值不存在,Option[T]就是对象None
val myMap: Map[String, String] = Map("key1" -> "value1")
val value1: Option[String] = myMap.get("key1")
val value2: Option[String] = myMap.get("key2")
println(value1)  // Some("value1")
println(value2)  // None

getOrElse

val a: Option[Int] = Some(5)
val b: Option[Int] = None
println(a.getOrElse(0))  // 5
println(b.getOrElse(0))  // 0

isEmpty

val a: Option[Int] = Some(5)
val b: Option[Int] = None
println(a.isEmpty)  // false
println(b.isEmpty)  // true

这篇关于大数据(8q)流计算updateStateByKey的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/386990

相关文章

SpringBoot集成Milvus实现数据增删改查功能

《SpringBoot集成Milvus实现数据增删改查功能》milvus支持的语言比较多,支持python,Java,Go,node等开发语言,本文主要介绍如何使用Java语言,采用springboo... 目录1、Milvus基本概念2、添加maven依赖3、配置yml文件4、创建MilvusClient

SpringValidation数据校验之约束注解与分组校验方式

《SpringValidation数据校验之约束注解与分组校验方式》本文将深入探讨SpringValidation的核心功能,帮助开发者掌握约束注解的使用技巧和分组校验的高级应用,从而构建更加健壮和可... 目录引言一、Spring Validation基础架构1.1 jsR-380标准与Spring整合1

MySQL 中查询 VARCHAR 类型 JSON 数据的问题记录

《MySQL中查询VARCHAR类型JSON数据的问题记录》在数据库设计中,有时我们会将JSON数据存储在VARCHAR或TEXT类型字段中,本文将详细介绍如何在MySQL中有效查询存储为V... 目录一、问题背景二、mysql jsON 函数2.1 常用 JSON 函数三、查询示例3.1 基本查询3.2

SpringBatch数据写入实现

《SpringBatch数据写入实现》SpringBatch通过ItemWriter接口及其丰富的实现,提供了强大的数据写入能力,本文主要介绍了SpringBatch数据写入实现,具有一定的参考价值,... 目录python引言一、ItemWriter核心概念二、数据库写入实现三、文件写入实现四、多目标写入

使用Python将JSON,XML和YAML数据写入Excel文件

《使用Python将JSON,XML和YAML数据写入Excel文件》JSON、XML和YAML作为主流结构化数据格式,因其层次化表达能力和跨平台兼容性,已成为系统间数据交换的通用载体,本文将介绍如何... 目录如何使用python写入数据到Excel工作表用Python导入jsON数据到Excel工作表用

Mysql如何将数据按照年月分组的统计

《Mysql如何将数据按照年月分组的统计》:本文主要介绍Mysql如何将数据按照年月分组的统计方式,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录mysql将数据按照年月分组的统计要的效果方案总结Mysql将数据按照年月分组的统计要的效果方案① 使用 DA

鸿蒙中Axios数据请求的封装和配置方法

《鸿蒙中Axios数据请求的封装和配置方法》:本文主要介绍鸿蒙中Axios数据请求的封装和配置方法,本文给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友参考下吧... 目录1.配置权限 应用级权限和系统级权限2.配置网络请求的代码3.下载在Entry中 下载AxIOS4.封装Htt

Python获取中国节假日数据记录入JSON文件

《Python获取中国节假日数据记录入JSON文件》项目系统内置的日历应用为了提升用户体验,特别设置了在调休日期显示“休”的UI图标功能,那么问题是这些调休数据从哪里来呢?我尝试一种更为智能的方法:P... 目录节假日数据获取存入jsON文件节假日数据读取封装完整代码项目系统内置的日历应用为了提升用户体验,

Java利用JSONPath操作JSON数据的技术指南

《Java利用JSONPath操作JSON数据的技术指南》JSONPath是一种强大的工具,用于查询和操作JSON数据,类似于SQL的语法,它为处理复杂的JSON数据结构提供了简单且高效... 目录1、简述2、什么是 jsONPath?3、Java 示例3.1 基本查询3.2 过滤查询3.3 递归搜索3.4

MySQL大表数据的分区与分库分表的实现

《MySQL大表数据的分区与分库分表的实现》数据库的分区和分库分表是两种常用的技术方案,本文主要介绍了MySQL大表数据的分区与分库分表的实现,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有... 目录1. mysql大表数据的分区1.1 什么是分区?1.2 分区的类型1.3 分区的优点1.4 分