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

Spark自定義累加器的使用實(shí)例詳解

 更新時(shí)間:2017年09月29日 11:32:21   作者:willian_zhang  
這篇文章主要介紹了Spark累加器的相關(guān)內(nèi)容,首先介紹了累加器的簡(jiǎn)單使用,然后向大家分享了自定義累加器的實(shí)例代碼,需要的朋友可以參考下。

累加器(accumulator)是Spark中提供的一種分布式的變量機(jī)制,其原理類(lèi)似于mapreduce,即分布式的改變,然后聚合這些改變。累加器的一個(gè)常見(jiàn)用途是在調(diào)試時(shí)對(duì)作業(yè)執(zhí)行過(guò)程中的事件進(jìn)行計(jì)數(shù)。

累加器簡(jiǎn)單使用

Spark內(nèi)置的提供了Long和Double類(lèi)型的累加器。下面是一個(gè)簡(jiǎn)單的使用示例,在這個(gè)例子中我們?cè)谶^(guò)濾掉RDD中奇數(shù)的同時(shí)進(jìn)行計(jì)數(shù),最后計(jì)算剩下整數(shù)的和。

val sparkConf = new SparkConf().setAppName("Test").setMaster("local[2]") 
val sc = new SparkContext(sparkConf) 
val accum = sc.longAccumulator("longAccum") //統(tǒng)計(jì)奇數(shù)的個(gè)數(shù) 
val sum = sc.parallelize(Array(1,2,3,4,5,6,7,8,9),2).filter(n=>{ 
 if(n%2!=0) accum.add(1L)  
 n%2==0 
}).reduce(_+_) 
println("sum: "+sum) 
println("accum: "+accum.value) 
sc.stop() 

結(jié)果為:

sum: 20
accum: 5

這是結(jié)果正常的情況,但是在使用累加器的過(guò)程中如果對(duì)于spark的執(zhí)行過(guò)程理解的不夠深入就會(huì)遇到兩類(lèi)典型的錯(cuò)誤:少加(或者沒(méi)加)、多加。

自定義累加器

自定義累加器類(lèi)型的功能在1.X版本中就已經(jīng)提供了,但是使用起來(lái)比較麻煩,在2.0版本后,累加器的易用性有了較大的改進(jìn),而且官方還提供了一個(gè)新的抽象類(lèi):AccumulatorV2來(lái)提供更加友好的自定義類(lèi)型累加器的實(shí)現(xiàn)方式。官方同時(shí)給出了一個(gè)實(shí)現(xiàn)的示例:CollectionAccumulator類(lèi),這個(gè)類(lèi)允許以集合的形式收集spark應(yīng)用執(zhí)行過(guò)程中的一些信息。例如,我們可以用這個(gè)類(lèi)收集Spark處理數(shù)據(jù)時(shí)的一些細(xì)節(jié),當(dāng)然,由于累加器的值最終要匯聚到driver端,為了避免 driver端的outofmemory問(wèn)題,需要對(duì)收集的信息的規(guī)模要加以控制,不宜過(guò)大。

繼承AccumulatorV2類(lèi),并復(fù)寫(xiě)它的所有方法

package spark
import constant.Constant
import org.apache.spark.util.AccumulatorV2
import util.getFieldFromConcatString
import util.setFieldFromConcatString
open class SessionAccmulator : AccumulatorV2<String, String>() {
  private var result = Constant.SESSION_COUNT + "=0|"+
      Constant.TIME_PERIOD_1s_3s + "=0|"+
      Constant.TIME_PERIOD_4s_6s + "=0|"+
      Constant.TIME_PERIOD_7s_9s + "=0|"+
      Constant.TIME_PERIOD_10s_30s + "=0|"+
      Constant.TIME_PERIOD_30s_60s + "=0|"+
      Constant.TIME_PERIOD_1m_3m + "=0|"+
      Constant.TIME_PERIOD_3m_10m + "=0|"+
      Constant.TIME_PERIOD_10m_30m + "=0|"+
      Constant.TIME_PERIOD_30m + "=0|"+
      Constant.STEP_PERIOD_1_3 + "=0|"+
      Constant.STEP_PERIOD_4_6 + "=0|"+
      Constant.STEP_PERIOD_7_9 + "=0|"+
      Constant.STEP_PERIOD_10_30 + "=0|"+
      Constant.STEP_PERIOD_30_60 + "=0|"+
      Constant.STEP_PERIOD_60 + "=0"
  override fun value(): String {
    return this.result
  }
  /**
   * 合并數(shù)據(jù)
   */
  override fun merge(other: AccumulatorV2<String, String>?) {
    if (other == null) return else {
      if (other is SessionAccmulator) {
        var newResult = ""
        val resultArray = arrayOf(Constant.SESSION_COUNT,Constant.TIME_PERIOD_1s_3s, Constant.TIME_PERIOD_4s_6s, Constant.TIME_PERIOD_7s_9s,
            Constant.TIME_PERIOD_10s_30s, Constant.TIME_PERIOD_30s_60s, Constant.TIME_PERIOD_1m_3m,
            Constant.TIME_PERIOD_3m_10m, Constant.TIME_PERIOD_10m_30m, Constant.TIME_PERIOD_30m,
            Constant.STEP_PERIOD_1_3, Constant.STEP_PERIOD_4_6, Constant.STEP_PERIOD_7_9,
            Constant.STEP_PERIOD_10_30, Constant.STEP_PERIOD_30_60, Constant.STEP_PERIOD_60)
        resultArray.forEach {
          val oldValue = other.result.getFieldFromConcatString("|", it)
          if (oldValue.isNotEmpty()) {
            val newValue = oldValue.toInt() + 1
            //找到原因,一直在循環(huán)賦予值,debug30分鐘 很煩
            if (newResult.isEmpty()){
              newResult = result.setFieldFromConcatString("|", it, newValue.toString())
            }
            //問(wèn)題就在于這里,自定義沒(méi)有寫(xiě)錯(cuò),合并錯(cuò)了
            newResult = newResult.setFieldFromConcatString("|", it, newValue.toString())
          }
        }
        result = newResult
      }
    }
  }
  override fun copy(): AccumulatorV2<String, String> {
    val sessionAccmulator = SessionAccmulator()
    sessionAccmulator.result = this.result
    return sessionAccmulator
  }
  override fun add(p0: String?) {
    val v1 = this.result
    val v2 = p0
    if (v2.isNullOrEmpty()){
      return
    }else{
      var newResult = ""
      val oldValue = v1.getFieldFromConcatString("|", v2!!)
      if (oldValue.isNotEmpty()){
        val newValue = oldValue.toInt() + 1
        newResult = result.setFieldFromConcatString("|", v2, newValue.toString())
      }
      result = newResult
    }
  }
  override fun reset() {
    val newResult = Constant.SESSION_COUNT + "=0|"+
        Constant.TIME_PERIOD_1s_3s + "=0|"+
        Constant.TIME_PERIOD_4s_6s + "=0|"+
        Constant.TIME_PERIOD_7s_9s + "=0|"+
        Constant.TIME_PERIOD_10s_30s + "=0|"+
        Constant.TIME_PERIOD_30s_60s + "=0|"+
        Constant.TIME_PERIOD_1m_3m + "=0|"+
        Constant.TIME_PERIOD_3m_10m + "=0|"+
        Constant.TIME_PERIOD_10m_30m + "=0|"+
        Constant.TIME_PERIOD_30m + "=0|"+
        Constant.STEP_PERIOD_1_3 + "=0|"+
        Constant.STEP_PERIOD_4_6 + "=0|"+
        Constant.STEP_PERIOD_7_9 + "=0|"+
        Constant.STEP_PERIOD_10_30 + "=0|"+
        Constant.STEP_PERIOD_30_60 + "=0|"+
        Constant.STEP_PERIOD_60 + "=0"
    result = newResult
  }
  override fun isZero(): Boolean {
    val newResult = Constant.SESSION_COUNT + "=0|"+
        Constant.TIME_PERIOD_1s_3s + "=0|"+
        Constant.TIME_PERIOD_4s_6s + "=0|"+
        Constant.TIME_PERIOD_7s_9s + "=0|"+
        Constant.TIME_PERIOD_10s_30s + "=0|"+
        Constant.TIME_PERIOD_30s_60s + "=0|"+
        Constant.TIME_PERIOD_1m_3m + "=0|"+
        Constant.TIME_PERIOD_3m_10m + "=0|"+
        Constant.TIME_PERIOD_10m_30m + "=0|"+
        Constant.TIME_PERIOD_30m + "=0|"+
        Constant.STEP_PERIOD_1_3 + "=0|"+
        Constant.STEP_PERIOD_4_6 + "=0|"+
        Constant.STEP_PERIOD_7_9 + "=0|"+
        Constant.STEP_PERIOD_10_30 + "=0|"+
        Constant.STEP_PERIOD_30_60 + "=0|"+
        Constant.STEP_PERIOD_60 + "=0"
    return this.result == newResult
  }
}

方法介紹

value方法:獲取累加器中的值

       merge方法:該方法特別重要,一定要寫(xiě)對(duì),這個(gè)方法是各個(gè)task的累加器進(jìn)行合并的方法(下面介紹執(zhí)行流程中將要用到)

        iszero方法:判斷是否為初始值

        reset方法:重置累加器中的值

        copy方法:拷貝累加器

spark中累加器的執(zhí)行流程:

          首先有幾個(gè)task,spark engine就調(diào)用copy方法拷貝幾個(gè)累加器(不注冊(cè)的),然后在各個(gè)task中進(jìn)行累加(注意在此過(guò)程中,被最初注冊(cè)的累加器的值是不變的),執(zhí)行最后將調(diào)用merge方法和各個(gè)task的結(jié)果累計(jì)器進(jìn)行合并(此時(shí)被注冊(cè)的累加器是初始值)

總結(jié)

以上就是本文關(guān)于Spark自定義累加器的使用實(shí)例詳解的全部?jī)?nèi)容,希望對(duì)大家有所幫助。有什么問(wèn)題可以隨時(shí)留言,小編會(huì)及時(shí)回復(fù)大家的。

相關(guān)文章

  • zerotier搭建免費(fèi)moon服務(wù)器的部署流程

    zerotier搭建免費(fèi)moon服務(wù)器的部署流程

    ZeroTier是一種基于P2P的虛擬組網(wǎng)工具,通過(guò)搭建Moon服務(wù)器?可大幅提升跨運(yùn)營(yíng)商/跨國(guó)節(jié)點(diǎn)的連接質(zhì)量,本文介紹了如何使用云服務(wù)部署ZeroTier的Moon服務(wù)器,并詳細(xì)步驟包括登錄服務(wù)器、安裝ZeroTier、生成Moon配置文件、配置Moon服務(wù)器和重啟服務(wù),感興趣的朋友一起看看吧
    2025-03-03
  • Cisco網(wǎng)絡(luò)防火墻配置方法

    Cisco網(wǎng)絡(luò)防火墻配置方法

    這篇文章主要介紹了Cisco網(wǎng)絡(luò)防火墻配置方法,需要的朋友可以參考下
    2016-04-04
  • 忘記Grafana不要緊2種Grafana重置admin密碼方法詳細(xì)步驟

    忘記Grafana不要緊2種Grafana重置admin密碼方法詳細(xì)步驟

    這篇文章主要介紹了忘記Grafana不要緊2種Grafana重置admin密碼方法詳細(xì)步驟,需要的朋友可以參考下
    2022-04-04
  • pgpool-II搭建集群,實(shí)現(xiàn)高可用與讀寫(xiě)分離

    pgpool-II搭建集群,實(shí)現(xiàn)高可用與讀寫(xiě)分離

    pgpool-II是開(kāi)源的PostgreSQL數(shù)據(jù)庫(kù)連接池、負(fù)載均衡和高可用解決方案,支持多種工作模式,包括原始模式、內(nèi)置復(fù)制模式和主/備模式,本文介紹pgpool-II的架構(gòu)、進(jìn)程、工作模式以及配置步驟,包括環(huán)境規(guī)劃、系統(tǒng)準(zhǔn)備、軟件安裝、數(shù)據(jù)庫(kù)主節(jié)點(diǎn)配置、pgpool配置和集群?jiǎn)?dòng)
    2025-04-04
  • Centos中VNC遠(yuǎn)程桌面程序的安裝與使用教程

    Centos中VNC遠(yuǎn)程桌面程序的安裝與使用教程

    這篇文章主要介紹了Centos中VNC遠(yuǎn)程桌面程序的安裝與使用的方法,較為詳細(xì)的分析了CentOS的VNC遠(yuǎn)程桌面程序安裝、配置、連接、啟動(dòng)等命令與相關(guān)操作技巧,需要的朋友可以參考下
    2016-07-07
  • rsync?常見(jiàn)錯(cuò)誤與解決方法整理

    rsync?常見(jiàn)錯(cuò)誤與解決方法整理

    由于我們經(jīng)常使用rsync進(jìn)行服務(wù)器文件的同步工作,但在配置過(guò)程中,會(huì)出現(xiàn)很多問(wèn)題,下面的錯(cuò)誤基本上都是通過(guò)客戶(hù)端返回的錯(cuò)誤進(jìn)行分析
    2012-11-11
  • XAMPP下使用頂級(jí)域名綁定虛擬主機(jī)的配置方法和示例

    XAMPP下使用頂級(jí)域名綁定虛擬主機(jī)的配置方法和示例

    這篇文章主要介紹了XAMPP下使用頂級(jí)域名綁定虛擬主機(jī)的配置方法和示例,XAMPP是Windows下非常好用的一款集成開(kāi)發(fā)環(huán)境,需要的朋友可以參考下
    2014-07-07
  • RSync實(shí)現(xiàn)文件備份同步詳解

    RSync實(shí)現(xiàn)文件備份同步詳解

    rsync,remote synchronize顧名思意就知道它是一款實(shí)現(xiàn)遠(yuǎn)程同步功能的軟件,它在同步文件的同時(shí),可以保持原來(lái)文件的權(quán)限、時(shí)間、軟硬鏈接等附加信息
    2016-03-03
  • windows系統(tǒng)搭建zookeeper服務(wù)器的教程

    windows系統(tǒng)搭建zookeeper服務(wù)器的教程

    這篇文章主要介紹了windows系統(tǒng)搭建zookeeper服務(wù)器的教程,本文圖文并茂給大家介紹的非常詳細(xì),具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2019-10-10
  • curl.exe安裝使用的最全參數(shù)詳解以及常用命令匯總

    curl.exe安裝使用的最全參數(shù)詳解以及常用命令匯總

    Curl是一個(gè)功能強(qiáng)大的命令行工具,可以看做是命令行瀏覽器,用于與服務(wù)器進(jìn)行數(shù)據(jù)交互,支持多種數(shù)據(jù)傳輸協(xié)議,如HTTP、HTTPS、FTP等,它支持文件的上傳和下載,它是一款開(kāi)源軟件,在多個(gè)操作系統(tǒng)上均可運(yùn)行,包括Windows、Linux、macOS等
    2024-04-04

最新評(píng)論

黎平县| 乌拉特中旗| 无极县| 娱乐| 宾阳县| 辽阳县| 伽师县| 阿城市| 葫芦岛市| 周口市| 新河县| 永清县| 洞口县| 喀喇| 高青县| 元阳县| 齐河县| 泗阳县| 固安县| 东丽区| 阿拉善左旗| 昌都县| 吉木萨尔县| 济阳县| 阿拉善左旗| 怀仁县| 延边| 昔阳县| 东阳市| 太谷县| 建始县| 额济纳旗| 安乡县| 舞钢市| 松桃| 疏勒县| 大英县| 许昌市| 卓资县| 固原市| 镇坪县|