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

使用golang實現(xiàn)一個MapReduce的示例代碼

 更新時間:2023年09月21日 10:57:26   作者:寫代碼的lorre  
這篇文章主要給大家介紹了關于如何使用golang實現(xiàn)一個MapReduce,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧

背景

在日常業(yè)務開發(fā)中,我們經(jīng)常遇到需要并發(fā)處理的場景。例如:

  • 依據(jù)id列表查詢db,獲取數(shù)據(jù)。為了保證查詢性能,單次查詢的id列表長度最好不要超過50(依據(jù)業(yè)務來判斷),當id列表長度超過50時,拆分成并發(fā)請求,減少耗時和提高性能,返回聚合后的結果
  • 外部提供的接口不支持批量寫入/讀取數(shù)據(jù),當需要批量處理數(shù)據(jù)時,為了減少耗時和提高性能,并發(fā)請求外部接口

以上處理數(shù)據(jù)的場景,都可以分成兩個階段:

  • 請求階段。基本都是IO操作,請求db,或者是調用外部接口
  • 處理階段。對返回的數(shù)據(jù)進行轉換,過濾,聚合等操作

同步調用,調用耗時增長明顯

并發(fā)調用,可以減少調用耗時

分析

上面說的處理數(shù)據(jù)的場景,都可以分成兩個階段:

  • 請求階段。IO操作,可以并發(fā)的去進行,互不干擾
  • 處理階段。同步進行,保證聚合結果的正確性

這種是一種特殊的MapReduce

為了處理這類場景,我們需要明確以下幾個部分:

  • 列表長度。代表有多少數(shù)據(jù)需要進行處理
  • map函數(shù)。并發(fā)處理的函數(shù),互不干擾
  • reduce函數(shù)。同步處理的函數(shù)
  • 最大并發(fā)數(shù)。決定需要開多少線程/協(xié)程來處理
  • 拆分長度。列表長度 / 拆分長度 = 子任務數(shù)

由于我在日常開發(fā)中常使用golang語言,下面梳理下使用golang來解決這類問題的一個思路

函數(shù)簽名

func ChunkProcess(length int, procedure func(start, end int) (interface{}, error),
   reduce func(partialResult interface{}, partialErr error, start, end int), maxConcurrent int, chunkSize int) 

核心邏輯:

  • 當最大并發(fā)數(shù) <= 1 或者子任務數(shù)(列表長度 / 拆分長度) <= 1時,同步執(zhí)行map函數(shù)和reduce函數(shù)即可

  • 其余情況,并發(fā)處理map函數(shù),同步執(zhí)行reduce函數(shù)

    • 獲取并發(fā)處理的子任務數(shù)量:lengthTask := int(math.Ceil(float64(length) / float64(chunkSize)))
    • 通過sync.Mutex保證reduce同步執(zhí)行
    • 通過sync.WaitGroup保證等待子任務全部執(zhí)行完成
    • 通過chan控制最大并發(fā)數(shù)

代碼實現(xiàn)

package test
import (
   "math"
   "sync"
)
func ChunkProcess(length int, procedure func(start, end int) (interface{}, error),
   reduce func(partialResult interface{}, partialErr error, start, end int), maxConcurrent int, chunkSize int) {
   if length < 1 {
      return
   }
   if maxConcurrent <= 1 || length <= chunkSize {
      doChunkProcessSerially(length, procedure, reduce, chunkSize)
   } else {
      doChunkProcessConcurrently(length, procedure, reduce, maxConcurrent, chunkSize)
   }
}
// 同步處理
func doChunkProcessSerially(length int, procedure func(start, end int) (interface{}, error),
   reduce func(partialResult interface{}, partialErr error, start, end int), chunkSize int) {
   // 拆分的子任務數(shù)
   chunkNums := int(math.Ceil(float64(length) / float64(chunkSize)))
   for i := 0; i < chunkNums; i++ {
      func(chunkIndex int) {
         defer func() {
            if err := recover(); err != nil {
               // 自定義錯誤處理
            }
         }()
         start := chunkIndex * chunkSize
         end := start + chunkSize
         if end > length {
            end = length
         }
         // 執(zhí)行map
         response, err := procedure(start, end)
         // 執(zhí)行reduce
         if reduce != nil {
            reduce(response, err, start, end)
         }
      }(i)
   }
}
// 并發(fā)處理
func doChunkProcessConcurrently(length int, procedure func(start, end int) (interface{}, error),
   reduce func(partialResult interface{}, partialErr error, start, end int), maxConcurrent int, chunkSize int) {
   index := 0
   chunkIndex := 0
   // 拆分的子任務數(shù)
   lengthTask := int(math.Ceil(float64(length) / float64(chunkSize)))
   // 保證reduce同步執(zhí)行
   var lock sync.Mutex
   // 保證子任務全部執(zhí)行完成
   var wg sync.WaitGroup
   wg.Add(lengthTask)
   // 控制并發(fā)數(shù)
   throttleChan := make(chan struct{}, maxConcurrent)
   for {
      start := index
      end := index + chunkSize
      if end > length {
         end = length
      }
      throttleChan <- struct{}{}
      go func(chunkIndex int) {
         defer func() {
            <-throttleChan
            if err := recover(); err != nil {
               // 自定義錯誤處理
            }
            wg.Done()
         }()
         // 執(zhí)行map
         response, err := procedure(start, end)
         // 執(zhí)行reduce
         if reduce != nil {
            lock.Lock()
            defer lock.Unlock()
            reduce(response, err, start, end)
         }
      }(chunkIndex)
      chunkIndex++
      index = index + chunkSize
      if index >= length {
         break
      }
   }
   wg.Wait()
   close(throttleChan)
}

測試:

func TestChunkProcess(t *testing.T) {
   trackIDs := []int64{123, 456, 789}
   results := make([]int64, 0)
   ChunkProcess(len(trackIDs), func(start, end int) (interface{}, error) {
      result := trackIDs[start] + 100
      return result, nil
   }, func(partialResult interface{}, partialErr error, start, end int) {
      results = append(results, partialResult.(int64))
   }, 2, 1)
   fmt.Println(results)
}

總結

多對業(yè)務場景進行抽象分析,為這一類場景提供解決方案

以上就是使用golang實現(xiàn)一個MapReduce的詳細內(nèi)容,更多關于golang實現(xiàn)MapReduce的資料請關注腳本之家其它相關文章!

相關文章

  • 初學Go必備的vscode插件及最常用快捷鍵和代碼自動補全

    初學Go必備的vscode插件及最常用快捷鍵和代碼自動補全

    這篇文章主要給大家介紹了關于初學vscode寫Go必備的vscode插件及最常用快捷鍵和代碼自動補全的相關資料,由于vscode是開源免費的,而且開發(fā)支持vscode的插件相對比較容易,更新速度也很快,需要的朋友可以參考下
    2023-07-07
  • Go語言執(zhí)行系統(tǒng)命令行命令的方法

    Go語言執(zhí)行系統(tǒng)命令行命令的方法

    這篇文章主要介紹了Go語言執(zhí)行系統(tǒng)命令行命令的方法,實例分析了Go語言操作系統(tǒng)命令行的技巧,具有一定參考借鑒價值,需要的朋友可以參考下
    2015-02-02
  • Go使用fmt包輸出與格式化核心庫的完整指南

    Go使用fmt包輸出與格式化核心庫的完整指南

    在 Go 語言中,fmt 是最基礎也是使用頻率最高的標準庫之一,幾乎每一個 Go 程序都會用到它,下面小編就和大家詳細介紹一下fmt包的具體使用吧
    2026-03-03
  • 詳解為什么說Golang中的字符串類型不能修改

    詳解為什么說Golang中的字符串類型不能修改

    在接觸Go這么語言,可能你經(jīng)常會聽到這樣一句話。對于字符串不能修改,可能你很納悶,日常開發(fā)中我們對字符串進行修改也是很正常的,為什么又說Go中的字符串不能進行修改呢?本文就來通過實際案例給大家演示一下
    2023-03-03
  • 一文帶你了解Golang中的泛型

    一文帶你了解Golang中的泛型

    泛型是一種可以編寫獨立于使用的特定類型的代碼的方法,可以通過編寫函數(shù)或類型來使用一組類型中的任何一個,下面就來和大家聊聊Golang中泛型的使用吧
    2023-07-07
  • 使用Go實現(xiàn)健壯的內(nèi)存型緩存的方法

    使用Go實現(xiàn)健壯的內(nèi)存型緩存的方法

    這篇文章主要介紹了使用Go實現(xiàn)健壯的內(nèi)存型緩存,本文比較了字節(jié)緩存和結構體緩存的優(yōu)劣勢,介紹了緩存穿透、緩存錯誤、緩存預熱、緩存?zhèn)鬏敗⒐收限D移、緩存淘汰等問題,并對一些常見的緩存庫進行了基準測試,需要的朋友可以參考下
    2022-05-05
  • Golang中Slice 底層機制的實現(xiàn)

    Golang中Slice 底層機制的實現(xiàn)

    本文主要介紹了Golang中Slice 底層機制的實現(xiàn),文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2026-03-03
  • 詳解Golang 中的并發(fā)限制與超時控制

    詳解Golang 中的并發(fā)限制與超時控制

    這篇文章主要介紹了詳解Golang 中的并發(fā)限制與超時控制,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-02-02
  • Go-Gin Web框架的實現(xiàn)示例

    Go-Gin Web框架的實現(xiàn)示例

    本文主要介紹了Go-Gin Web框架的實現(xiàn)示例,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2025-11-11
  • GoLand如何設置中文

    GoLand如何設置中文

    這篇文章主要介紹了GoLand如何設置中文,本文通過圖文并茂的形式給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-12-12

最新評論

榆社县| 鲁甸县| 扎鲁特旗| 桦川县| 公安县| 广灵县| 周宁县| 耒阳市| 广汉市| 闻喜县| 任丘市| 慈溪市| 安陆市| 永兴县| 周宁县| 疏勒县| 东阳市| 福建省| 芮城县| 腾冲县| 屏南县| 安乡县| 汤原县| 土默特左旗| 清徐县| 秦安县| 长治县| 佛坪县| 诏安县| 东城区| 宿松县| 东源县| 利川市| 佛学| 怀仁县| 黔西| 万安县| 芦溪县| 波密县| 定安县| 皮山县|