【发布时间】: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