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

使用Python實(shí)現(xiàn)MapReduce的示例代碼

 更新時(shí)間:2024年05月06日 09:38:37   作者:changwanyi  
MapReduce是一個(gè)用于大規(guī)模數(shù)據(jù)處理的分布式計(jì)算模型,最初由Google工程師設(shè)計(jì)并實(shí)現(xiàn)的,Google已經(jīng)將完整的MapReduce論文公開(kāi)發(fā)布了,本文給大家介紹了使用Python實(shí)現(xiàn)MapReduce的示例代碼,需要的朋友可以參考下

一、MapReduce

將這個(gè)單詞分解為Map、Reduce。

  • Map階段:在這個(gè)階段,輸入數(shù)據(jù)集被分割成小塊,并由多個(gè)Map任務(wù)處理。每個(gè)Map任務(wù)將輸入數(shù)據(jù)映射為一系列(key, value)對(duì),并生成中間結(jié)果。

  • Reduce階段:在這個(gè)階段,中間結(jié)果被重新分組和排序,以便相同key的中間結(jié)果被傳遞到同一個(gè)Reduce任務(wù)。每個(gè)Reduce任務(wù)將具有相同key的中間結(jié)果合并、計(jì)算,并生成最終的輸出。

舉個(gè)例子,在一個(gè)很長(zhǎng)的字符串中統(tǒng)計(jì)某個(gè)字符出現(xiàn)的次數(shù)。

from collections import defaultdict
def mapper(word):
    return word, 1
 
def reducer(key_value_pair):
    key, values = key_value_pair
    return key, sum(values)
def map_reduce_function(input_list, mapper, reducer):
    '''
    - input_list: 字符列表
    - mapper: 映射函數(shù),將輸入列表中的每個(gè)元素映射到一個(gè)鍵值對(duì)
    - reducer: 聚合函數(shù),將映射結(jié)果中的每個(gè)鍵值對(duì)聚合到一個(gè)鍵值對(duì)
    - return: 聚合結(jié)果
    '''
    map_results = map(mapper, input_list)
    shuffler = defaultdict(list)
    for key, value in map_results:
        shuffler[key].append(value)
    return map(reducer, shuffler.items())
 
if __name__ == "__main__":
    words = "python best language".split(" ")
    result = list(map_reduce_function(words, mapper, reducer))
    print(result)

輸出結(jié)果為

[('python', 1), ('best', 1), ('language', 1)]

但是這里并沒(méi)有體現(xiàn)出MapReduce的特點(diǎn)。只是展示了MapReduce的運(yùn)行原理。

二、基于多線程實(shí)現(xiàn)MapReduce

from collections import defaultdict
import threading
 
class MapReduceThread(threading.Thread):
    def __init__(self, input_list, mapper, shuffler):
        super(MapReduceThread, self).__init__()
        self.input_list = input_list
        self.mapper = mapper
        self.shuffler = shuffler
 
    def run(self):
        map_results = map(self.mapper, self.input_list)
        for key, value in map_results:
            self.shuffler[key].append(value)
 
def reducer(key_value_pair):
    key, values = key_value_pair
    return key, sum(values)
def mapper(word):
    return word, 1
def map_reduce_function(input_list, num_threads):
    shuffler = defaultdict(list)
    threads = []
    chunk_size = len(input_list) // num_threads
    
    for i in range(0, len(input_list), chunk_size):
        chunk = input_list[i:i+chunk_size]
        thread = MapReduceThread(chunk, mapper, shuffler)
        thread.start()
        threads.append(thread)
    
    for thread in threads:
        thread.join()
 
    return map(reducer, shuffler.items())
 
if __name__ == "__main__":
    words = "python is the best language for programming and python is easy to learn".split(" ")
    result = list(map_reduce_function(words, num_threads=4))
    for i in result:
        print(i)

這里的本質(zhì)一模一樣,將字符串分割為四份,并且分發(fā)這四個(gè)字符串到不同的線程執(zhí)行,最后將執(zhí)行結(jié)果歸約。只不過(guò)由于Python的GIL機(jī)制,導(dǎo)致Python中的線程是并發(fā)執(zhí)行的,而不能實(shí)現(xiàn)并行,所以在python中使用線程來(lái)實(shí)現(xiàn)MapReduce是不合理的。(GIL機(jī)制:搶占式線程,不能在同一時(shí)間運(yùn)行多個(gè)線程)。

三、基于多進(jìn)程實(shí)現(xiàn)MapReduce

由于Python中GIL機(jī)制的存在,無(wú)法實(shí)現(xiàn)真正的并行。這里有兩種解決方案,一種是使用其他語(yǔ)言,例如C語(yǔ)言,這里我們不考慮;另一種就是利用多核,CPU的多任務(wù)處理能力。

from collections import defaultdict
import multiprocessing
 
def mapper(chunk):
    word_count = defaultdict(int)
    for word in chunk.split():
        word_count[word] += 1
    return word_count
 
def reducer(word_counts):
    result = defaultdict(int)
    for word_count in word_counts:
        for word, count in word_count.items():
            result[word] += count
    return result
 
def chunks(lst, n):
    for i in range(0, len(lst), n):
        yield lst[i:i + n]
 
def map_reduce_function(text, num_processes):
    chunk_size = (len(text) + num_processes - 1) // num_processes
    chunks_list = list(chunks(text, chunk_size))
 
    with multiprocessing.Pool(processes=num_processes) as pool:
        word_counts = pool.map(mapper, chunks_list)
 
    result = reducer(word_counts)
    return result
 
if __name__ == "__main__":
    text = "python is the best language for programming and python is easy to learn"
    num_processes = 4
    result = map_reduce_function(text, num_processes)
    for i in result:
        print(i, result[i])

這里使用多進(jìn)程來(lái)實(shí)現(xiàn)MapReduce,這里就是真正意義上的并行,依然是將數(shù)據(jù)切分,采用并行處理這些數(shù)據(jù),這樣才可以體現(xiàn)出MapReduce的高效特點(diǎn)。但是在這個(gè)例子中可能看不出來(lái)很大的差異,因?yàn)閿?shù)據(jù)量太小。在實(shí)際應(yīng)用中,如果數(shù)據(jù)集太小,是不適用的,可能無(wú)法帶來(lái)任何收益,甚至產(chǎn)生更大的開(kāi)銷(xiāo)導(dǎo)致性能的下降。

四、在100GB的文件中檢索數(shù)據(jù)

這里依然使用MapReduce的思想,但是有兩個(gè)問(wèn)題

讀取速度慢

解決方法:

使用分塊讀取,但是在分區(qū)時(shí)不宜過(guò)小。因?yàn)樵趧?chuàng)建分區(qū)時(shí)會(huì)被序列化到進(jìn)程,在進(jìn)程中又需要將其解開(kāi),這樣反復(fù)的序列化和反序列化會(huì)占用大量時(shí)間。不宜過(guò)大,因?yàn)檫@樣創(chuàng)建的進(jìn)程會(huì)變少,可能無(wú)法充分利用CPU的多核能力。

內(nèi)存消耗特別大

解決方法:

使用生成器和迭代器,但需獲取。例如分塊為8塊,生成器會(huì)一次讀取一塊的內(nèi)容并且返回對(duì)應(yīng)的迭代器,以此類(lèi)推,這樣就避免了讀取內(nèi)存過(guò)大的問(wèn)題。

from datetime import datetime
import multiprocessing
def chunked_file_reader(file_path:str, chunk_size:int):
    """
    生成器函數(shù):分塊讀取文件內(nèi)容
    - file_path: 文件路徑
    - chunk_size: 塊大小,默認(rèn)為1MB
    """
    with open(file_path, 'r', encoding='utf-8') as file:
        while True:
            chunk = file.read(chunk_size)
            if not chunk:
                break
            yield chunk
 
def search_in_chunk(chunk:str, keyword:str):
    """在文件塊中搜索關(guān)鍵字
    - chunk: 文件塊
    - keyword: 要搜索的關(guān)鍵字
    """
    lines = chunk.split('\n')
    for line in lines:
        if keyword in line:
            print(f"找到了:", line)
 
def search_in_file(file_path:str, keyword:str, chunk_size=1024*1024):
    """在文件中搜索關(guān)鍵字
    file_path: 文件路徑
    keyword: 要搜索的關(guān)鍵字
    chunk_size: 文件塊大小,為1MB
    """
    with multiprocessing.Pool() as pool:
        for chunk in chunked_file_reader(file_path, chunk_size):
            pool.apply_async(search_in_chunk, args=(chunk, keyword))
        
if __name__ == "__main__":
    start = datetime.now()
    file_path = "file.txt"
    keyword = "張三"
    search_in_file(file_path, keyword)
    end = datetime.now()
    print(f"搜索完成,耗時(shí) {end - start}")

最后程序運(yùn)行時(shí)間為兩分半左右。

到此這篇關(guān)于使用Python實(shí)現(xiàn)MapReduce的示例代碼的文章就介紹到這了,更多相關(guān)Python實(shí)現(xiàn)MapReduce內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Python中np.linalg.norm()用法實(shí)例總結(jié)

    Python中np.linalg.norm()用法實(shí)例總結(jié)

    在線性代數(shù)中一個(gè)向量通過(guò)矩陣轉(zhuǎn)換成另一個(gè)向量時(shí),原有向量的大小就是向量的范數(shù),這個(gè)變化過(guò)程的大小就是矩陣的范數(shù),下面這篇文章主要給大家介紹了關(guān)于Python中np.linalg.norm()用法的相關(guān)資料,需要的朋友可以參考下
    2022-07-07
  • Python中實(shí)現(xiàn)列表的逆序、復(fù)制與清除的幾種常見(jiàn)方法

    Python中實(shí)現(xiàn)列表的逆序、復(fù)制與清除的幾種常見(jiàn)方法

    本文介紹了Python中列表的逆序、復(fù)制和清除操作,通過(guò)reverse()方法、切片、copy()方法和clear()方法,我們可以輕松地對(duì)列表進(jìn)行這些操作
    2024-12-12
  • 使用Python的Flask框架表單插件Flask-WTF實(shí)現(xiàn)Web登錄驗(yàn)證

    使用Python的Flask框架表單插件Flask-WTF實(shí)現(xiàn)Web登錄驗(yàn)證

    Flask處理表單除了本身的WTForms包,使用Flask-WTF擴(kuò)展來(lái)增強(qiáng)表單功能也是很多開(kāi)發(fā)者的選擇,這里我們就來(lái)講解如何使用Python的Flask框架表單插件Flask-WTF實(shí)現(xiàn)Web登錄驗(yàn)證
    2016-07-07
  • 2021年最新用于圖像處理的Python庫(kù)總結(jié)

    2021年最新用于圖像處理的Python庫(kù)總結(jié)

    為了快速地處理大量信息,科學(xué)家需要利用圖像準(zhǔn)備工具來(lái)完成人工智能和深度學(xué)習(xí)任務(wù).在本文中,我將深入研究Python中最有用的圖像處理庫(kù),這些庫(kù)正在人工智能和深度學(xué)習(xí)任務(wù)中得到大力利用.我們開(kāi)始吧,需要的朋友可以參考下
    2021-06-06
  • python數(shù)據(jù)分析之DataFrame內(nèi)存優(yōu)化

    python數(shù)據(jù)分析之DataFrame內(nèi)存優(yōu)化

    pandas處理幾百兆的dataframe是沒(méi)有問(wèn)題的,但是我們?cè)谔幚韼讉€(gè)G甚至更大的數(shù)據(jù)時(shí),就會(huì)特別占用內(nèi)存,對(duì)內(nèi)存小的用戶(hù)特別不好,所以對(duì)數(shù)據(jù)進(jìn)行壓縮是很有必要的,本文就介紹了python DataFrame內(nèi)存優(yōu)化,感興趣的可以了解一下
    2021-07-07
  • Python使用selenium模擬鍵盤(pán)輸入的常見(jiàn)方法

    Python使用selenium模擬鍵盤(pán)輸入的常見(jiàn)方法

    本文主要介紹了在使用Selenium進(jìn)行模擬登錄時(shí),賬號(hào)和密碼輸入失敗的問(wèn)題,并提出了幾種解決方法,包括使用Selenium的send_keys函數(shù)、pyautogui庫(kù)、pynput庫(kù)和JavaScript輸入操作,并有相關(guān)的代碼示例供大家參考,需要的朋友可以參考下
    2025-04-04
  • 使用Python程序計(jì)算鋼琴88個(gè)鍵的音高

    使用Python程序計(jì)算鋼琴88個(gè)鍵的音高

    這篇文章介紹了使用Python程序計(jì)算鋼琴88個(gè)鍵的音高,文中通過(guò)示例代碼介紹的非常詳細(xì)。對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2022-04-04
  • 在Django中使用MQTT的方法

    在Django中使用MQTT的方法

    這篇文章主要介紹了在Django中使用MQTT的方法,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2021-05-05
  • Python運(yùn)行出現(xiàn)DeprecationWarning的問(wèn)題及解決

    Python運(yùn)行出現(xiàn)DeprecationWarning的問(wèn)題及解決

    這篇文章主要介紹了Python運(yùn)行出現(xiàn)DeprecationWarning的問(wèn)題及解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-07-07
  • PyQt界面阻塞卡死問(wèn)題的解決

    PyQt界面阻塞卡死問(wèn)題的解決

    當(dāng)用PyQt5開(kāi)發(fā)一個(gè)GUI界面 ,需要執(zhí)行業(yè)務(wù)邏輯時(shí),后臺(tái)邏輯執(zhí)行時(shí)間長(zhǎng),界面就容易出現(xiàn)卡死、未響應(yīng)等問(wèn)題,本文主要介紹了PyQt界面阻塞卡死問(wèn)題的解決
    2024-01-01

最新評(píng)論

喜德县| 应用必备| 永德县| 应用必备| 阿拉尔市| 西平县| 中阳县| 遂溪县| 建德市| 钟祥市| 淮滨县| 安乡县| 平邑县| 勐海县| 稻城县| 和田市| 色达县| 金乡县| 泰兴市| 东城区| 仁寿县| 新安县| 襄汾县| 堆龙德庆县| 南和县| 梁河县| 白城市| 东乡县| 高雄市| 工布江达县| 日照市| 通道| 那坡县| 厦门市| 宜兰市| 湖南省| 黄梅县| 仁化县| 长兴县| 康定县| 拜城县|