【问题标题】:Zeromq (pyzmq) ROUTER procession of multiple clients' data and subsequent timeout handlingZeromq (pyzmq) ROUTER 处理多个客户端的数据和后续的超时处理
【发布时间】:2016-05-11 07:37:33
【问题描述】:

我有一个ROUTER,其目的是收集来自多个DEALER 客户端的图像数据并对完整图像执行OCR。我发现处理 OCR 最有效的方法是利用 Python 的多处理库。累积的图像字节将put 转换为Queue,以便在单独的Process 中进行适当的处​​理。但是,我需要确保当客户端遇到超时时,Process 会被正确终止,并且不会无意义地逗留和占用资源。

在我当前的解决方案中,我将每个新连接的客户端插入到 dict 中,其中 value 是我的 ClientHandler 类,它拥有所有图像数据并生成一个 Thread ,它设置了一个名为 " 的 boolean 变量超时”到True 5 秒后。如果在 5 秒帧内收到新消息,则调用 bump 并将计时器重置回 0,否则我会在线程终止之前进行清理,并从主循环中的 dict 中删除引用:

import threading
import time
import zmq

class ClientHandler(threading.Thread):
    def __init__(self, socket):
        self.elapsed = time.time()
        self.timeout = False

        self.socket = socket

        super(ClientHandler, self).__init__()

    def run(self):
        while time.time() - self.elapsed < 5.0:
            pass

        self.timeout = True

        # CLIENT TIMED OUT
        # HANDLE TERMINATION AND CLEAN UP HERE

    def bump(self):
        self.elapsed = time.time()

    def handle(self, id, header, data):
        # HANDLE CLIENT DATA HERE
        # ACCUMULATE IMAGE BYTES, ETC

        self.socket.send_multipart([id, str(0)])

def server_task():
    clients = dict()

    context = zmq.Context.instance()
    server = context.socket(zmq.ROUTER)

    server.setsockopt(zmq.RCVTIMEO, 0)

    server.bind("tcp://127.0.0.1:7777")

    while True:
        try:
            id, header, data = server.recv_multipart()

            client = clients.get(id)

            if client == None:
                client = clients[id] = ClientHandler(server)

                client.start()

            client.bump()
            client.handle(id, header, data)
        except zmq.Again:
            for id in clients.keys():
                if clients[id].timeout:
                    del clients[id]

    context.term()

if __name__ == "__main__":
    server_task()

但是整个方法感觉不太对。我这样做不正确吗?如果是这样,如果有人能指出我正确的方向,我将不胜感激。

【问题讨论】:

  • 我的兄弟建议我为每个连接的客户端分配一个端口,并在那里简单地处理每个单独的超时,而不是处理整个字典混乱。收集代码后,我将提供并标记答案。

标签: python multithreading multiprocessing zeromq


【解决方案1】:

自己想出来的,希望对别人有帮助。

我在分配的端口上有一个路由器,它将唯一的端口分配给每个客户端,然后连接到所述唯一端口上新绑定的套接字。当客户端断开连接时,端口会被回收以重新分配。

import sys
import zmq
from multiprocessing import Process, Queue, Value

def server_task():
    context = zmq.Context.instance()

    server = context.socket(zmq.ROUTER)

    server.bind("tcp://127.0.0.1:7777")

    timeout_queue = Queue()
    port_list = [ 1 ]

    proc_list = [ ]

    while True:
        try:
            id = server.recv_multipart()[0]

            # Get an unused port from the list
            # Ports from clients that have timed out are recycled here

            while not timeout_queue.empty():
                port_list.append(timeout_queue.get())

            port = port_list.pop()

            if len(port_list) == 0:
                port_list.append(port + 1)

            # Spawn a new worker task, binding the port to a socket

            proc_running = Value("b", True)

            proc_list.append(proc_running)

            Process(target=worker_task, args=(proc_running, port, timeout_queue)).start()

            # Send the new port to the client

            server.send_multipart([id, str(7777 + port)])

        except KeyboardInterrupt:
            break

    # Safely allow our worker processes to terminate
    for proc_running in proc_list:
        proc_running.value = False

    context.term()

def worker_task(proc_running, port, timeout_queue):
    context = zmq.Context.instance()

    worker = context.socket(zmq.ROUTER)

    worker.setsockopt(zmq.RCVTIMEO, 5000)
    worker.bind("tcp://127.0.0.1:%d" % (7777 + port, ))

    while proc_running.value:
        try:
            id, data = worker.recv_multipart()

            worker.send_multipart([id, data])
        except zmq.Again:
            timeout_queue.put(port)

            context.term()

            break

    print("Client on port %d disconnected" % (7777 + port, ))

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-04-09
    • 1970-01-01
    • 1970-01-01
    • 2012-04-17
    • 1970-01-01
    • 2017-11-29
    相关资源
    最近更新 更多