【问题标题】:Python multiprocessing: Parallel Pipeline implementationPython 多处理:并行管道实现
【发布时间】:2017-07-16 18:14:50
【问题描述】:

我正在尝试创建一个管道,但我遇到了糟糕的退出问题(僵尸)和性能问题。我创建了这个通用类:

class Generator(Process):
'''
<function>: function to call. None value means that the current class will
    be used as a template for another class, with <function> being defined
    there
<input_queues> : Queue or list of Queue objects , which refer to the input
    to <function>.
<output_queues> : Queue or list of Queue objects , which are used to pass
    output
<sema_to_acquire> : Condition or list of Condition objects, which are
    blocking generation while not notified
<sema_to_release> : Condition or list of Condition objects, which will be
    notified after <function> is called
'''

def __init__(self, function=None, input_queues=None, output_queues=None, sema_to_acquire=None,
             sema_to_release=None):
    Process.__init__(self)
    self.input_queues = input_queues
    self.output_queues = output_queues
    self.sema_to_acquire = sema_to_acquire
    self.sema_to_release = sema_to_release
    if function is not None:
            self.function = function

def run(self):
    if self.sema_to_release is not None:
            try:
                self.sema_to_release.release()
            except AttributeError:
                [sema.release() for sema in self.sema_to_release]

    while True:
        if self.sema_to_acquire is not None:
            try:
                self.sema_to_acquire.acquire()
            except AttributeError:
                [sema.acquire() for sema in self.sema_to_acquire]


        if self.input_queues is not None:
            try:
                data = self.input_queues.get()
            except AttributeError:
                data = [queue.get() for queue in self.input_queues]
            isiterable = True
            try:
                iter(data)
                res = self.function(*tuple(data))
            except TypeError, te:
                res = self.function(data)
        else:
            res = self.function()
        if self.output_queues is not None:
            try:
                if self.output_queues.full():
                    self.output_queues.get(res)
                self.output_queues.put(res)
            except AttributeError:
                [queue.put(res) for queue in self.output_queues]
        if self.sema_to_release is not None:
            if self.sema_to_release is not None:
                try:
                    self.sema_to_release.release()
                except AttributeError:
                    [sema.release() for sema in self.sema_to_release]

模拟管道内的工作人员。生成器希望运行一个无限的 while 循环,其中使用来自 n 个队列的输入执行一个函数,并将结果写入 m 个队列。在一个迭代发生之前,有一些信号量需要被进程获取,当迭代完成时,其他一些信号量被释放。因此,对于需要并行运行并为另一个进程生成输入的进程,我发送“交叉”信号量作为参数,以强制它们一起执行单次迭代。对于不需要并行运行的进程,我不使用任何条件。一个例子(如果有人忽略输入函数,我实际使用)如下:

import time
from multiprocess import Lock
print_lock = Lock()
_t_=0.5
def func0(data):
    time.sleep(_t_)
    print_lock.acquire()
    print 'func0 sends',data
    print_lock.release()
    return data
def func1(data):
    time.sleep(_t_)
    print_lock.acquire()
    print 'func1 receives and sends',data
    print_lock.release()
    return data
def func2(data):
    time.sleep(_t_)
    print_lock.acquire()
    print 'func2 receives and sends',data
    print_lock.release()
    return data
def func3(*data):
    print_lock.acquire()
    print 'func3 receives',data
    print_lock.release()


run_svm = Semaphore()
run_rf = Semaphore()
inp_rf = Queue()
inp_svm = Queue()
out_rf = Queue()
out_svm = Queue()
kin_stream = Queue()
res_mixed = Queue()
streamproc = Generator(func0,
                       input_queues=kin_stream,
                       output_queues=[inp_rf,
                                       inp_svm])
streamproc.daemon = True
streamproc.start()
svm_class = Generator(func1,
                       input_queues=inp_svm,
                       output_queues=out_svm,
                       sema_to_acquire=run_svm,
                       sema_to_release=run_rf)
svm_class.daemon=True
svm_class.start()
rf_class = Generator(func2,
                      input_queues=inp_rf,
                      output_queues=out_rf,
                      sema_to_acquire=run_rf,
                      sema_to_release=run_svm)
rf_class.daemon=True
rf_class.start()
mixed_class = Generator(func3,
                         input_queues=[out_rf, out_svm])
mixed_class.daemon = True
mixed_class.start()
count = 1
while True:
    kin_stream.put([count])
    count+=1
    time.sleep(1)
streamproc.join()
svm_class.join()
rf_class.join()
mixed_class.join()

这个例子给出:

func0 sends 1
func2 receives and sends 1
func1 receives and sends 1
func3 receives (1, 1)
func0 sends 2
func2 receives and sends 2
func1 receives and sends 2
func3 receives (2, 2)
func0 sends 3
func2 receives and sends 3
func1 receives and sends 3
func3 receives (3, 3)
...

一切都好。 但是,如果我尝试杀死 main,则不能保证其他子进程终止:终端可能会冻结,或者 python 编译器可能会继续在后台运行(可能是僵尸),我不知道为什么这正在发生,因为我已将相应的守护进程设置为 True。 有没有人对实现这种类型的管道有更好的想法,或者可以提出解决这个邪恶问题的方法?谢谢大家。

编辑

固定测试。然而,僵尸仍然确实存在。

【问题讨论】:

    标签: python performance subprocess pipeline python-multiprocessing


    【解决方案1】:

    我能够克服这个问题,方法是引入一个终止队列作为给定类的附加参数,并为 SIGINT 中断设置一个信号处理程序,以停止管道执行。我不知道这是否是让它工作的最优雅的方式,但它确实有效。此外,信号处理程序的设置方式很重要,因为某些原因必须在 process.start() 之前设置,如果有人知道原因,他可以发表评论。此外,信号处理程序由子进程继承,因此我必须将join 放入try:..except AssertionError:pass 模式中,否则会抛出错误(再次,如果有人知道如何绕过它,请详细说明)。无论如何,它有效。

    源代码

    class Generator(Process):
        '''
        <term_queue>: Queue to write termination events, must be same for all
                    processes spawned
        <function>: function to call. None value means that the current class will
            be used as a template for another class, with <function> being defined
            there
        <input_queues> : Queue or list of Queue objects , which refer to the input
            to <function>.
        <output_queues> : Queue or list of Queue objects , which are used to pass
            output
        <sema_to_acquire> : Semaphore or list of Semaphore objects, which are
            blocking function execution
        <sema_to_release> : Semaphore or list of Semaphore objects, which will be
            released after <function> is called
        '''
    
        def __init__(self, term_queue,
                     function=None, input_queues=None, output_queues=None, sema_to_acquire=None,
                     sema_to_release=None):
            Process.__init__(self)
            self.term_queue = term_queue
            self.input_queues = input_queues
            self.output_queues = output_queues
            self.sema_to_acquire = sema_to_acquire
            self.sema_to_release = sema_to_release
            if function is not None:
                self.function = function
    
        def run(self):
            if self.sema_to_release is not None:
                try:
                    self.sema_to_release.release()
                except AttributeError:
                    deb = [sema.release() for sema in self.sema_to_release]
            while True:
                if not self.term_queue.empty():
                    self.term_queue.put((self.name, 0))
                    break
                try:
                    if self.sema_to_acquire is not None:
                        try:
                            self.sema_to_acquire.acquire()
                        except AttributeError:
                            deb = [sema.acquire() for sema in self.sema_to_acquire]
    
                    if self.input_queues is not None:
                        try:
                            data = self.input_queues.get()
                        except AttributeError:
                            data = tuple([queue.get()
                                          for queue in self.input_queues])
                        res = self.function(data)
                    else:
                        res = self.function()
                    if self.output_queues is not None:
                        try:
                            if self.output_queues.full():
                                self.output_queues.get(res)
                            self.output_queues.put(res)
                        except AttributeError:
                            deb = [queue.put(res) for queue in self.output_queues]
                    if self.sema_to_release is not None:
                        if self.sema_to_release is not None:
                            try:
                                self.sema_to_release.release()
                            except AttributeError:
                                deb = [sema.release() for sema in self.sema_to_release]
                except Exception as exc:
                    self.term_queue.put((self.name, exc))
                    break
    
    
    
    def signal_handler(sig, frame, term_queue, processes):
        '''
        <term_queue> is the queue to write termination of the __main__
        <processes> is a dicitonary holding all running processes
        '''
        term_queue.put((__name__, 'SIGINT'))
        try:
            [processes[key].join() for key in processes]
        except AssertionError:
            pass
        sys.exit(0)
    
    term_queue = Queue()
    '''
    initialize some Generators and add them to <processes> dicitonary
    '''
    signal.signal(signal.SIGINT, lambda sig,frame: signal_handler(sig,frame,
                                                                  term_queue,processes))
    [processes[key].start() for key in processes]
    while True:
        if not term_queue.empty():
            [processes[key].join() for key in processes]
            break
    

    示例也相应更改(如果您希望我添加,请评论)

    【讨论】:

      【解决方案2】:

      我也不得不解决这个问题,事实上,将一些通信管道或队列传递给进程似乎是告诉它们终止的最简单方法。

      但是,终止代码可以利用主进程中的 finally: 块,它会处理包括信号在内的任何事件。

      如果您的进程应该与对象同时终止,您可能还想使用weakref.finalize,但它可以是tricky。

      【讨论】:

        猜你喜欢
        • 2011-12-17
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-02-27
        • 2018-06-30
        • 2013-02-25
        相关资源
        最近更新 更多