【问题标题】:Asynchronous Client/Server pattern in Python ZeroMQPython ZeroMQ 中的异步客户端/服务器模式
【发布时间】:2015-03-09 12:32:39
【问题描述】:

我有 3 个用 Python 编写的程序,需要连接。 2 个程序 X 和 Y 收集一些信息,由它们发送给程序 Z。程序 Z 分析数据并将一些决定发送给程序 X 和 Y。未来将扩展类似于 X 和 Y 的程序数量。最初我使用命名管道来允许从 X、Y 到 Z 的通信。但正如您所看到的,我需要双向关系。我的老板告诉我使用 ZeroMQ。我刚刚为我的用例找到了模式,称为异步客户端/服务器。请参阅下面的 ZMQ 书 (http://zguide.zeromq.org/py:all) 中的代码。

问题是我的老板不想使用任何线程、分叉等。我将客户端和服务器任务移动到单独的程序中,但我不确定如何处理 ServerWorker 类。这可以在没有线程的情况下以某种方式使用吗?另外,我想知道如何确定最佳工人数量。

import zmq
import sys
import threading
import time
from random import randint, random

__author__ = "Felipe Cruz <felipecruz@loogica.net>"
__license__ = "MIT/X11"

def tprint(msg):
    """like print, but won't get newlines confused with multiple threads"""
    sys.stdout.write(msg + '\n')
    sys.stdout.flush()

class ClientTask(threading.Thread):
    """ClientTask"""
    def __init__(self, id):
        self.id = id
        threading.Thread.__init__ (self)

    def run(self):
        context = zmq.Context()
        socket = context.socket(zmq.DEALER)
        identity = u'worker-%d' % self.id
        socket.identity = identity.encode('ascii')
        socket.connect('tcp://localhost:5570')
        print('Client %s started' % (identity))
        poll = zmq.Poller()
        poll.register(socket, zmq.POLLIN)
        reqs = 0
        while True:
            reqs = reqs + 1
            print('Req #%d sent..' % (reqs))
            socket.send_string(u'request #%d' % (reqs))
            for i in range(5):
                sockets = dict(poll.poll(1000))
                if socket in sockets:
                    msg = socket.recv()
                    tprint('Client %s received: %s' % (identity, msg))

        socket.close()
        context.term()

class ServerTask(threading.Thread):
    """ServerTask"""
    def __init__(self):
        threading.Thread.__init__ (self)

    def run(self):
        context = zmq.Context()
        frontend = context.socket(zmq.ROUTER)
        frontend.bind('tcp://*:5570')

        backend = context.socket(zmq.DEALER)
        backend.bind('inproc://backend')

        workers = []
        for i in range(5):
            worker = ServerWorker(context)
            worker.start()
            workers.append(worker)

        poll = zmq.Poller()
        poll.register(frontend, zmq.POLLIN)
        poll.register(backend,  zmq.POLLIN)

        while True:
            sockets = dict(poll.poll())
            if frontend in sockets:
                ident, msg = frontend.recv_multipart()
                tprint('Server received %s id %s' % (msg, ident))
                backend.send_multipart([ident, msg])
            if backend in sockets:
                ident, msg = backend.recv_multipart()
                tprint('Sending to frontend %s id %s' % (msg, ident))
                frontend.send_multipart([ident, msg])

        frontend.close()
        backend.close()
        context.term()

class ServerWorker(threading.Thread):
    """ServerWorker"""
    def __init__(self, context):
        threading.Thread.__init__ (self)
        self.context = context

    def run(self):
        worker = self.context.socket(zmq.DEALER)
        worker.connect('inproc://backend')
        tprint('Worker started')
        while True:
            ident, msg = worker.recv_multipart()
            tprint('Worker received %s from %s' % (msg, ident))
            replies = randint(0,4)
            for i in range(replies):
                time.sleep(1. / (randint(1,10)))
                worker.send_multipart([ident, msg])

        worker.close()

def main():
    """main function"""
    server = ServerTask()
    server.start()
    for i in range(3):
        client = ClientTask(i)
        client.start()

    server.join()

if __name__ == "__main__":
    main()

【问题讨论】:

    标签: python zeromq


    【解决方案1】:

    所以,您从这里获取了代码:Asynchronous Client/Server Pattern

    密切注意显示此代码所针对的模型的图像。特别是,请查看“图 38 - 异步服务器的详细信息”。 ServerWorker 类正在旋转 5 个“Worker”节点。在代码中,这些节点是线程,但您可以使它们完全独立的程序。在这种情况下,您的服务器程序(可能)将不负责启动它们,它们会单独启动并与您的服务器通信它们已准备好接收工作。

    您会在 ZMQ 示例中经常看到这一点,这是在单个可执行文件中的线程中模拟的多节点拓扑。这只是为了让阅读整个内容变得容易,并不总是打算这样使用。

    对于您的特定情况,让工作线程成为线程或将它们分解为单独的程序可能是有意义的......但如果这是您老板的业务需求,那么只需将它们分解为单独的程序即可。

    当然,要回答您的第二个问题,如果不了解他们将执行的工作量以及他们需要多快做出响应,就无法知道有多少工人是最佳的……您的目标是工人完成工作的速度比收到新工作的速度快。在许多情况下,这很有可能通过一个工人来完成。如果是这样,您可以让您的服务器本身成为工作人员,而只需跳过架构的整个“工作人员层”。为了简单起见,您应该从那里开始,然后进行一些负载测试,看看它是否真的能有效地处理您的工作负载。如果不是,请了解完成一项任务需要多长时间,以及任务执行的速度。假设一个工人可以在 15 秒内完成一项任务。那是每分钟 4 个任务。如果每分钟有 5 个任务,你需要 2 个工人,你会有一点增长空间。如果情况变化很大,那么您将不得不就资源与可靠性做出决定。

    在您走得更远之前,请务必阅读第 4 章,可靠的请求/回复模式,它将为处理异常提供一些见解,并可能为您提供更好的模式。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-04-04
      • 1970-01-01
      • 2014-08-20
      • 1970-01-01
      • 1970-01-01
      • 2015-06-26
      • 2012-08-07
      • 2019-04-08
      相关资源
      最近更新 更多