最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Spark的廣播變量和累加器使用方法代碼示例

 更新時(shí)間:2017年09月29日 10:57:24   作者:唐予之_  
這篇文章主要介紹了Spark的廣播變量和累加器使用方法代碼示例,文中介紹了廣播變量和累加器的含義,然后通過(guò)實(shí)例演示了其用法,需要的朋友可以參考下。

一、廣播變量和累加器

通常情況下,當(dāng)向Spark操作(如map,reduce)傳遞一個(gè)函數(shù)時(shí),它會(huì)在一個(gè)遠(yuǎn)程集群節(jié)點(diǎn)上執(zhí)行,它會(huì)使用函數(shù)中所有變量的副本。這些變量被復(fù)制到所有的機(jī)器上,遠(yuǎn)程機(jī)器上并沒(méi)有被更新的變量會(huì)向驅(qū)動(dòng)程序回傳。在任務(wù)之間使用通用的,支持讀寫的共享變量是低效的。盡管如此,Spark提供了兩種有限類型的共享變量,廣播變量和累加器。

1.1 廣播變量:

廣播變量允許程序員將一個(gè)只讀的變量緩存在每臺(tái)機(jī)器上,而不用在任務(wù)之間傳遞變量。廣播變量可被用于有效地給每個(gè)節(jié)點(diǎn)一個(gè)大輸入數(shù)據(jù)集的副本。Spark還嘗試使用高效地廣播算法來(lái)分發(fā)變量,進(jìn)而減少通信的開銷。

Spark的動(dòng)作通過(guò)一系列的步驟執(zhí)行,這些步驟由分布式的shuffle操作分開。Spark自動(dòng)地廣播每個(gè)步驟每個(gè)任務(wù)需要的通用數(shù)據(jù)。這些廣播數(shù)據(jù)被序列化地緩存,在運(yùn)行任務(wù)之前被反序列化出來(lái)。這意味著當(dāng)我們需要在多個(gè)階段的任務(wù)之間使用相同的數(shù)據(jù),或者以反序列化形式緩存數(shù)據(jù)是十分重要的時(shí)候,顯式地創(chuàng)建廣播變量才有用。

通過(guò)在一個(gè)變量v上調(diào)用SparkContext.broadcast(v)可以創(chuàng)建廣播變量。廣播變量是圍繞著v的封裝,可以通過(guò)value方法訪問(wèn)這個(gè)變量。舉例如下:

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)

在創(chuàng)建了廣播變量之后,在集群上的所有函數(shù)中應(yīng)該使用它來(lái)替代使用v.這樣v就不會(huì)不止一次地在節(jié)點(diǎn)之間傳輸了。另外,為了確保所有的節(jié)點(diǎn)獲得相同的變量,對(duì)象v在被廣播之后就不應(yīng)該再修改。

1.2 累加器:

累加器是僅僅被相關(guān)操作累加的變量,因此可以在并行中被有效地支持。它可以被用來(lái)實(shí)現(xiàn)計(jì)數(shù)器和總和。Spark原生地只支持?jǐn)?shù)字類型的累加器,編程者可以添加新類型的支持。如果創(chuàng)建累加器時(shí)指定了名字,可以在Spark的UI界面看到。這有利于理解每個(gè)執(zhí)行階段的進(jìn)程。(對(duì)于python還不支持)

累加器通過(guò)對(duì)一個(gè)初始化了的變量v調(diào)用SparkContext.accumulator(v)來(lái)創(chuàng)建。在集群上運(yùn)行的任務(wù)可以通過(guò)add或者”+=”方法在累加器上進(jìn)行累加操作。但是,它們不能讀取它的值。只有驅(qū)動(dòng)程序能夠讀取它的值,通過(guò)累加器的value方法。

下面的代碼展示了如何把一個(gè)數(shù)組中的所有元素累加到累加器上:

scala> val accum = sc.accumulator(0, "My Accumulator")
accum: spark.Accumulator[Int] = 0
scala> sc.parallelize(Array(1, 2, 3, 4)).foreach(x => accum += x)
...
10/09/29 18:41:08 INFO SparkContext: Tasks finished in 0.317106 s
scala> accum.value
res2: Int = 10

盡管上面的例子使用了內(nèi)置支持的累加器類型Int,但是開發(fā)人員也可以通過(guò)繼承AccumulatorParam類來(lái)創(chuàng)建它們自己的累加器類型。AccumulatorParam接口有兩個(gè)方法:
zero方法為你的類型提供一個(gè)0值。
addInPlace方法將兩個(gè)值相加。
假設(shè)我們有一個(gè)代表數(shù)學(xué)vector的Vector類。我們可以向下面這樣實(shí)現(xiàn):

object VectorAccumulatorParam extends AccumulatorParam[Vector] {
 def zero(initialValue: Vector): Vector = {
  Vector.zeros(initialValue.size)
 }
 def addInPlace(v1: Vector, v2: Vector): Vector = {
  v1 += v2
 }
}
// Then, create an Accumulator of this type:
val vecAccum = sc.accumulator(new Vector(...))(VectorAccumulatorParam)

在Scala里,Spark提供更通用的累加接口來(lái)累加數(shù)據(jù),盡管結(jié)果的類型和累加的數(shù)據(jù)類型可能不一致(例如,通過(guò)收集在一起的元素來(lái)創(chuàng)建一個(gè)列表)。同時(shí),SparkContext..accumulableCollection方法來(lái)累加通用的Scala的集合類型。

累加器僅僅在動(dòng)作操作內(nèi)部被更新,Spark保證每個(gè)任務(wù)在累加器上的更新操作只被執(zhí)行一次,也就是說(shuō),重啟任務(wù)也不會(huì)更新。在轉(zhuǎn)換操作中,用戶必須意識(shí)到每個(gè)任務(wù)對(duì)累加器的更新操作可能被不只一次執(zhí)行,如果重新執(zhí)行了任務(wù)和作業(yè)的階段。

累加器并沒(méi)有改變Spark的惰性求值模型。如果它們被RDD上的操作更新,它們的值只有當(dāng)RDD因?yàn)閯?dòng)作操作被計(jì)算時(shí)才被更新。因此,當(dāng)執(zhí)行一個(gè)惰性的轉(zhuǎn)換操作,比如map時(shí),不能保證對(duì)累加器值的更新被實(shí)際執(zhí)行了。下面的代碼片段演示了此特性:

val accum = sc.accumulator(0)
data.map { x => accum += x; f(x) }
//在這里,accum的值仍然是0,因?yàn)闆](méi)有動(dòng)作操作引起map被實(shí)際的計(jì)算.

二.Java和Scala版本的實(shí)戰(zhàn)演示

2.1 Java版本:

/**
 * 實(shí)例:利用廣播進(jìn)行黑名單過(guò)濾!
 * 檢查新的數(shù)據(jù) 根據(jù)是否在廣播變量-黑名單內(nèi),從而實(shí)現(xiàn)過(guò)濾數(shù)據(jù)。
 */
public class BroadcastAccumulator {
 /**
  * 創(chuàng)建一個(gè)List的廣播變量
  *
  */
 private static volatile Broadcast<List<String>> broadcastList = null;
 /**
  * 計(jì)數(shù)器!
  */
 private static volatile Accumulator<Integer> accumulator = null;
 public static void main(String[] args) {
  SparkConf conf = new SparkConf().setMaster("local[2]").
    setAppName("WordCountOnlineBroadcast");
  JavaStreamingContext jsc = new JavaStreamingContext(conf, Durations.seconds(5));

  /**
   * 注意:分發(fā)廣播需要一個(gè)action操作觸發(fā)。
   * 注意:廣播的是Arrays的asList 而非對(duì)象的引用。廣播Array數(shù)組的對(duì)象引用會(huì)出錯(cuò)。
   * 使用broadcast廣播黑名單到每個(gè)Executor中!
   */
  broadcastList = jsc.sc().broadcast(Arrays.asList("Hadoop","Mahout","Hive"));
  /**
   * 累加器作為全局計(jì)數(shù)器!用于統(tǒng)計(jì)在線過(guò)濾了多少個(gè)黑名單!
   * 在這里實(shí)例化。
   */
  accumulator = jsc.sparkContext().accumulator(0,"OnlineBlackListCounter");

  JavaReceiverInputDStream<String> lines = jsc.socketTextStream("Master", 9999);

  /**
   * 這里省去flatmap因?yàn)槊麊问且粋€(gè)個(gè)的!
   */
  JavaPairDStream<String, Integer> pairs = lines.mapToPair(new PairFunction<String, String, Integer>() {
   @Override
   public Tuple2<String, Integer> call(String word) {
    return new Tuple2<String, Integer>(word, 1);
   }
  });
  JavaPairDStream<String, Integer> wordsCount = pairs.reduceByKey(new Function2<Integer, Integer, Integer>() {
   @Override
   public Integer call(Integer v1, Integer v2) {
    return v1 + v2;
   }
  });
  /**
   * Funtion里面 前幾個(gè)參數(shù)是 入?yún)ⅰ?
   * 后面的出參。
   * 體現(xiàn)在call方法里面!
   *
   */
  wordsCount.foreach(new Function2<JavaPairRDD<String, Integer>, Time, Void>() {
   @Override
   public Void call(JavaPairRDD<String, Integer> rdd, Time time) throws Exception {
    rdd.filter(new Function<Tuple2<String, Integer>, Boolean>() {
     @Override
     public Boolean call(Tuple2<String, Integer> wordPair) throws Exception {
      if (broadcastList.value().contains(wordPair._1)) {
       /**
        * accumulator不僅僅用來(lái)計(jì)數(shù)。
        * 可以同時(shí)寫進(jìn)數(shù)據(jù)庫(kù)或者緩存中。
        */
       accumulator.add(wordPair._2);
       return false;
      }else {
       return true;
      }
     };
     /**
      * 廣播和計(jì)數(shù)器的執(zhí)行,需要進(jìn)行一個(gè)action操作!
      */
    }).collect();
    System.out.println("廣播器里面的值"+broadcastList.value());
    System.out.println("計(jì)時(shí)器里面的值"+accumulator.value());
    return null;
   }
  });

  jsc.start();
  jsc.awaitTermination();
  jsc.close();
 }
 }

2.2 Scala版本

package com.Streaming
import java.util
import org.apache.spark.streaming.{Duration, StreamingContext}
import org.apache.spark.{Accumulable, Accumulator, SparkContext, SparkConf}
import org.apache.spark.broadcast.Broadcast
/**
 * Created by lxh on 2016/6/30.
 */
object BroadcastAccumulatorStreaming {
 /**
 * 聲明一個(gè)廣播和累加器!
 */
 private var broadcastList:Broadcast[List[String]] = _
 private var accumulator:Accumulator[Int] = _
 def main(args: Array[String]) {
 val sparkConf = new SparkConf().setMaster("local[4]").setAppName("broadcasttest")
 val sc = new SparkContext(sparkConf)
 /**
  * duration是ms
  */
 val ssc = new StreamingContext(sc,Duration(2000))
 // broadcastList = ssc.sparkContext.broadcast(util.Arrays.asList("Hadoop","Spark"))
 broadcastList = ssc.sparkContext.broadcast(List("Hadoop","Spark"))
 accumulator= ssc.sparkContext.accumulator(0,"broadcasttest")
 /**
  * 獲取數(shù)據(jù)!
  */
 val lines = ssc.socketTextStream("localhost",9999)
 /**
  * 1.flatmap把行分割成詞。
  * 2.map把詞變成tuple(word,1)
  * 3.reducebykey累加value
  * (4.sortBykey排名)
  * 4.進(jìn)行過(guò)濾。 value是否在累加器中。
  * 5.打印顯示。
  */
 val words = lines.flatMap(line => line.split(" "))
 val wordpair = words.map(word => (word,1))
 wordpair.filter(record => {broadcastList.value.contains(record._1)})
 val pair = wordpair.reduceByKey(_+_)
 /**
  * 這個(gè)pair 是PairDStream<String, Integer>
  * 查看這個(gè)id是否在黑名單中,如果是的話,累加器就+1
  */
/* pair.foreachRDD(rdd => {
  rdd.filter(record => {
  if (broadcastList.value.contains(record._1)) {
   accumulator.add(1)
   return true
  } else {
   return false
  }
  })
 })*/
 val filtedpair = pair.filter(record => {
  if (broadcastList.value.contains(record._1)) {
   accumulator.add(record._2)
   true
  } else {
   false
  }
  }).print
 println("累加器的值"+accumulator.value)
 // pair.filter(record => {broadcastList.value.contains(record._1)})
 /* val keypair = pair.map(pair => (pair._2,pair._1))*/
 /**
  * 如果DStream自己沒(méi)有某個(gè)算子操作。就通過(guò)轉(zhuǎn)化transform!
  */
 /* keypair.transform(rdd => {
  rdd.sortByKey(false)//TODO
 })*/
 pair.print()
 ssc.start()
 ssc.awaitTermination()
 }
}

總結(jié)

以上就是本文關(guān)于Spark的廣播變量和累加器使用方法代碼示例的全部?jī)?nèi)容,希望對(duì)大家有所幫助。感興趣的朋友可以參閱:詳解Java編寫并運(yùn)行spark應(yīng)用程序的方法  、 Spark入門簡(jiǎn)介等,有什么問(wèn)題可以隨時(shí)留言,小編會(huì)及時(shí)回復(fù)大家。感謝朋友們對(duì)腳本之家網(wǎng)站的支持。

相關(guān)文章

  • Linux下Web性能壓力測(cè)試工具h(yuǎn)ttp_load使用教程

    Linux下Web性能壓力測(cè)試工具h(yuǎn)ttp_load使用教程

    http_load基于linux平臺(tái)的一種性能測(cè)工具。以并行復(fù)用的方式運(yùn)行,用以測(cè)試web服務(wù)器的吞吐量與負(fù)載,測(cè)試web頁(yè)面的性能。
    2014-11-11
  • nfs和web服務(wù)器的搭建過(guò)程

    nfs和web服務(wù)器的搭建過(guò)程

    這篇文章主要介紹了nfs和web服務(wù)器的搭建過(guò)程,本文通過(guò)圖文并茂的形式給大家介紹的非常詳細(xì),感興趣的朋友跟隨小編一起看看吧
    2024-07-07
  • TortoiseSvn小烏龜安裝最新圖文詳細(xì)教程

    TortoiseSvn小烏龜安裝最新圖文詳細(xì)教程

    在使用TortoiseSvn安裝過(guò)程中經(jīng)常出現(xiàn)各種奇葩問(wèn)題,干脆換成svn吧,在這里小編把我的安裝過(guò)程記錄下來(lái),有需要的朋友直接拿去用吧
    2021-05-05
  • VPS主機(jī)快速搬家方法:邊打包邊傳輸邊解壓適合大中型論壇網(wǎng)站

    VPS主機(jī)快速搬家方法:邊打包邊傳輸邊解壓適合大中型論壇網(wǎng)站

    本篇文章給大家分享如何在VPS主機(jī)之間快速搬家,一邊打包壓縮原主機(jī)上的文件,一邊傳輸文件數(shù)據(jù)到新的主機(jī)上,一邊在新的VPS主機(jī)上解壓文件,因?yàn)樗械牟僮鞫际窃赩PS主機(jī)上之間進(jìn)行,傳輸速度可以達(dá)到幾MB/s以上,特別適合一些大中型的論壇和網(wǎng)站搬家
    2017-07-07
  • web壓力測(cè)試工具_(dá)動(dòng)力節(jié)點(diǎn)Java 學(xué)院整理

    web壓力測(cè)試工具_(dá)動(dòng)力節(jié)點(diǎn)Java 學(xué)院整理

    本文給大家分享幾個(gè)web 壓力測(cè)試工具,非常不錯(cuò),具有參考借鑒價(jià)值,需要的的朋友參考下吧
    2017-08-08
  • 基于IntelliJ IDEA運(yùn)行慢的解決方法

    基于IntelliJ IDEA運(yùn)行慢的解決方法

    下面小編就為大家分享一篇基于IntelliJ IDEA運(yùn)行慢的解決方法,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2017-11-11
  • git忽略特殊文件_動(dòng)力節(jié)點(diǎn)Java學(xué)院整理

    git忽略特殊文件_動(dòng)力節(jié)點(diǎn)Java學(xué)院整理

    這篇文章主要為大家詳細(xì)介紹了git忽略特殊文件的相關(guān)資料,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2017-08-08
  • zlmediakit實(shí)現(xiàn) rtsp流服務(wù)器的方法

    zlmediakit實(shí)現(xiàn) rtsp流服務(wù)器的方法

    這篇文章主要介紹了zlmediakit實(shí)現(xiàn) rtsp流服務(wù)器,本次實(shí)現(xiàn)是將內(nèi)存中的H264數(shù)據(jù)經(jīng)過(guò)zlmediakit實(shí)現(xiàn)為rtsp流,我用的是CAPI的方式,將zlmediakit作為一個(gè)sdk嵌入到自己的程序中而不是作為一個(gè)獨(dú)立的進(jìn)進(jìn)程服務(wù),需要的朋友可以參考下
    2024-04-04
  • windows安裝OpenSSL的方法小結(jié)

    windows安裝OpenSSL的方法小結(jié)

    openssl是一個(gè)強(qiáng)大的安全套接字密碼庫(kù),囊括主要的密碼算法、常用的密鑰和證書封裝管理功能及SSL協(xié)議,并提供豐富的應(yīng)用程序供測(cè)試或其他目的使用
    2023-09-09
  • 阿里云服務(wù)器?jdk1.8?安裝配置教程

    阿里云服務(wù)器?jdk1.8?安裝配置教程

    這篇文章主要介紹了阿里云服務(wù)器?jdk1.8?安裝配置,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2022-12-12

最新評(píng)論

泽库县| 渝中区| 茶陵县| 临猗县| 曲水县| 麻阳| 遂溪县| 昌宁县| 申扎县| 新闻| 凉山| 唐山市| 大厂| 铜川市| 灵寿县| 正定县| 津南区| 盐池县| 马尔康县| 潼关县| 公安县| 肇州县| 鄂托克旗| 叶城县| 富蕴县| 清镇市| 奉化市| 叙永县| 垣曲县| 从化市| 深泽县| 岳阳市| 华阴市| 全南县| 岑巩县| 抚宁县| 新野县| 武定县| 三江| 四平市| 高淳县|