【问题标题】:How do I abort a socket.recvfrom() from another thread in python?如何从 python 中的另一个线程中止 socket.recvfrom()?
【发布时间】:2011-09-16 19:03:28
【问题描述】:

这看起来像 How do I abort a socket.recv() from another thread in Python 的副本,但它不是,因为我想在一个线程中中止 recvfrom(),它是 UDP,而不是 TCP。

这可以通过 poll() 或 select.select() 解决吗?

【问题讨论】:

  • UDP 或 TCP 与是否线程无关。
  • 您可以使用旧的 C 技巧:创建一个假管道并在假和套接字上使用 select。当你想停止简单的发送消息给假的......如果我有时间我会提交一个包含所有细节的答案。
  • @qarma 看看我的评论,如果您对这种解决方案感兴趣,请告诉我。
  • @Micheled'Amico 是的,这是一个有效的解决方案
  • @qarma 我会写这个我们

标签: python multithreading udp pthreads


【解决方案1】:

如果您想解除阻塞从另一个线程读取的 UDP,请向其发送数据报!

Rgds, 马丁

【讨论】:

  • 如果(addr, port) 组合是独一无二的,那就太好了。如果REUSEADDR 处于开启状态并且有多个进程正在侦听,该怎么办?
【解决方案2】:

处理这种异步中断的好方法是旧的 C 管道技巧。您可以创建一个管道并在套接字和管道上使用select/poll:现在,当您需要中断接收器时,您只需向管道发送一个字符。

  • 优点:
    • 可以同时用于 UDP 和 TCP
    • 与协议无关
  • 缺点:
    • 管道上的选择/轮询在 Windows 上不可用,在这种情况下,您应该将其替换为另一个用作通知的 UDP 套接字管道

起点

interruptable_socket.py

import os
import socket
import select


class InterruptableUdpSocketReceiver(object):
    def __init__(self, host, port):
        self._host = host
        self._port = port
        self._socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
        self._socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
        self._r_pipe, self._w_pipe = os.pipe()
        self._interrupted = False

    def bind(self):
        self._socket.bind((self._host, self._port))

    def recv(self, buffersize, flags=0):
        if self._interrupted:
            raise RuntimeError("Cannot be reused")
        read, _w, errors = select.select([self._r_pipe, self._socket], [], [self._socket])
        if self._socket in read:
            return self._socket.recv(buffersize, flags)
        return ""

    def interrupt(self):
        self._interrupted = True
        os.write(self._w_pipe, "I".encode())

一个测试套件:

test_interruptable_socket.py

import socket
from threading import Timer
import time
from interruptable_socket import InterruptableUdpSocketReceiver
import unittest


class Sender(object):
    def __init__(self, destination_host, destination_port):
        self._socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM, socket.IPPROTO_UDP)
        self._dest = (destination_host, destination_port)

    def send(self, message):
        self._socket.sendto(message, self._dest)

class Test(unittest.TestCase):
    def create_receiver(self, host="127.0.0.1", port=3010):
        receiver = InterruptableUdpSocketReceiver(host, port)
        receiver.bind()
        return receiver

    def create_sender(self, host="127.0.0.1", port=3010):
        return Sender(host, port)

    def create_sender_receiver(self, host="127.0.0.1", port=3010):
        return self.create_sender(host, port), self.create_receiver(host, port)

    def test_create(self):
        self.create_receiver()

    def test_recv_async(self):
        sender, receiver = self.create_sender_receiver()
        start = time.time()
        send_message = "TEST".encode('UTF-8')
        Timer(0.1, sender.send, (send_message, )).start()
        message = receiver.recv(128)
        elapsed = time.time()-start
        self.assertGreaterEqual(elapsed, 0.095)
        self.assertLess(elapsed, 0.11)
        self.assertEqual(message, send_message)

    def test_interrupt_async(self):
        receiver = self.create_receiver()
        start = time.time()
        Timer(0.1, receiver.interrupt).start()
        message = receiver.recv(128)
        elapsed = time.time()-start
        self.assertGreaterEqual(elapsed, 0.095)
        self.assertLess(elapsed, 0.11)
        self.assertEqual(0, len(message))

    def test_exception_after_interrupt(self):
        sender, receiver = self.create_sender_receiver()
        receiver.interrupt()
        with self.assertRaises(RuntimeError):
            receiver.recv(128)


if __name__ == '__main__':
    unittest.main()

进化

现在这段代码只是一个起点。为了使其更通用,我认为我们应该解决以下问题:

  1. 接口:中断情况下返回空消息不好处理,最好使用异常处理
  2. 泛化:我们应该在socket.recv()之前只调用一个函数,将中断扩展到其他recv方法变得非常简单
  3. 可移植性:为了简单地将其移植到 Windows,我们应该将异步通知隔离在一个对象中,以便为我们的操作系统选择正确的实现

首先我们更改test_interrupt_async() 以检查异常而不是空消息:

from interruptable_socket import InterruptException

def test_interrupt_async(self):
    receiver = self.create_receiver()
    start = time.time()
    with self.assertRaises(InterruptException):
        Timer(0.1, receiver.interrupt).start()
        receiver.recv(128)
    elapsed = time.time()-start
    self.assertGreaterEqual(elapsed, 0.095)
    self.assertLess(elapsed, 0.11)

在此之后,我们可以将return '' 替换为raise InterruptException 并且测试再次通过。

准备好扩展的版本可以是:

interruptable_socket.py

import os
import socket
import select


class InterruptException(Exception):
    pass


class InterruptableUdpSocketReceiver(object):
    def __init__(self, host, port):
        self._host = host
        self._port = port
        self._socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
        self._socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
        self._async_interrupt = AsycInterrupt(self._socket)

    def bind(self):
        self._socket.bind((self._host, self._port))

    def recv(self, buffersize, flags=0):
        self._async_interrupt.wait_for_receive()
        return self._socket.recv(buffersize, flags)

    def interrupt(self):
        self._async_interrupt.interrupt()


class AsycInterrupt(object):
    def __init__(self, descriptor):
        self._read, self._write = os.pipe()
        self._interrupted = False
        self._descriptor = descriptor

    def interrupt(self):
        self._interrupted = True
        self._notify()

    def wait_for_receive(self):
        if self._interrupted:
            raise RuntimeError("Cannot be reused")
        read, _w, errors = select.select([self._read, self._descriptor], [], [self._descriptor])
        if self._descriptor not in read:
            raise InterruptException

    def _notify(self):
        os.write(self._write, "I".encode())

现在包装更多recv函数,实现windows版本或处理socket超时变得非常简单。

【讨论】:

  • 代码量太大了,为什么不直接继承socket呢?
  • 因为在这个例子中我不想处理所有 4 个recvs 方法并覆盖(和测试)所有这些方法。此外,如果我将socket 子类化,我应该在所有recvs 方法中处理套接字超时。这是一个示例和一组测试:用子类替换合作者很简单(但您应该解决上述问题)。但在此之前最好 1- 使用异常而不是返回空消息,2- 创建一个装饰器以检查并等待传入​​消息在所有 4 recvs 中使用它,3- 提取一个执行异步通知的对象以替换它是通过 Windows 中的 UDP 实现来实现的。
  • @qarma 还有一件事:子类化是一种强大的工具,但有时使用协作者会更好,因为当您对新对象进行子类化时,它就是父对象,您应该注意所有父对象的行为和责任。使用协作者使更简单遵循单一职责原则并使您的代码更干净:如果您在对象中仅使用recv_from,为什么要创建一个处理所有socket 方法的子类?当您需要一些其他方法时,您可以考虑在子类中对其进行更改,您的测试将帮助您安全地进行。
  • 哇,你多产了!虽然我更喜欢与套接字相同的 API,并且代码库更小,但您的答案是最好的,ergo +50 是您的。
【解决方案3】:

在服务器和客户端套接字上执行退出命令。应该像这样工作:

Thread1: 
    status: listening
    handler: quit

Thread2: client
    exec: socket.send "quit"  ---> Thread1.socket @ host:port

Thread1: 
    status: socket closed()

【讨论】:

    【解决方案4】:

    这里的解决方法是强制关闭socket。问题在于执行此操作的方法是特定于操作系统的,而 Python 在抽象执行此操作的方式或结果方面做得不好。基本上,您需要在套接字上执行 shutdown(),然后执行 close()。在 POSIX 系统(如 Linux)上,关闭是强制 recvfrom 停止的关键元素(单独调用 close() 不会这样做)。在 Windows 上,shutdown() 不会影响 recvfrom,而 close() 是关键元素。如果您在 C 中实现此代码并使用本机 POSIX 套接字或 Winsock 套接字,这正是您会看到的行为,因此 Python 在这些调用之上提供了一个非常薄的层。

    在 POSIX 和 Windows 系统上,这一系列调用会导致引发 OSError。但是,异常的位置及其详细信息是特定于操作系统的。在 POSIX 系统上,调用 shutdown() 时会引发异常,并且异常的 errno 值设置为 107(未连接传输端点)。在 Windows 系统上,调用 recvfrom() 时会引发异常,并且异常的 winerror 值设置为 10038(尝试对非套接字的操作进行操作)。这意味着无法以与操作系统无关的方式执行此操作,代码必须同时考虑 Windows 和 POSIX 行为和错误。这是我写的一个简单的例子:

    import socket
    import threading
    import time
    
    class MyServer(object):
        def __init__(self, port:int=0):
            if port == 0:
                raise AttributeError('Invalid port supplied.')
    
            self.port = port
            self.socket = socket.socket(family=socket.AF_INET,
                    type=socket.SOCK_DGRAM)
            self.socket.bind(('0.0.0.0', port))
    
            self.exit_now = False
    
            print('Starting server.')
            self.thread = threading.Thread(target=self.run_server,
                    args=[self.socket])
            self.thread.start()
    
        def run_server(self, socket:socket.socket=None):
            if socket is None:
                raise AttributeError('No socket provided.')
    
            buffer_size = 4096
    
            while self.exit_now == False:
                data = b''
                try:
                    data, address = socket.recvfrom(buffer_size)
                except OSError as e:
                    if e.winerror == 10038:
                        # Error is, "An operation was attempted on something that
                        # is not a socket".  We don't care.
                        pass
                    else:
                        raise e
                if len(data) > 0:
                    print(f'Received {len(data)} bytes from {address}.')
    
        def stop(self):
            self.exit_now = True
            try:
                self.socket.shutdown(socket.SHUT_RDWR)
            except OSError as e:
                if e.errno == 107:
                    # Error is, "Transport endpoint is not connected".
                    # We don't care.
                    pass
                else:
                    raise e
            self.socket.close()
            self.thread.join()
            print('Server stopped.')
    
    
    if __name__ == '__main__':
        server = MyServer(5555)
        time.sleep(2)
        server.stop()
        exit(0)
    

    【讨论】:

    • 单独调用close()执行此操作,除非套接字 FD 已被继承。
    【解决方案5】:

    要在 python 中正确关闭一个 tcp 套接字,你必须在调用 socket.close() 之前调用 socket.shutdown(arg)。请参阅 python 套接字文档,关于关闭的部分。

    如果套接字是UDP,你不能调用socket.shutdown(...),它会引发异常。并且单独调用 socket.close() 会像 tcp 一样保持被阻塞的操作阻塞。单独的 close() 不会打断他们。

    许多建议的解决方案(不是全部)不起作用或被视为繁琐,因为它们涉及 3rd 方库。我还没有测试过 poll() 或 select()。确实有效的方法如下:

    首先,为任何正在运行 socket.recv() 的线程创建一个官方的 Thread 对象,并保存它的句柄。其次,导入信号。 Signal 是一个官方库,它可以向进程发送/接收 linux/posix 信号(阅读其文档)。第三,中断,假设你的线程句柄被称为 udpThreadHandle:

    signal.pthread_kill(udpthreadHandle.ident, signal.SIGINT)
    

    当然,在实际的线程/循环中进行接收:

    try:
        while True:
           myUdpSocket.recv(...)
    except KeyboardInterrupt:
        pass
    

    注意,KeyboardInterrupt 的异常处理程序(由 SIGINT 生成)位于接收循环之外。这会静默终止接收循环及其线程。

    【讨论】:

      猜你喜欢
      • 2011-02-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-12-29
      • 1970-01-01
      • 2014-11-29
      • 1970-01-01
      • 2010-12-13
      相关资源
      最近更新 更多