【问题标题】:Chunking data from a large file for multiprocessing?从大文件中分块数据以进行多处理?
【发布时间】:2012-01-03 18:52:55
【问题描述】:

我正在尝试使用多处理并行化应用程序 一个非常大的 csv 文件(64MB 到 500MB),逐行执行一些工作,然后输出一个固定大小的小文件 文件。

目前我做了一个list(file_obj),不幸的是它完全加载了 进入记忆(我想)然后我把这个列表分成 n 个部分,n 是 我要运行的进程数。然后我在分手时做一个pool.map() 列表。

与单个相比,这似乎有一个非常非常糟糕的运行时间 线程化,只需打开文件并迭代它的方法。有人可以 建议更好的解决方案?

此外,我需要按组处理文件的行,以保留 某列的值。这些行组本身可以拆分, 但任何组都不应包含此列的多个值。

【问题讨论】:

    标签: python parallel-processing


    【解决方案1】:

    fileobj 很大时,list(file_obj) 可能需要大量内存。我们可以通过使用itertools 来根据需要提取行块来减少内存需求。

    特别是,我们可以使用

    reader = csv.reader(f)
    chunks = itertools.groupby(reader, keyfunc)
    

    将文件分割成可处理的块,并且

    groups = [list(chunk) for key, chunk in itertools.islice(chunks, num_chunks)]
    result = pool.map(worker, groups)
    

    让多处理池一次处理num_chunks 块。

    通过这样做,我们只需要大约足够的内存来在内存中保存几个 (num_chunks) 块,而不是整个文件。


    import multiprocessing as mp
    import itertools
    import time
    import csv
    
    def worker(chunk):
        # `chunk` will be a list of CSV rows all with the same name column
        # replace this with your real computation
        # print(chunk)
        return len(chunk)  
    
    def keyfunc(row):
        # `row` is one row of the CSV file.
        # replace this with the name column.
        return row[0]
    
    def main():
        pool = mp.Pool()
        largefile = 'test.dat'
        num_chunks = 10
        results = []
        with open(largefile) as f:
            reader = csv.reader(f)
            chunks = itertools.groupby(reader, keyfunc)
            while True:
                # make a list of num_chunks chunks
                groups = [list(chunk) for key, chunk in
                          itertools.islice(chunks, num_chunks)]
                if groups:
                    result = pool.map(worker, groups)
                    results.extend(result)
                else:
                    break
        pool.close()
        pool.join()
        print(results)
    
    if __name__ == '__main__':
        main()
    

    【讨论】:

    • 当我说这些行不相关时,我撒了谎——在 csv 中,有一个列需要被分割(一个名称列,并且所有具有该名称的行都不能分开)。但是,我认为我可以根据此标准将其调整为分组。谢谢!我对 itertools 一无所知,现在我几乎一无所知。
    • 我的原始代码有错误。对pool.apply_async 的所有调用都是非阻塞的,因此整个文件会立即排队。这将导致没有内存节省。所以我稍微改变了循环,一次排队num_chunks。对pool.map 的调用是阻塞的,这将阻止整个文件一次排队。
    • @HappyLeapSecond 一个用户正在尝试在这里实现您的方法stackoverflow.com/questions/31164731/… 并且遇到了麻烦。或许你能帮忙?
    【解决方案2】:

    我会保持简单。让一个程序打开文件并逐行读取。您可以选择将其拆分为多少个文件,打开多少个输出文件,并将每一行写入下一个文件。这会将文件分成 n 个相等的部分。然后,您可以对每个文件并行运行 Python 程序。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-08-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-06-01
      • 1970-01-01
      相关资源
      最近更新 更多