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

Spark處理數(shù)據(jù)排序問題如何避免OOM

 更新時間:2020年05月21日 11:00:38   作者:Sheep Sun  
這篇文章主要介紹了Spark處理數(shù)據(jù)排序問題如何避免OOM,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下

錯誤思想

舉個列子,當我們想要比較 一個 類型為 RDD[(Long, (String, Int))] 的RDD,讓它先按Long分組,然后按int的值進行倒序排序,最容易想到的思維就是先分組,然后把Iterable 轉(zhuǎn)換為 list,然后sortby,但是這樣卻有一個致命的缺點,就是Iterable 在內(nèi)存中是一個指針,不占內(nèi)存,而list是一個容器,占用內(nèi)存,如果Iterable 含有元素過多,那么極易引起OOM

 val cidAndSidCountGrouped: RDD[(Long, Iterable[(String, Int)])] = cidAndSidCount.groupByKey()
    // 4. 排序, 取top10
    val result: RDD[(Long, List[(String, Int)])] = cidAndSidCountGrouped.map {
      case (cid, sidCountIt) =>
        // sidCountIt 排序, 取前10
        // Iterable轉(zhuǎn)成容器式集合的時候, 如果數(shù)據(jù)量過大, 極有可能導(dǎo)致oom
        (cid, sidCountIt.toList.sortBy(-_._2).take(5))
    }

首先,我們要知道,RDD 的排序需要 shuffle, 是采用了內(nèi)存+磁盤來完成的排序.這樣能有效避免OOM的風(fēng)險,但是RDD是全部排序,所以需要針對性的過濾Key值來進行排序

方法一 利用RDD排序特點

 //把long(即key值)提取出來
    val cids: List[Long] = categoryCountList.map(_.cid.toLong)
    val buffer: ListBuffer[(Long, List[(String, Int)])] = ListBuffer[(Long, List[(String, Int)])]()
    //根據(jù)每個key來過濾RDD
    for (cid <- cids) {
      /*
      List((15,(632972a4-f811-4000-b920-dc12ea803a41,10)), (15,(f34878b8-1784-4d81-a4d1-0c93ce53e942,8)), (15,(5e3545a0-1521-4ad6-91fe-e792c20c46da,8)), (15,(66a421b0-839d-49ae-a386-5fa3ed75226f,8)), (15,(9fa653ec-5a22-4938-83c5-21521d083cd0,8)))
      目標:
      (9,List((199f8e1d-db1a-4174-b0c2-ef095aaef3ee,9), (329b966c-d61b-46ad-949a-7e37142d384a,8), (5e3545a0-1521-4ad6-91fe-e792c20c46da,8), (e306c00b-a6c5-44c2-9c77-15e919340324,7), (bed60a57-3f81-4616-9e8b-067445695a77,7)))
       */
      val arr: Array[(String, Int)] = cidAndSidCount.filter(cid == _._1)
        .sortBy(-_._2._2)
        .take(5)
        .map(_._2)
      buffer += ((cid, arr.toList))
    }
    buffer.foreach(println)

這樣做也有缺點:即有多少個key,就有多少個Job,占用資源

方法二 利用TreeSet自動排序特性

 def statCategoryTop10Session_3(sc: SparkContext,
                  categoryCountList: List[CategroyCount],
                  userVisitActionRDD: RDD[UserVisitAction]) = {
    // 1. 過濾出來 top10品類的所有點擊記錄
    // 1.1 先map出來top10的品類id
    val cids = categoryCountList.map(_.cid.toLong)
    val topCategoryActionRDD: RDD[UserVisitAction] = userVisitActionRDD.filter(action => cids.contains(action.click_category_id))


    // 2. 計算每個品類 下的每個session 的點擊量 rdd ((cid, sid) ,1)
    val cidAndSidCount: RDD[(Long, (String, Int))] = topCategoryActionRDD
      .map(action => ((action.click_category_id, action.session_id), 1))
      // 使用自定義分區(qū)器 重點理解分區(qū)器的原理
      .reduceByKey(new CategoryPartitioner(cids), _ + _)
      .map {
        case ((cid, sid), count) => (cid, (sid, count))
      }
    
    // 3. 排序取top10
//因為已經(jīng)按key分好了區(qū),所以用Mappartitions ,在每個分區(qū)中新建一個TreeSet即可
    val result: RDD[(Long, List[SessionInfo])] = cidAndSidCount.mapPartitions((it: Iterator[(Long, (String, Int))]) => {
//new 一個TreeSet,并同時指定排序規(guī)則
   var treeSet: mutable.TreeSet[CategorySession] = new mutable.TreeSet[CategorySession]()(new Ordering[CategorySession] {
          override def compare(x: CategorySession, y: CategorySession): Int = {
            if (x.clickCount >= y.clickCount) -1 else 1
          }
        })
   var id = 0l
  iter.foreach({
    case (l, session) => {
      id = l
      treeSet.add(session)
    if (treeSet.size > 10) treeSet = treeSet.take(10)
          }
        })
        Iterator(id, treeSet)
      })
  
    result.collect.foreach(println)
    
    Thread.sleep(1000000)
  }
}

/*
根據(jù)傳入的key值來決定分區(qū)號,讓相同key進入相同的分區(qū),能夠避免多次shuffle
 */
class CategoryPartitioner(cids: List[Long]) extends Partitioner {
  // 用cid索引, 作為將來他的分區(qū)索引.
  private val cidWithIndex: Map[Long, Int] = cids.zipWithIndex.toMap
  
  // 返回集合的長度
  override def numPartitions: Int = cids.length
  
  // 根據(jù)key返回分區(qū)的索引
  override def getPartition(key: Any): Int = {
    key match {
      // 根據(jù)品類id返回分區(qū)的索引!  0-9
      case (cid: Long, _) =>
        cidWithIndex(cid)
    }
  }
}

以上就是本文的全部內(nèi)容,希望對大家的學(xué)習(xí)有所幫助,也希望大家多多支持腳本之家。

相關(guān)文章

  • Python中type()函數(shù)的具體使用

    Python中type()函數(shù)的具體使用

    在Python中,type()函數(shù)是一個非常有用的工具,它可以查看變量或?qū)ο蟮臄?shù)據(jù)類型,本文主要介紹了Python中type()函數(shù)的具體使用,感興趣的可以一起來了解一下
    2024-01-01
  • python生成式的send()方法(詳解)

    python生成式的send()方法(詳解)

    下面小編就為 大家?guī)硪黄猵ython生成式的send()方法(詳解)。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-05-05
  • 用Python編寫一個國際象棋AI程序

    用Python編寫一個國際象棋AI程序

    在這篇文章中我會介紹這個AI如何工作,每一個部分做什么,它為什么能那樣工作起來。你可以直接通讀本文,或者去下載代碼,邊讀邊看代碼。雖然去看看其他文件中有什么AI依賴的類也可能有幫助,但是AI部分全都在AI.py文件中
    2014-11-11
  • jupyter notebook 多環(huán)境conda kernel配置方式

    jupyter notebook 多環(huán)境conda kernel配置方式

    這篇文章主要介紹了jupyter notebook 多環(huán)境conda kernel配置方式,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2020-04-04
  • paramiko使用tail實時獲取服務(wù)器的日志輸出詳解

    paramiko使用tail實時獲取服務(wù)器的日志輸出詳解

    這篇文章主要給大家介紹了關(guān)于paramiko使用tail實時獲取服務(wù)器的日志輸出的相關(guān)資料,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-12-12
  • Python rindex()方法案例詳解

    Python rindex()方法案例詳解

    這篇文章主要介紹了Python rindex()方法案例詳解,本篇文章通過簡要的案例,講解了該項技術(shù)的了解與使用,以下就是詳細內(nèi)容,需要的朋友可以參考下
    2021-09-09
  • python matplotlib坐標軸設(shè)置的方法

    python matplotlib坐標軸設(shè)置的方法

    本篇文章主要介紹了python matplotlib坐標軸設(shè)置的方法,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-12-12
  • 基于Pytorch版yolov5的滑塊驗證碼破解思路詳解

    基于Pytorch版yolov5的滑塊驗證碼破解思路詳解

    這篇文章主要介紹了基于Pytorch版yolov5的滑塊驗證碼破解思路詳解,本文給大家介紹的非常詳細,對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-02-02
  • Python入門

    Python入門

    Python入門...
    2007-02-02
  • Python+OpenCV解決彩色圖亮度不均衡問題

    Python+OpenCV解決彩色圖亮度不均衡問題

    當我們換新頭像時,常常會遇到圖片過暗導(dǎo)致看不到圖片內(nèi)容的情況,本文將介紹如何通過Python和OpenCV解決色彩圖亮度不均衡的問題,需要的可以參考一下
    2021-12-12

最新評論

中牟县| 静安区| 阿坝县| 五大连池市| 华亭县| 石楼县| 临猗县| 桐城市| 富川| 扶绥县| 灵台县| 临颍县| 拉孜县| 正镶白旗| 庐江县| 巴青县| 石阡县| 六盘水市| 徐闻县| 淮滨县| 河北省| 中方县| 广饶县| 夏河县| 慈利县| 宁晋县| 玉环县| 阿勒泰市| 潞城市| 高青县| 广灵县| 盘山县| 玉山县| 穆棱市| 顺昌县| 南乐县| 资兴市| 台山市| 田阳县| 彭水| 石楼县|