【问题标题】:Why does this python multiprocessing script slow down after a while?为什么这个 python 多处理脚本会在一段时间后变慢?
【发布时间】:2013-12-21 09:31:29
【问题描述】:

script from this answer 的基础上,我有以下场景:一个包含 2500 个大文本文件(每个约 55Mb)的文件夹,所有文件用制表符分隔。基本上是网络日志。

我需要对每个文件的每一行中的第二个“列”进行 md5 哈希处理,将修改后的文件保存在其他地方。源文件位于机械磁盘上​​,目标文件位于 SSD 上。

脚本处理前 25 个(左右)文件的速度非常快。然后它会减慢速度。根据前 25 个文件,它应该在 2 分钟左右完成所有文件。但是,根据之后的表现,全部完成需要 15 分钟左右。

它在具有 32 Gb RAM 的服务器上运行,并且任务管理器很少显示超过 6 Gb 的正在使用。我将它设置为启动 6 个进程,但内核上的 CPU 使用率很低,很少超过 15%。

为什么会变慢?磁盘读/写问题?垃圾收集器?代码不好?关于如何加快速度的任何想法?

这是脚本

import os

import multiprocessing
from multiprocessing import Process
import threading
import hashlib

class ThreadRunner(threading.Thread):
    """ This class represents a single instance of a running thread"""
    def __init__(self, fileset, filedirectory):
        threading.Thread.__init__(self)
        self.files_to_process = fileset
        self.filedir          = filedirectory

    def run(self):
        for current_file in self.files_to_process:

            # Open the current file as read only
            active_file_name = self.filedir + "/" + current_file
            output_file_name = "D:/hashed_data/" + "hashed_" + current_file

            active_file = open(active_file_name, "r")
            output_file = open(output_file_name, "ab+")

            for line in active_file:
                # Load the line, hash the username, save the line
                lineList = line.split("\t")

                if not lineList[1] == "-":
                    lineList[1] = hashlib.md5(lineList[1]).hexdigest()

                lineOut = '\t'.join(lineList)
                output_file.write(lineOut)

            # Always close files after you open them
            active_file.close()
            output_file.close()

            print "\nCompleted " + current_file

class ProcessRunner:
    """ This class represents a single instance of a running process """
    def runp(self, pid, numThreads, fileset, filedirectory):
        mythreads = []
        for tid in range(numThreads):
            th = ThreadRunner(fileset, filedirectory)
            mythreads.append(th) 
        for i in mythreads:
            i.start()
        for i in mythreads:
            i.join()

class ParallelExtractor:    
    def runInParallel(self, numProcesses, numThreads, filedirectory):
        myprocs = []
        prunner = ProcessRunner()

        # Store the file names from that directory in a list that we can iterate
        file_names = os.listdir(filedirectory)

        file_sets = []
        for i in range(numProcesses):
            file_sets.append([])

        for index, name in enumerate(file_names):
            num = index % numProcesses
            file_sets[num].append(name)


        for pid in range(numProcesses):
            pr = Process(target=prunner.runp, args=(pid, numThreads, file_sets[pid], filedirectory)) 
            myprocs.append(pr) 
        for i in myprocs:
            i.start()

        for i in myprocs:
            i.join()

if __name__ == '__main__':    

    file_directory = "E:/original_data"

    processes = 6
    threads   = 1

    extractor = ParallelExtractor()
    extractor.runInParallel(numProcesses=processes, numThreads=threads, filedirectory=file_directory)

【问题讨论】:

  • 您可能会获得性能提升,因为操作系统会将第一个文件缓存在内存中,因此不会发生磁盘 I/O。您可以通过重新启动服务器轻松检查这一点,并查看处理速度是否减慢。如果您无法重新启动,您应该通过从磁盘读取足够的文件来填充物理内存来填充缓存。如果您具有本地访问权限,则可以简单地仔细聆听磁盘搜索。值得一提的是,对文件执行散列肯定是受磁盘约束的,而不是 CPU,因此在最好的情况下并行执行它是无用的
  • 此外,如果您的源文件位于机械磁盘上​​,那么同时读取它们的 6 个进程可能会大大降低您的速度,尤其是在 I/O 调度非常糟糕的 Windows 上。将源文件移动到 SSD 会发生什么?
  • 实际上,如果 I/O 调度是问题(如果您的 CPU 使用率一直很低,这很可能),您应该通过将 numProcesses 降低到 1 来提高性能。
  • @Max Noel 我相信源文件(实际上是编译后的字节码)只会被读取一次并保存在磁盘缓存和内存映射文件中(除非字节码真的,真的 大)
  • @MaxNoel 将进程数减少到 1 肯定会加快速度。我想这最终会受到硬盘上有多少读/写磁头的限制?

标签: python performance multiprocessing


【解决方案1】:

散列是一项相对简单的任务,与旋转磁盘的速度相比,现代 CPU 的速度非常快。 i7 上的快速基准测试表明,它使用 MD5 可以散列大约 450 MB/s,或者使用 SHA-1 可以散列 290 MB/s。相比之下,旋转磁盘的典型(顺序原始读取)速度约为 70-150 MB/s。这意味着,即使忽略文件系统的开销和最终的磁盘寻道,CPU 散列文件的速度也比磁盘读取它的速度快 3 倍。

处理第一个文件时获得的性能提升可能是因为操作系统将第一个文件缓存在内存中,因此不会发生磁盘 I/O。这可以通过以下任一方式确认:

  • 重启服务器,从而刷新缓存
  • 通过从磁盘读取足够大的文件来填充缓存
  • 在处理第一个文件时仔细聆听是否没有磁盘寻道

现在,由于散列文件的性能瓶颈是磁盘,因此在多个进程或线程中执行散列是没有用的,因为它们都将使用同一个磁盘。正如@Max Noel 所提到的,它实际上会降低性能,因为您将并行读取多个文件,因此您的磁盘必须在文件之间进行查找。正如他所提到的,性能也会因您使用的操作系统的 I/O 调度程序而异。

现在,如果您仍在生成数据,您有一些可能的解决方案:

  • 按照@Max Noel 的建议,使用更快的磁盘或 SSD。
  • 从多个磁盘读取 - 在不同的文件系统中或在 RAID 上的单个文件系统中
  • 将任务拆分到多台机器上(每台机器有一个或多个磁盘)

但是,如果您只想散列这 2500 个文件并且您已经将它们放在一个磁盘上,那么这些解决方案将毫无用处。将它们从磁盘读取到其他磁盘然后执行散列更慢,因为您将读取文件两次,并且可以尽可能快地进行散列他们。

最后,根据@yaccz 的想法,如果您安装了findxargsmd5sum 的cygwin 二进制文件,我想您可以避免编写程序来执行散列的麻烦。

【讨论】:

    【解决方案2】:

    既然可以让事情变得复杂,为什么还要简单呢?

    通过 smbfs 或其他方式在 linux 主机上安装驱动器并运行

    #! /bin/sh
    
    SRC="" # FIXME
    DST="" # FIXME
    
    convert_line() {
        new_line=`echo $i | cut -f 1 -d "\t"`
        f2=`echo $i | cut -f 2 -d "\t"`
        frest=`echo $i | cut -f 1,2 --complement -d "\t"`
    
        if [ ! "x${f2}" = "-" ] ; then
            f2=`echo "${f2}" | md5sum | head -c-1`
            # might wanna throw in some memoization
        fi
    
        echo "${new_line}\t$f2\t${frest}"
    }
    
    convert_file() {
        for i in `cat $1`; do
            convert_line "${i}" >> $DST/hashed-$1
        done
    }
    
    for i in $SRC/*; do
        convert_file $i
    done
    

    未测试。可能需要打磨一些粗糙的边缘。

    【讨论】:

    • 好主意,但是服务器被锁定了,管理员不可能让我这样做。
    猜你喜欢
    • 1970-01-01
    • 2018-10-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-01-07
    • 2017-12-05
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多