【问题标题】:Read large file in parallel?并行读取大文件?
【发布时间】:2013-08-07 13:19:01
【问题描述】:

我有一个大文件,我需要从中读取并制作字典。我希望这尽可能快。但是我在 python 中的代码太慢了。这是一个显示问题的最小示例。

先做一些假数据

paste <(seq 20000000) <(seq 2 20000001)  > largefile.txt

现在这是一段最小的 Python 代码,用于读取它并制作字典。

import sys
from collections import defaultdict
fin = open(sys.argv[1])

dict = defaultdict(list)

for line in fin:
    parts = line.split()
    dict[parts[0]].append(parts[1])

时间安排:

time ./read.py largefile.txt
real    0m55.746s

但是可以更快地读取整个文件:

time cut -f1 largefile.txt > /dev/null    
real    0m1.702s

我的 CPU 有 8 个核心,是否可以将这个程序并行化 python来加快速度?

一种可能是读取输入的大块,然后在不同的非重叠子块上并行运行 8 个进程,从内存中的数据并行创建字典,然后读取另一个大块。这在 python 中是否可能以某种方式使用多处理?

更新。假数据不是很好,因为每个键只有一个值。更好的是

perl -E 'say int rand 1e7, $", int rand 1e4 for 1 .. 1e7' > largefile.txt

(与Read in large file and make dictionary有关。)

【问题讨论】:

  • 这个列表可能会拖慢你的速度。
  • @sihrc 根据分析,大部分时间都花在 dict[parts[0]].append(parts[1]) 中,但 append 确实需要大约 10% 的时间。请注意,在这个假案例中,列表的长度只有 1,但总的来说我希望它们更长。
  • 这就是我的意思。那 .append 部分正在减慢您的速度。您是否想过附加大型列表的替代方法?您要附加什么样的数据?
  • 您不能仅使用一个磁盘并行化 I/O。如果您有很多 RAM,您可以使用 readlines() 将整个文件读入内存并直接在 RAM 中处理它们。
  • 是的,但是读取文件并不是慢的部分。就是字典中的追加列表,可以并行化

标签: python performance multiprocessing


【解决方案1】:

也许可以将其并行化以加快速度,但并行执行多个读取不太可能有帮助。

您的操作系统不太可能有效地并行执行多次读取(例外情况是像条带化 RAID 阵列,在这种情况下,您仍然需要知道步幅才能充分利用它)。

您可以做的是在读取的同时运行相对昂贵的字符串/字典/列表操作。

因此,一个线程读取(大)块并将其推送到同步队列,一个或多个消费者线程从队列中拉取块,将它们分成行,然后填充字典。

(如果您使用多个消费者线程,正如 Pappnese 所说,请为每个线程构建一个字典,然后加入它们)。


提示:


回复。赏金:

C 显然没有 GIL 可抗衡,因此多个消费者可能会更好地扩展。读取行为并没有改变。不利的一面是 C 缺乏对哈希映射(假设您仍然需要 Python 样式的字典)和同步队列的内置支持,因此您必须找到合适的组件或编写自己的组件。 多个消费者各自构建自己的字典然后最后合并它们的基本策略仍然可能是最好的。

使用strtok_r 而不是str.split 可能更快,但请记住,您还需要手动管理所有字符串的内存。哦,你也需要逻辑来管理行片段。老实说,C 为您提供了很多选择,我认为您只需要对其进行分析并查看即可。

【讨论】:

  • 我已经给出了一些提示来帮助您入门。写一些代码,问你有没有问题。
  • 多线程可以工作吗?我很担心 GIL。
  • 生产者线程是 I/O 绑定的,所以如果 GIL 在 read 系统调用的线程阻塞时被释放,一个消费者线程可以进行计算哈希和更新它的字典。不过,多个消费者线程很可能会争夺 GIL。
  • 这个对类似问题的回答可能会对您有所帮助:stackoverflow.com/a/43079667/196732 我不确定连接是否会在另一端一起产生...?
【解决方案2】:

认为使用处理池可以解决这样的问题似乎很诱人,但最终会比这更复杂一些,至少在纯 Python 中是这样。

因为 OP 提到每个输入行上的列表实际上会比两个元素长,所以我使用以下方法制作了一个更真实的输入文件:

paste <(seq 20000000) <(seq 2 20000001) <(seq 3 20000002) |
  head -1000000 > largefile.txt

在分析原始代码后,我发现该过程中最慢的部分是行分割例程。 (.split() 在我的机器上花费的时间大约是 .append() 的 2 倍。)

1000000    0.333    0.000    0.333    0.000 {method 'split' of 'str' objects}
1000000    0.154    0.000    0.154    0.000 {method 'append' of 'list' objects}

所以我将拆分分解为另一个函数,并使用池来分配拆分字段的工作:

import sys
import collections
import multiprocessing as mp

d = collections.defaultdict(list)

def split(l):
    return l.split()

pool = mp.Pool(processes=4)
for keys in pool.map(split, open(sys.argv[1])):
    d[keys[0]].append(keys[1:])

不幸的是,添加池会使速度降低 2 倍以上。原来的版本是这样的:

$ time python process.py smallfile.txt 
real    0m7.170s
user    0m6.884s
sys     0m0.260s

与并行版本:

$ time python process-mp.py smallfile.txt 
real    0m16.655s
user    0m24.688s
sys     0m1.380s

因为.map()调用基本上要序列化(pickle)每个输入,将其发送到远程进程,然后反序列化(unpickle)来自远程进程的返回值,这种方式使用池要慢得多。通过向池中添加更多内核确实可以获得一些改进,但我认为这从根本上是分配这项工作的错误方式。

要真正加快跨内核的速度,我的猜测是您需要使用某种固定块大小来读取大块输入。然后你可以将整个块发送到一个工作进程并取回序列化列表(尽管仍然不知道这里的反序列化会花费你多少)。以固定大小的块读取输入听起来对于预期的输入可能很棘手,但是,因为我的猜测是每行不一定是相同的长度。

【讨论】:

  • 我很惊讶字典创建对你来说并不是最慢的部分,因为它在我的分析中。关于行长我想我不清楚。每条线的长度不会完全相同,但会分成相同数量的部分。我刚刚添加的新假数据示例更加真实。
  • pool.map(split,open(sys.argv[1]),10000) 将按 10000 长度的块对行进行分组,但执行时间与非分块版本没有区别。
【解决方案3】:

几年前,Tim Bray 的网站 [1] 上有一篇关于此的博文系列“Wide Finder Project”。您可以找到 ElementTree [3] 和 PIL [4] 的 Fredrik Lundh 的解决方案 [2]。我知道在这个网站上发布链接通常是不鼓励的,但我认为这些链接比复制粘贴他的代码给你更好的答案。

[1]http://www.tbray.org/ongoing/When/200x/2007/10/30/WF-Results
[2]http://effbot.org/zone/wide-finder.htm
[3]http://docs.python.org/3/library/xml.etree.elementtree.html
[4]http://www.pythonware.com/products/pil/

【讨论】:

【解决方案4】:

您可以尝试的一件事是从文件中获取行数,然后生成 8 个线程,每个线程从文件的 1/8 生成一个字典,然后在所有线程完成后加入字典。如果附加需要时间而不是读取行,这可能会加快速度。

【讨论】:

  • 由于这是一个 CPU 密集型问题,CPython 上的threading isn't going to help much。正常的解决方案是改用多处理,但这不允许在进程之间共享 python dict。因此,必须有一个具有添加任务的 dict 进程,以及其他用于读取和解析的进程,这些进程使用队列或管道沿其结果发送。由于开销,我认为使用此方法无法实现超过 3 倍的加速。
  • @LauritzV.Thaulow 将整个内容加载到内存中然后构建 8 个字典然后合并它们怎么样?这听起来合理吗?
  • @Anush 好吧,当从一个进程移动到另一个进程时,必须对字典进行序列化和重建,并且由于重建几乎需要与构建一样多的工作,因此它的大部分目的都失败了。然后必须合并重建的dicts,这里我的python知识不足,因为我不知道d1.update(d2)是否可以做一些快捷方式,使其比循环d2并一一添加键值对快得多.无论如何,我怀疑它会有多大帮助。我知道 IronPython 没有 CPU 绑定的线程问题,所以如果你可以使用它,那么 Useless 的答案应该可以工作。
【解决方案5】:

对于慢速字典追加的更多基本解决方案:用字符串对数组替换字典。填写然后排序。

【讨论】:

    【解决方案6】:

    如果您存档的数据不经常更改,您可以选择对其进行序列化。 Python 解释器将更快地反序列化它。 您可以使用 cPickle 模块。

    或者创建 8 个单独的进程是另一种选择。 因为,拥有唯一的 dict 使它更有可能。 您可以通过“多处理”模块或“套接字”模块中的管道在这些进程之间进行交互。

    最好的问候

    BarışÇUHADAR.

    【讨论】:

      猜你喜欢
      • 2020-07-20
      • 1970-01-01
      • 2021-07-20
      • 1970-01-01
      • 2018-07-21
      • 1970-01-01
      • 1970-01-01
      • 2022-06-24
      相关资源
      最近更新 更多