【问题标题】:Python Chunking CSV File MultiproccessingPython 分块 CSV 文件多处理
【发布时间】:2015-09-18 19:28:56
【问题描述】:

我正在使用以下代码将 CSV 文件拆分为多个块(来自 here

def worker(chunk):
    print len(chunk)

def keyfunc(row):
    return row[0]

def main():
    pool = mp.Pool()
    largefile = 'Counseling.csv'
    num_chunks = 10
    start_time = time.time()
    results = []
    with open(largefile) as f:
        reader = csv.reader(f)
        reader.next()
        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()

但是,无论我选择使用多少块,块的数量似乎总是保持不变。例如,无论我选择 1 块还是 10 块,在处理示例文件时我总是会得到这个输出。理想情况下,我想对文件进行分块,以便公平分配。

请注意,我正在分块的真实文件超过 1300 万行,这就是我逐个处理它的原因。这是必须的!

6
7
1
...
1
1
94
--- 0.101687192917 seconds ---

【问题讨论】:

  • 假设您选择将文件分成 10 个块。您是希望一个工作进程处理文件的 1 块,还是希望将该 1 块均匀地分布在池中的工作器中,等到它们全部完成,然后将下一个块发送给池?
  • @HappyLeapSecond 每个工作进程 1 个块会更有效(所以我不必阻塞并等待所有其他进程也完成)在问这个问题之前,我查看了 Python文档相当广泛。我的理解是您正在使用 groupby 将一行中的每个值映射到一个键(相应的列)。这将返回一个迭代器。然后你将它传递给从 0 开始的 islice,然后取出 num_chunks(即 10)。这将是正确的行数?理想情况下,我想让流程处理 10,000 个行块。
  • 在另一个问题中,“有一列需要按...[分组],并且所有具有该名称的行都不能拆分”。这就是使用itertools.groupby 的原因。这里不需要按某列的值对行进行分组,所以我们可以跳过使用itertools.groupby

标签: python csv numpy multiprocessing python-multiprocessing


【解决方案1】:

首先,如果记录尚未在键列上排序,那么 itertools.groupby 将没有任何实际意义。 此外,如果您的要求只是将 csv 文件分块为预定数量的行并将其提供给 worker ,那么您不必执行所有这些操作。

一个简单的实现是:

import csv
from multiprocessing import Pool


def worker(chunk):
    print len(chunk)

def emit_chunks(chunk_size, file_path):
    lines_count = 0
    with open(file_path) as f:
        reader = csv.reader(f)
        chunk = []
        for line in reader:
            lines_count += 1
            chunk.append(line)
            if lines_count == chunk_size:
                lines_count = 0
                yield chunk
                chunk = []
            else:
                continue
        if chunk : yield chunk

def main():
    chunk_size = 10
    gen = emit_chunks(chunk_size, 'c:/Temp/in.csv')
    p = Pool(5)
    p.imap(worker, gen)
    print 'Completed..'

*编辑:改为 pool.imap 而不是 pool.map

【讨论】:

  • pool.imap 在内存方面不会更好,如果该列已排序,则调整 if lines_count == chunk_size 以确保要求特定列具有不同的值
  • @deinonychusaur 当然,pool.imap 是正确的方法,否则我们会遇到内存问题。我正在改变我的答案以使用它。谢谢。
  • 我明白了。您不是将它们存储在内存中,而是使用 yield 从生成器中生成这些值,对吗?我选择了另一个答案,因为 yield 关键字有点复杂,我花了一点时间才明白你在做什么。无论如何,我赞成你的回答,我非常感谢你的帮助。继续做你所做的事情:-)!
【解决方案2】:

根据the comments, 我们希望每个进程都在一个 10000 行的块上工作。这并不难 去做;请参阅下面的iter/islice 配方。但是,使用的问题

pool.map(worker, ten_thousand_row_chunks)

pool.map 是否会尝试将所有块放入任务队列 一次。如果这需要比可用内存更多的内存,那么你会得到一个 MemoryError。 (注:pool.imapsuffers from the same problem。)

因此,我们需要在每个块的片段上迭代地调用pool.map

import itertools as IT
import multiprocessing as mp
import csv

def worker(chunk):
    return len(chunk)

def main():
    # num_procs is the number of workers in the pool
    num_procs = mp.cpu_count()
    # chunksize is the number of lines in a chunk
    chunksize = 10**5

    pool = mp.Pool(num_procs)
    largefile = 'Counseling.csv'
    results = []
    with open(largefile, 'rb') as f:
        reader = csv.reader(f)
        for chunk in iter(lambda: list(IT.islice(reader, chunksize*num_procs)), []):
            chunk = iter(chunk)
            pieces = list(iter(lambda: list(IT.islice(chunk, chunksize)), []))
            result = pool.map(worker, pieces)
            results.extend(result)
    print(results)
    pool.close()
    pool.join()

main()

每个chunk 最多由文件中的chunksize*num_procs 行组成。 这些数据足以为池中的所有工作人员提供一些工作,但不会太大而导致 MemoryError——前提是 chunksize 没有设置太大。

然后将每个chunk 分成几块,每块最多包含 chunksize 文件中的行。然后将这些片段发送到pool.map


iter(lambda: list(IT.islice(iterator, chunksize)), []) 是如何工作的

这是一个习惯用法,用于将迭代器分组为长度为 chunksize 的块。 让我们通过一个例子来看看它是如何工作的:

In [111]: iterator = iter(range(10))

请注意,每次调用 IT.islice(iterator, 3) 时,都会生成一个包含 3 个项目的新块 从迭代器中被切掉:

In [112]: list(IT.islice(iterator, 3))
Out[112]: [0, 1, 2]

In [113]: list(IT.islice(iterator, 3))
Out[113]: [3, 4, 5]

In [114]: list(IT.islice(iterator, 3))
Out[114]: [6, 7, 8]

当迭代器中剩余的项少于 3 个时,只返回剩余的项:

In [115]: list(IT.islice(iterator, 3))
Out[115]: [9]

如果你再次调用它,你会得到一个空列表:

In [116]: list(IT.islice(iterable, 3))
Out[116]: []

lambda: list(IT.islice(iterator, chunksize)) 是一个在调用时返回list(IT.islice(iterator, chunksize)) 的函数。它是一个“单线”,相当于

def func():
    return  list(IT.islice(iterator, chunksize))

最后,iter(callable, sentinel) 返回另一个迭代器。此迭代器产生的值是可调用对象返回的值。它继续产生值,直到可调用返回等于哨兵的值。所以

iter(lambda: list(IT.islice(iterator, chunksize)), [])

将继续返回值list(IT.islice(iterator, chunksize)),直到该值是空列表:

In [121]: iterator = iter(range(10))

In [122]: list(iter(lambda: list(IT.islice(iterator, 3)), []))
Out[122]: [[0, 1, 2], [3, 4, 5], [6, 7, 8], [9]]

【讨论】:

  • 哇!伟大而描述性的答案。太感谢了。我现在明白了很多。如果我可以问你一个问题,你是如何在这些事情上如此擅长并直观地理解这些 Python 原则的?你有可以推荐的书或资源吗?
  • 有很多其他人比我知道的多得多,所以我更认同你,提出问题的人,而不是试图回答问题的人。而且,可能没有a royal road。有一件事,也许真的对我有帮助——我收集了一些简短的例子,展示了 Python 中每个特性和函数的使用。
  • 我认为您阅读的文档并不重要。网上有很多很棒的免费文档和教程。重要的是你练习和使用语言。具体的例子使语言的含义和行为变得清晰。所以我能给出的最好建议是享受编程并参与a lot of practice/play
  • 如果我有一个函数说func1,它只接受Counseling.csv 文件中特定列的一行var1 作为输入,这个函数将生成一个列表,它将写入名为“output.csv”的新csv 文件?
猜你喜欢
  • 1970-01-01
  • 2015-10-10
  • 1970-01-01
  • 1970-01-01
  • 2013-04-22
  • 2016-02-24
  • 1970-01-01
  • 2019-01-17
  • 1970-01-01
相关资源
最近更新 更多