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

Spark調(diào)優(yōu)多線程并行處理任務(wù)實現(xiàn)方式

 更新時間:2020年08月06日 10:21:41   作者:lshan  
這篇文章主要介紹了Spark調(diào)優(yōu)多線程并行處理任務(wù)實現(xiàn)方式,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下

方式1:

1. 明確 Spark中Job 與 Streaming中 Job 的區(qū)別

1.1 Spark Core

一個 RDD DAG Graph 可以生成一個或多個 Job(Action操作)

一個Job可以認(rèn)為就是會最終輸出一個結(jié)果RDD的一條由RDD組織而成的計算

Job在spark里應(yīng)用里是一個被調(diào)度的單位

1.2 Streaming

一個 batch 的數(shù)據(jù)對應(yīng)一個 DStreamGraph

而一個 DStreamGraph 包含一或多個關(guān)于 DStream 的輸出操作

每一個輸出對應(yīng)于一個Job,一個 DStreamGraph 對應(yīng)一個JobSet,里面包含一個或多個Job

2. Streaming Job的并行度

Job的并行度由兩個配置決定:

spark.scheduler.mode(FIFO/FAIR)
spark.streaming.concurrentJobs

一個 Batch 可能會有多個 Action 執(zhí)行,比如注冊了多個 Kafka 數(shù)據(jù)流,每個Action都會產(chǎn)生一個Job

所以一個 Batch 有可能是一批 Job,也就是 JobSet 的概念

這些 Job 由 jobExecutor 依次提交執(zhí)行

而 JobExecutor 是一個默認(rèn)池子大小為1的線程池,所以只能執(zhí)行完一個Job再執(zhí)行另外一個Job

這里說的池子,大小就是由spark.streaming.concurrentJobs 控制的

concurrentJobs 決定了向 Spark Core 提交Job的并行度

提交一個Job,必須等這個執(zhí)行完了,才會提交第二個

假設(shè)我們把它設(shè)置為2,則會并發(fā)的把 Job 提交給 Spark Core

Spark 有自己的機(jī)制決定如何運行這兩個Job,這個機(jī)制其實就是FIFO或者FAIR(決定了資源的分配規(guī)則)

默認(rèn)是 FIFO,也就是先進(jìn)先出,把 concurrentJobs 設(shè)置為2,但是如果底層是FIFO,那么會優(yōu)先執(zhí)行先提交的Job

雖然如此,如果資源夠兩個job運行,還是會并行運行兩個Job

Spark Streaming 不同Batch任務(wù)可以并行計算么 https://developer.aliyun.com/article/73004

conf.setMaster("local[4]")
conf.set("spark.streaming.concurrentJobs", "3") //job 并行對
conf.set("spark.scheduler.mode", "FIFO")
val sc = new StreamingContext(conf, Seconds(5))

你會發(fā)現(xiàn),不同batch的job其實也可以并行運行的,這里需要有幾個條件:

有延時發(fā)生了,batch無法在本batch完成

concurrentJobs > 1

如果scheduler mode 是FIFO則需要某個Job無法一直消耗掉所有資源

Mode是FAIR則盡力保證你的Job是并行運行的,毫無疑問是可以并行的。

方式2:

場景1:

程序每次處理的數(shù)據(jù)量是波動的,比如周末比工作日多很多,晚八點比凌晨四點多很多。

一個spark程序處理的時間在1-2小時波動是OK的。而spark streaming程序不可以,如果每次處理的時間是1-10分鐘,就很蛋疼。
設(shè)置10分鐘吧,實際上10分鐘的也就那一段高峰時間,如果設(shè)置每次是1分鐘,很多時候會出現(xiàn)程序處理不過來,排隊過多的任務(wù)延遲更久,還可能出現(xiàn)程序崩潰的可能。

場景2:

  • 程序需要處理的相似job數(shù)隨著業(yè)務(wù)的增長越來越多
  • 我們知道spark的api里無相互依賴的stage是并行處理的,但是job之間是串行處理的。
  • spark程序通常是離線處理,比如T+1之類的延遲,時間變長是可以容忍的。而spark streaming是準(zhǔn)實時的,如果業(yè)務(wù)增長導(dǎo)致延遲增加就很不合理。

spark雖然是串行執(zhí)行job,但是是可以把job放到線程池里多線程執(zhí)行的。如何在一個SparkContext中提交多個任務(wù)

DStream.foreachRDD{
   rdd =>
    //創(chuàng)建線程池
    val executors=Executors.newFixedThreadPool(rules.length)
    //將規(guī)則放入線程池
    for( ru <- rules){
     val task= executors.submit(new Callable[String] {
      override def call(): String ={
       //執(zhí)行規(guī)則
       runRule(ru,spark)
      }
     })
    }
    //每次創(chuàng)建的線程池執(zhí)行完所有規(guī)則后shutdown
    executors.shutdown()
  }

注意點

1.最后需要executors.shutdown()。

  • 如果是executors.shutdownNow()會發(fā)生未執(zhí)行完的task強(qiáng)制關(guān)閉線程。
  • 如果使用executors.awaitTermination()則會發(fā)生阻塞,不是我們想要的結(jié)果。
  • 如果沒有這個shutdowm操作,程序會正常執(zhí)行,但是長時間會產(chǎn)生大量無用的線程池,因為每次foreachRDD都會創(chuàng)建一個線程池。

2.可不可以將創(chuàng)建線程池放到foreachRDD外面?

不可以,這個關(guān)系到對于scala閉包到理解,經(jīng)測試,第一次或者前幾次batch是正常的,后面的batch無線程可用。

3.線程池executor崩潰了就會導(dǎo)致數(shù)據(jù)丟失

原則上是這樣的,但是正常的代碼一般不會發(fā)生executor崩潰。至少我在使用的時候沒遇到過。

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

相關(guān)文章

  • Spring boot @RequestBody數(shù)據(jù)傳遞過程詳解

    Spring boot @RequestBody數(shù)據(jù)傳遞過程詳解

    這篇文章主要介紹了Spring boot @RequestBody數(shù)據(jù)傳遞過程詳解,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2019-12-12
  • JsonFormat與@DateTimeFormat注解實例解析

    JsonFormat與@DateTimeFormat注解實例解析

    這篇文章主要介紹了JsonFormat與@DateTimeFormat注解實例解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2019-12-12
  • 詳解阿里云maven鏡像庫配置(gradle,maven)

    詳解阿里云maven鏡像庫配置(gradle,maven)

    這篇文章主要介紹了詳解阿里云maven鏡像庫配置(gradle,maven),小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-02-02
  • IntelliJ IDEA2021.1 配置大全(超詳細(xì)教程)

    IntelliJ IDEA2021.1 配置大全(超詳細(xì)教程)

    這篇文章主要介紹了IntelliJ IDEA2021.1 配置大全(超詳細(xì)教程),需要的朋友可以參考下
    2021-04-04
  • SpringBoot中時間格式化的五種方法匯總

    SpringBoot中時間格式化的五種方法匯總

    時間格式化在項目中使用頻率是非常高的,這篇文章主要給大家介紹了關(guān)于SpringBoot中時間格式化的五種方法,文中通過示例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2021-07-07
  • springmvc接口接收參數(shù)與請求參數(shù)格式的整理

    springmvc接口接收參數(shù)與請求參數(shù)格式的整理

    這篇文章主要介紹了springmvc接口接收參數(shù)與請求參數(shù)格式的整理,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-11-11
  • 利用EasyPOI實現(xiàn)多sheet和列數(shù)的動態(tài)生成

    利用EasyPOI實現(xiàn)多sheet和列數(shù)的動態(tài)生成

    EasyPoi功能如同名字,主打的功能就是容易,讓一個沒見接觸過poi的人員就可以方便的寫出Excel導(dǎo)出,Excel導(dǎo)入等功能,本文主要來講講如何利用EasyPOI實現(xiàn)多sheet和列數(shù)的動態(tài)生成,需要的可以了解下
    2025-03-03
  • Java開發(fā)神器Lombok安裝與使用詳解

    Java開發(fā)神器Lombok安裝與使用詳解

    Lombok的安裝分兩部分:Idea插件的安裝和maven中pom文件的導(dǎo)入,本文重點給大家介紹Java開發(fā)神器Lombok安裝與使用詳解,感興趣的朋友跟隨小編一起看看吧
    2022-02-02
  • Spring?Boot實現(xiàn)微信掃碼登錄功能流程分析

    Spring?Boot實現(xiàn)微信掃碼登錄功能流程分析

    這篇文章主要介紹了Spring?Boot?實現(xiàn)微信掃碼登錄功能,介紹了授權(quán)流程代碼和用戶登錄和登出的操作代碼,代碼簡單易懂,對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2022-04-04
  • Java 中 getClass() 方法的使用與原理深入分析(對象類型信息)

    Java 中 getClass() 方法的使用與原理深入分析(對象類型信息)

    在 Java 編程中,getClass() 是一個非常重要的方法,它用于獲取對象的運行時類信息,無論是調(diào)試代碼、反射操作,還是類型檢查,getClass() 都扮演著關(guān)鍵角色,本文將深入探討 getClass() 的使用方法、底層原理以及實際應(yīng)用場景,感興趣的朋友一起看看吧
    2024-12-12

最新評論

博白县| 南靖县| 奉化市| 晋州市| 涿鹿县| 蓝山县| 通州区| 且末县| 徐汇区| 丹巴县| 临西县| 南京市| 太白县| 仪陇县| 板桥市| 松溪县| 北碚区| 子长县| 巴中市| 保定市| 扶余县| 西青区| 黄龙县| 宜阳县| 衢州市| 天门市| 邹城市| 厦门市| 富平县| 成都市| 大冶市| 富锦市| 安义县| 蛟河市| 乡城县| 湖口县| 中宁县| 黔西县| 郯城县| 高邑县| 光山县|