代码之家  ›  专栏  ›  技术社区  ›  Géry Ogam

在接收到SIGINT并在Python中使用代理对象后,如何防止BrokenPipeErrors?

  •  0
  • Géry Ogam  · 技术社区  · 7 年前

    此Python程序:

    import concurrent.futures
    import multiprocessing
    import time
    
    class A:
    
        def __init__(self):
            self.event = multiprocessing.Manager().Event()
    
        def start(self):
            try:
                while True:
                    if self.event.is_set():
                        break
                    print("processing")
                    time.sleep(1)
            except BaseException as e:
                print(type(e).__name__ + " (from pool thread):", e)
    
        def shutdown(self):
            self.event.set()
    
    if __name__ == "__main__":
        try:
            a = A()
            pool = concurrent.futures.ThreadPoolExecutor(1)
            future = pool.submit(a.start)
            while not future.done():
                concurrent.futures.wait([future], timeout=0.1)
        except BaseException as e:
            print(type(e).__name__ + " (from main thread):", e)
        finally:
            a.shutdown()
            pool.shutdown()
    

    输出:

    processing
    processing
    processing
    KeyboardInterrupt (from main thread):
    BrokenPipeError (from pool thread): [WinError 232] The pipe is being closed
    Traceback (most recent call last):
      File "C:\Program Files\Python37\lib\multiprocessing\managers.py", line 788, in _callmethod
        conn = self._tls.connection
    AttributeError: 'ForkAwareLocal' object has no attribute 'connection'
    
    During handling of the above exception, another exception occurred:
    
    Traceback (most recent call last):
      File ".\foo.py", line 34, in <module>
        a.shutdown()
      File ".\foo.py", line 21, in shutdown
        self.event.set()
      File "C:\Program Files\Python37\lib\multiprocessing\managers.py", line 1067, in set
        return self._callmethod('set')
      File "C:\Program Files\Python37\lib\multiprocessing\managers.py", line 792, in _callmethod
        self._connect()
      File "C:\Program Files\Python37\lib\multiprocessing\managers.py", line 779, in _connect
        conn = self._Client(self._token.address, authkey=self._authkey)
      File "C:\Program Files\Python37\lib\multiprocessing\connection.py", line 490, in Client
        c = PipeClient(address)
      File "C:\Program Files\Python37\lib\multiprocessing\connection.py", line 691, in PipeClient
        _winapi.WaitNamedPipe(address, 1000)
    FileNotFoundError: [WinError 2] The system cannot find the file specified
    

    当它运行时 SIGINT 三秒钟后发送信号(按 Ctrl键 + C )。

    分析 –The SIGINT公司 信号被发送到每个进程的主线程。在这种情况下,有两个流程:主流程和经理的子流程。

    • 在主进程的主线程中:收到 SIGINT公司 信号,默认值 SIGINT公司 信号处理程序引发 KeyboardInterrupt 捕获并打印异常。
    • 在管理器子进程的主线程中:同时,在收到 SIGINT公司 信号,默认值 SIGINT公司 信号处理程序引发 键盘中断 异常,终止子进程。因此,其他进程随后对管理器共享对象的所有使用都会引发 BrokenPipeError 例外
    • 在主进程池的子线程中:在本例中 断开管道错误 在该行引发异常 if self.event.is_set():
    • 在主进程的主线程中:最后,控制流到达行 a.shutdown() ,这提高了 AttributeError FileNotFoundError 例外情况。

    如何防止这种情况 断开管道错误 例外

    0 回复  |  直到 7 年前
        1
  •  1
  •   Géry Ogam    7 年前

    此问题的解决方案是覆盖默认值 SIGINT 具有将忽略信号的处理程序的信号处理程序,例如 signal.SIG_IGN 标准信号处理器。可以通过调用 signal.signal 经理子流程开始时的功能:

    import concurrent.futures
    import multiprocessing.managers
    import signal
    import time
    
    def init():
        signal.signal(signal.SIGINT, signal.SIG_IGN)
    
    class A:
    
        def __init__(self):
            manager = multiprocessing.managers.SyncManager()
            manager.start(init)
            self.event = manager.Event()
    
        def start(self):
            try:
                while True:
                    if self.event.is_set():
                        break
                    print("processing")
                    time.sleep(1)
            except BaseException as e:
                print(type(e).__name__ + " (from pool thread):", e)
    
        def shutdown(self):
            self.event.set()
    
    if __name__ == "__main__":
        try:
            a = A()
            pool = concurrent.futures.ThreadPoolExecutor(1)
            future = pool.submit(a.start)
            while not future.done():
                concurrent.futures.wait([future], timeout=0.1)
        except BaseException as e:
            print(type(e).__name__ + " (from main thread):", e)
        finally:
            a.shutdown()
            pool.shutdown()
    

    笔记 此程序还与 concurrent.futures.ProcessPoolExecutor