【问题标题】:python -> multiprocessing modulepython -> 多处理模块
【发布时间】:2011-04-04 22:38:15
【问题描述】:

这就是我想要完成的 -

  1. 我有大约一百万个文件需要解析并将解析后的内容附加到单个文件中。
  2. 由于单个进程需要很长时间,因此此选项已失效。
  3. 不使用 Python 中的线程,因为它本质上是运行单个进程(由于 GIL)。
  4. 因此使用多处理模块。即产生 4 个子进程来利用所有原始核心功能:)

到目前为止一切顺利,现在我需要一个所有子进程都可以访问的共享对象。我正在使用多处理模块中的队列。此外,所有子流程都需要将其输出写入单个文件。我猜是使用锁的潜在场所。当我运行这个设置时,我没有收到任何错误(所以父进程看起来很好),它只是停止。当我按 ctrl-C 时,我看到一个回溯(每个子进程一个)。也没有输出写入输出文件。这是代码(请注意,在没有多进程的情况下一切运行良好) -

import os
import glob
from multiprocessing import Process, Queue, Pool

data_file  = open('out.txt', 'w+')

def worker(task_queue):
    for file in iter(task_queue.get, 'STOP'):
        data = mine_imdb_page(os.path.join(DATA_DIR, file))
        if data:
            data_file.write(repr(data)+'\n')
    return

def main():
    task_queue = Queue()
    for file in glob.glob('*.csv'):
        task_queue.put(file)
    task_queue.put('STOP') # so that worker processes know when to stop

    # this is the block of code that needs correction.
    if multi_process:
        # One way to spawn 4 processes
        # pool = Pool(processes=4) #Start worker processes
        # res  = pool.apply_async(worker, [task_queue, data_file])

        # But I chose to do it like this for now.
        for i in range(4):
            proc = Process(target=worker, args=[task_queue])
            proc.start()
    else: # single process mode is working fine!
        worker(task_queue)
    data_file.close()
    return

我做错了什么?我还尝试在生成时将打开的 file_object 传递给每个进程。但是没有效果。例如-Process(target=worker, args=[task_queue, data_file])。但这并没有改变什么。我觉得子进程由于某种原因无法写入文件。 file_object 的实例没有被复制(在生成时)或其他一些怪癖......有人知道吗?

EXTRA: 还有有什么方法可以保持持久的 mysql_connection 打开并将其传递给 sub_processes?所以我在我的父进程中打开了一个 mysql 连接,并且我的所有子进程都应该可以访问打开的连接。基本上这相当于 python 中的 shared_memory 。这里有什么想法吗?

【问题讨论】:

  • 如果您不写入文件而是进行打印,那么它可以工作吗? (在 Linux 上,我会使用 python script.py > out.dat 来防止屏幕泛滥)。
  • 我认为 proc.start 是非阻塞的,所以你可能应该在某个地方等待,让进程有机会在你做 datafile.close() 之前做一些工作
  • data_file.close() 在最后完成。它应该在这里起作用吗?打印工作也很好。当我使用打印时,我在屏幕上看到输出......但我想使用文件。帮助!还有有什么方法可以保持持久的 mysql_connection 打开并将其传递给 sub_processes?
  • @extraneon: 很好,但是如果程序试图在一个关闭的文件上写入,应该引发一个异常。
  • 你是在读mysql还是写mysql,还是两者兼而有之?

标签: python queue multiprocessing


【解决方案1】:

虽然与 Eric 的讨论富有成果,但后来我找到了更好的方法。在多处理模块中,有一个名为“Pool”的方法非常适合我的需求。

它会根据我的系统拥有的核心数量进行自我优化。即只产生与否一样多的进程。的核心。当然,这是可定制的。所以这里是代码。以后可能会帮助别人-

from multiprocessing import Pool

def main():
    po = Pool()
    for file in glob.glob('*.csv'):
        filepath = os.path.join(DATA_DIR, file)
        po.apply_async(mine_page, (filepath,), callback=save_data)
    po.close()
    po.join()
    file_ptr.close()

def mine_page(filepath):
    #do whatever it is that you want to do in a separate process.
    return data

def save_data(data):
    #data is a object. Store it in a file, mysql or...
    return

仍在经历这个巨大的模块。不确定 save_data() 是由父进程执行还是由衍生的子进程使用。如果是孩子进行了保存,则在某些情况下可能会导致并发问题。如果有人有更多使用此模块的经验,您可以在这里了解更多知识...

【讨论】:

    【解决方案2】:

    多处理的文档指出了几种在进程之间共享状态的方法:

    http://docs.python.org/dev/library/multiprocessing.html#sharing-state-between-processes

    我确信每个进程都会获得一个新的解释器,然后将目标(函数)和参数加载到其中。在这种情况下,脚本中的全局命名空间将绑定到您的工作函数,因此 data_file 将在那里。但是,我不确定文件描述符在复制时会发生什么。您是否尝试过将文件对象作为参数之一传递?

    另一种方法是传递另一个队列来保存工作人员的结果。工人put结果和主代码gets结果并将其写入文件。

    【讨论】:

    • 是的!我可以那样做。我可以有另一个队列,它类似于进程写入的 out_queue。由于父进程可以访问此它可以继续读取此队列并写入文件。这可以工作!我也尝试将文件对象作为参数之一传递。它似乎不起作用。线程不写入文件。还有 Eric,知道如何将持久的 mysql 连接传递给子进程吗?
    • @Srikar,希望对您有所帮助。至于mysql连接,我不确定那个。我会说你最好为每个进程单独连接。即使您可以建立连接,我也不确定它有多“线程安全”。如果你真的只需要分享一个,那么你可能不得不做一些奇怪的事情。再说一次,您也可以在队列中代理连接的查询/响应机制。然后主进程(或单独的 mysql 处理程序进程)从队列中获取查询,运行它们,然后将结果放回......或类似的东西。
    猜你喜欢
    • 2017-12-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-11-21
    • 2020-06-25
    相关资源
    最近更新 更多