Python技术迷

如何在Python中实现并发的MapReduce?

那天中午在公司楼下抽烟,旁边那哥几个还在吐槽我们后端的任务调度慢得要命,说那个MapReduce跑一上午都没出结果,我一听就乐了,我说你们这肯定是串行跑的吧。他们说“是多线程啊”,我说多线程不代表并发跑得快啊你得看你怎么设计的——就像MapReduce这种模式,用不好真就是灾难现场。

MapReduce是个啥?为什么你总听说它但从来没人写出来给你看?

我跟你说,MapReduce其实不是啥高深玩意儿,咱们用Python实现起来也很简单。核心两个动作你记住:一个是 map —— 把每条数据处理一下,另一个是 reduce —— 把处理过的结果合并起来。

比如说你有个一堆文本文件,你想统计里面每个单词出现了几次。这时候你就可以这么干:

  1. 用 map 把每个文件拆成一个个单词 (word, 1);
  2. 然后用 reduce 把这些 (word, 1) 按照 word 聚合起来。

这就是 MapReduce 最基础的模型了。

但为啥慢?问题就出在你没有“并发”!

你想啊,如果你一共10个文件,结果你就一个线程一个一个读,然后map完之后又一条条去reduce,这不就白瞎了你那8核CPU吗?

所以我那天就在会议室临时撸了个并发版的MapReduce,用的是 Python 里的 concurrent.futures 这个模块,哥几个当场看了都说“卧槽怎么就这么丝滑”。

下面是我给他们演示的代码,写得很简洁,不搞那些装模作样的架构:

import os
from collections import Counter
from concurrent.futures import ThreadPoolExecutor, as_completed

defmap_function(filepath):
with open(filepath, 'r', encoding='utf-8') as f:
        words = f.read().split()
return Counter(words)

defreduce_function(results):
    final_count = Counter()
for result in results:
        final_count.update(result)
return final_count

defconcurrent_mapreduce(file_list, max_workers=4):
    mapped_results = []
with ThreadPoolExecutor(max_workers=max_workers) as executor:
        futures = [executor.submit(map_function, file) for file in file_list]
for future in as_completed(futures):
            mapped_results.append(future.result())
return reduce_function(mapped_results)

用法也简单,我那天直接拖了几个 .txt 文件进去,打印:

if __name__ == '__main__':
    files = [f'test_data/file_{i}.txt'for i in range(5)]  # 随便几个测试文件
    result = concurrent_mapreduce(files, max_workers=5)
    print(result.most_common(10))  # 看最常见的10个词

那效果是真的快,我那台老破笔记本都能在2秒钟内跑完,之前串行得跑个十几秒,差距一下就出来了。

注意点在哪儿?坑也不少

那天我们组的小王就问我:“那如果数据超大呢?文件十几G那种”,我说你这就不能用线程了,要用进程池——GIL你懂吧?Python多线程处理CPU密集型任务根本不行,得用 multiprocessing 才能榨干CPU性能。只不过进程池的IO开销会大一点,你得权衡。

还有就是 map 阶段不能太耗时,否则线程池里一个任务卡住,整个 pipeline 也就慢了。我建议是 map 阶段尽量轻量,IO操作要控制好,最好异步或者缓冲读。

如果你数据不是文件,而是数据库或者消息队列?

那就换一下 map_function 就好了啊,把 filepath 参数换成一个行数据就行。你完全可以把队列拿来的数据扔进线程池里 map,然后结果扔到某个共享队列里再 reduce。

我甚至跟你说,如果你想搞得再硬核点,可以用 asyncio 来做 Map,那是真正的非阻塞并发了。不过那种写法你得脑子清楚,不然容易把自己绕进去。

总之,MapReduce 并不神秘,你用好 Python 自带的这些并发工具,轻轻松松就能写出跑得飞快的版本。不要被“MapReduce”这仨字吓住,实际上就是“并发 for 循环 + 聚合结果”这么个事儿。

-END-

我为大家打造了一份RPA教程,完全免费:songshuhezi.com/rpa.html

🔥虎哥私藏精品🔥

虎哥作为一名老码农,整理了全网最全《python高级架构师资料合集》,总量高达650GB,点击下方公众号回复关键字 python 全部免费领