spark streaming中的广播变量应用

2024-06-16 19:58

本文主要是介绍spark streaming中的广播变量应用,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

1. 广播变量

我们知道spark 的广播变量允许缓存一个只读的变量在每台机器上面,而不是每个任务保存一份拷贝。常见于spark在一些全局统计的场景中应用。通过广播变量,能够以一种更有效率的方式将一个大数据量输入集合的副本分配给每个节点。Spark也尝试着利用有效的广播算法去分配广播变量,以减少通信的成本。 
一个广播变量可以通过调用SparkContext.broadcast(v)方法从一个初始变量v中创建。广播变量是v的一个包装变量,它的值可以通过value方法访问,下面的代码说明了这个过程:

scala> val broadcastVar = sc.broadcast(Array(1, 2, 3))
broadcastVar: org.apache.spark.broadcast.Broadcast[Array[Int]] = Broadcast(0)scala> broadcastVar.value
res0: Array[Int] = Array(1, 2, 3)

2. Spark Streaming 广播变量的更新

广播变量的声明很简单,调用broadcast就能搞定,并且scala中一切可序列化的对象都是可以进行广播的,这就给了我们很大的想象空间,可以利用广播变量将一些经常访问的大变量进行广播,而不是每个任务保存一份,这样可以减少资源上的浪费。

但是,现在项目中遇到一种这样的需求,用spark streaming 通过一些离线全局更新好的数据对用户进行实时推荐(当然这里基于一些spark streaming的内部机制,不能实现真正的时效性):(1)日志流通过kafka获取 (2) 解析日志流数据,融合离线的全局数据,对每个Dtream进行计算(3)计算结果最后发送到redis中。

其中就会涉及这样的问题:(1)离线全局的数据是需要全局获取的,不能局部进行计算 (2)这部分数据是离线定期更新的,而spark streaming一旦开始,就长时间运行。如果离线数据更新了,如何在开始的流计算中,获取到这部分更新后的数据。

针对上述问题,我们可以直接想的一种方法是,在driver端开启一个附属线程,周期性去获取离线的全局数据,然后通过diver分发到各个task中。但是考虑到这种方式:spark streaming整体的性能开销会很大,并且重新开启的后台线程的不易管理。结合spark中的广播变量,我们采用另一种方式来解决以上问题: 
1> spark中的广播变量是只读的,通过unpersist函数,可以内存中的相关序列化对象 
2> 通过Dstream的foreachRDD方法,做到定时更新 (官网上有说明,该方法是在driver端执行的)


import java.io.{ObjectInputStream, ObjectOutputStream}
import com.bf.dt.wireless.config.WirelessConfig
import com.bf.dt.wireless.formator.WirelessFormator
import com.bf.dt.wireless.storage.MysqlConnectionPool
import com.bf.dt.wireless.utils.DateUtils
import kafka.serializer.StringDecoder
import org.apache.spark.{SparkConf, SparkContext}
import org.apache.spark.broadcast.Broadcast
import org.apache.spark.streaming.kafka.KafkaUtils
import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.json4s._
import org.slf4j.LoggerFactory
import scala.collection.mutableobject WirelessLogAnalysis {object BroadcastWrapper {@volatile private var instance: Broadcast[Map[String, List[String]]] = nullprivate val map = mutable.LinkedHashMap[String, List[String]]()def getMysql(): Map[String, List[String]] = {//1.获取mysql连接池的一个连接val conn = MysqlConnectionPool.getConnection.get//2.查询新的数据val sql = "select aid_type,aids from cf_similarity"val ps = conn.prepareStatement(sql)val rs = ps.executeQuery()while (rs.next()) {val aid = rs.getString("aid_type")val aids = rs.getString("aids").split(",").toListmap += (aid -> aids)}//3.连接池回收连接MysqlConnectionPool.closeConnection(conn)map.toMap}def update(sc: SparkContext, blocking: Boolean = false): Unit = {if (instance != null)instance.unpersist(blocking)instance = sc.broadcast(getMysql())}def getInstance(sc: SparkContext): Broadcast[Map[String, List[String]]] = {if (instance == null) {synchronized {if (instance == null) {instance = sc.broadcast(getMysql)}}}instance}private def writeObject(out: ObjectOutputStream): Unit = {out.writeObject(instance)}private def readObject(in: ObjectInputStream): Unit = {instance = in.readObject().asInstanceOf[Broadcast[Map[String, List[String]]]]}}def main(args: Array[String]): Unit = {val logger = LoggerFactory.getLogger(this.getClass)val conf = new SparkConf().setAppName("wirelessLogAnalysis")val ssc = new StreamingContext(conf, Seconds(10))val kafkaConfig: Map[String, String] = Map("metadata.broker.list" -> WirelessConfig.getConf.get.getString("wireless.metadata.broker.list"),"group.id" -> WirelessConfig.getConf.get.getString("wireless.group.id"),"zookeeper.connect" -> WirelessConfig.getConf.get.getString("wireless.zookeeper.connect"),"auto.offset.reset" -> WirelessConfig.getConf.get.getString("wireless.auto.offset.reset"))val androidvvTopic = WirelessConfig.getConf.get.getString("wireless.topic1")val iphonevvToplic = WirelessConfig.getConf.get.getString("wireless.topic2")val kafkaDStream = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](ssc,kafkaConfig,Set(androidvvTopic, iphonevvToplic))//原始日志流打印kafkaDStream.print()val jsonDstream = kafkaDStream.map(x =>//解析日志流WirelessFormator.format(x._2))//解密的日志流打印jsonDstream.print()jsonDstream.foreachRDD {rdd => {// driver端运行,涉及操作:广播变量的初始化和更新// 可以自定义更新时间if ((DateUtils.getNowTime().split(" ")(1) >= "08:00:00") && (DateUtils.getNowTime().split(" ")(1) <= "10:10:00")) {BroadcastWrapper.update(rdd.sparkContext, true)println("广播变量更新成功: " + DateUtils.getNowTime())}//worker端运行,涉及操作:Dstream数据的处理和Redis更新rdd.foreachPartition {partitionRecords =>//1.获取redis连接,保证每个partition建立一次连接,避免每个记录建立/关闭连接的性能消耗partitionRecords.foreach(record => {//2.处理日志流val uid = record._1val aid_type = record._2 + "_" + record._3if (cf.value.keySet.contains(aid_type)) {(uid, cf.value.get(aid_type))println((uid, cf.value.get(aid_type)))}else(uid, "-1")}//3.redis更新数据)//4.关闭redis连接}}}ssc.start()ssc.awaitTermination()}
}

说明:以上是无线推荐项目中部分代码,其中离线全局数据存储在mysql中,MysqlConnectionPool是mysql连接池定义类,WirelessFormator是日志解密的定义类

这篇关于spark streaming中的广播变量应用的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



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

相关文章

Python之变量命名规则详解

《Python之变量命名规则详解》Python变量命名需遵守语法规范(字母开头、不使用关键字),遵循三要(自解释、明确功能)和三不要(避免缩写、语法错误、滥用下划线)原则,确保代码易读易维护... 目录1. 硬性规则2. “三要” 原则2.1. 要体现变量的 “实际作用”,拒绝 “无意义命名”2.2. 要让

利用Python操作Word文档页码的实际应用

《利用Python操作Word文档页码的实际应用》在撰写长篇文档时,经常需要将文档分成多个节,每个节都需要单独的页码,下面:本文主要介绍利用Python操作Word文档页码的相关资料,文中通过代码... 目录需求:文档详情:要求:该程序的功能是:总结需求:一次性处理24个文档的页码。文档详情:1、每个

Java中的分布式系统开发基于 Zookeeper 与 Dubbo 的应用案例解析

《Java中的分布式系统开发基于Zookeeper与Dubbo的应用案例解析》本文将通过实际案例,带你走进基于Zookeeper与Dubbo的分布式系统开发,本文通过实例代码给大家介绍的非常详... 目录Java 中的分布式系统开发基于 Zookeeper 与 Dubbo 的应用案例一、分布式系统中的挑战二

Java 缓存框架 Caffeine 应用场景解析

《Java缓存框架Caffeine应用场景解析》文章介绍Caffeine作为高性能Java本地缓存框架,基于W-TinyLFU算法,支持异步加载、灵活过期策略、内存安全机制及统计监控,重点解析其... 目录一、Caffeine 简介1. 框架概述1.1 Caffeine的核心优势二、Caffeine 基础2

使用Node.js和PostgreSQL构建数据库应用

《使用Node.js和PostgreSQL构建数据库应用》PostgreSQL是一个功能强大的开源关系型数据库,而Node.js是构建高效网络应用的理想平台,结合这两个技术,我们可以创建出色的数据驱动... 目录初始化项目与安装依赖建立数据库连接执行CRUD操作查询数据插入数据更新数据删除数据完整示例与最佳

SpringBoot中@Value注入静态变量方式

《SpringBoot中@Value注入静态变量方式》SpringBoot中静态变量无法直接用@Value注入,需通过setter方法,@Value(${})从属性文件获取值,@Value(#{})用... 目录项目场景解决方案注解说明1、@Value("${}")使用示例2、@Value("#{}"php

PHP应用中处理限流和API节流的最佳实践

《PHP应用中处理限流和API节流的最佳实践》限流和API节流对于确保Web应用程序的可靠性、安全性和可扩展性至关重要,本文将详细介绍PHP应用中处理限流和API节流的最佳实践,下面就来和小编一起学习... 目录限流的重要性在 php 中实施限流的最佳实践使用集中式存储进行状态管理(如 Redis)采用滑动

深入浅出Spring中的@Autowired自动注入的工作原理及实践应用

《深入浅出Spring中的@Autowired自动注入的工作原理及实践应用》在Spring框架的学习旅程中,@Autowired无疑是一个高频出现却又让初学者头疼的注解,它看似简单,却蕴含着Sprin... 目录深入浅出Spring中的@Autowired:自动注入的奥秘什么是依赖注入?@Autowired

GO语言短变量声明的实现示例

《GO语言短变量声明的实现示例》在Go语言中,短变量声明是一种简洁的变量声明方式,使用:=运算符,可以自动推断变量类型,下面就来具体介绍一下如何使用,感兴趣的可以了解一下... 目录基本语法功能特点与var的区别适用场景注意事项基本语法variableName := value功能特点1、自动类型推

PostgreSQL简介及实战应用

《PostgreSQL简介及实战应用》PostgreSQL是一种功能强大的开源关系型数据库管理系统,以其稳定性、高性能、扩展性和复杂查询能力在众多项目中得到广泛应用,本文将从基础概念讲起,逐步深入到高... 目录前言1. PostgreSQL基础1.1 PostgreSQL简介1.2 基础语法1.3 数据库