【问题标题】:Random freezing / hanging in Python ZeroMQPython ZeroMQ 中的随机冻结/挂起
【发布时间】:2015-03-23 13:52:15
【问题描述】:

我正在使用 ZeroMQ 编写一个用 python 编写的无代理、平衡、客户端工作者服务。

客户端获取worker的地址,建立连接(zmq.REQ / zmq.REP),发送单个请求,接收单个响应然后断开连接。

我选择了无代理架构,因为需要在客户端和工作人员之间传输的数据量相对较大,尽管每个连接只有一个 REQ/REP 对,并且使用代理作为“中间人”会造成瓶颈。

在测试系统时,我注意到客户端和工作人员之间的通信随机停止,有时会在几秒钟后恢复(通常是几分钟)。

我将问题范围缩小到客户对工人的.connect() / .disconnect()

我编写了两个重现该错误的小型 Python 脚本。

import zmq

class Site:

      def __init__(self):
        ctx = zmq.Context()
        self.pair_socket = ctx.socket(zmq.REQ)
        self.num = 0


      def __del__(self):
        print "closed"


      def run_site(self):
        print "running..."
        while True:
            self.pair_socket.connect('tcp://127.0.0.1:5555')
            print 'connected'
            self.pair_socket.send_pyobj(self.num)
            print 'sent', self.num
            print self.pair_socket.recv_pyobj()
            self.pair_socket.disconnect('tcp://127.0.0.1:5555')
            print 'disconnected'
            self.num += 1

s = Site()
s.run_site()

import zmq

class Server:

      def __init__(self):
          ctx = zmq.Context()
          self.pair_socket = ctx.socket(zmq.REP)
          self.pair_socket.bind('tcp://127.0.0.1:5555')


      def __del__(self):
          print " closed"


      def run_server(self):
          print "running..."
          while True:
              x =  self.pair_socket.recv_pyobj()
              print x
              self.pair_socket.send_pyobj(x)


s = Server()  
s.run_server()

我认为问题与内存或 gc 无关,因为我已尝试禁用 gc - 没有太大影响。

我已尝试使用 zmq.LINGER,如下所述:Zeromq with python hangs if connecting to invalid socket

什么可能导致这些随机数冻结?

【问题讨论】:

  • 使用数据包嗅探器查看这对的哪一侧挂起...消息在发送之前挂在客户端上,还是在收到消息后挂在服务器上。您如何确定冻结,即冻结开始前您看到的最后一条消息是什么?
  • 恕我直言,就底层资源和相关系统开销而言,设计 while True: .connect(); ... ; .disconnect() 是一种残酷 方式。肯定有更好和“更生态”/“更环保”的方式来用代码表达您的设计意图,不会浪费 CPU / 资源分配

标签: python zeromq pyzmq


【解决方案1】:

REP 套接字根据定义是同步的。所以你的服务器一次只能处理一个请求,其余的只会填满缓冲区并在某个时候丢失。

要解决根本原因,您需要改用ROUTER 套接字。

class Server:
    def __init__(self):
        ctx = zmq.Context()
        self.pair_socket = ctx.socket(zmq.ROUTER)
        self.pair_socket.bind('tcp://127.0.0.1:5555')
        self.poller = zmq.Poller()
        self.poller.register(self.pair_socket, zmq.POLLIN)

    def __del__(self):
        print " closed"

    def run_server(self):
        print "running..."
        while True:
            try:
                items = dict(self.poller.poll())
            except KeyboardInterrupt:
                break
            if self.pair_socket in items:
                x = self.pair_socket.recv_multipart()
                print x
                self.pair_socket.send_multipart(x)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2011-02-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-02-28
    • 1970-01-01
    相关资源
    最近更新 更多