【问题标题】:Threading Queue and multiprocessing assistance线程队列和多处理辅助
【发布时间】:2016-05-02 10:32:14
【问题描述】:

前言:这是我第一次尝试使用这些工具

上下文:我正在尝试处理一个非常大的文件。所以我试图把文件分成更小的块。然后将这些文件加载​​到队列中进行处理。

目标是加快一个非常缓慢的过程。

代码:

import lifetimes
import os
import pandas
import Queue
import threading
import multiprocessing
import glob
import subprocess


#move master to processing dir
os.system("cp /data/ltv-testing1.csv /data/out")

#break master csv into 1 million row chunks
subprocess.call(['bash', '/home/ddewberry/LTV_CSV_Split.sh'])

#remove master file
os.remove("/data/out/ltv-testing1.csv")

os.chdir("/data/out")


# Create List of Files
worker_data = glob.glob('split_*')

#build queue with file list
q = Queue.Queue(worker_data)

#import tools for data processing
from lifetimes.utils import summary_data_from_transaction_data

#define worker for threads

def worker(outfile = '/data/in/Worker.csv'):
    while True:
        item = q.get()
        data = pandas.read_csv(item)
        summary = summary_data_from_transaction_data(data, data[[2]], data[[1]])
        summary.to_csv(outfile%s % (item))
        q.task_done()

cpus=multiprocessing.cpu_count() #detect number of cores
print("Creating %d threads" % cpus)
for i in range(cpus):
     t = threading.Thread(target=worker)
     t.daemon = True
     t.start()

q.join()

#clean up
for row in worker_data:
    os.remove(row)

问题:

我没有收到任何错误消息,但它根本不起作用。 (它基本上什么都不做)

我对自己做错了什么或需要解决什么感到非常困惑。

【问题讨论】:

  • 对于初学者来说,Queue.Queue 接受一个参数maxsize,而不是一个可迭代的参数,所以q.get() 将无限期阻塞,因为其中没有任何内容......另外,对于这种问题线程不会给你很大的加速。
  • 感谢您的解释。你对我可以采取什么方法来帮助加快这个过程有什么建议吗?
  • 所以我应该这样做:q = Queue.Queue() for file in worker_data: q.put(file)
  • 应该可以,但我建议您查看multiprocessing.Pool,当使用 Pool.map 或类似时,您不需要处理队列等...
  • Mata - 非常感谢您的建议。我将探索该解决方案,感谢指导

标签: multithreading python-2.7 queue python-multiprocessing


【解决方案1】:
import lifetimes
import os
import pandas
import Queue
import threading
import multiprocessing
import glob
import subprocess


#move master to processing dir
os.system("cp /data/ltv-testing1.csv /data/out")

#break master csv into 1 million row chunks
subprocess.call(['bash', '/home/ddewberry/LTV_CSV_Split.sh'])

#remove master file
os.remove("/data/out/ltv-testing1.csv")

os.chdir("/data/out")


# Create List of Files
worker_data = glob.glob('split_*')

# rename to csv
for row in worker_data:
    os.rename(row, row+'.csv')

worker_data1 = glob.glob('split_*')

#build queue with file list
q = Queue.Queue()
for files in worker_data1:
    q.put(files)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2012-11-12
    • 2014-10-28
    • 2015-04-15
    • 1970-01-01
    • 1970-01-01
    • 2012-12-06
    • 2018-06-27
    • 2013-11-23
    相关资源
    最近更新 更多