【问题标题】:Limiting queue length with PyZMQ使用 PyZMQ 限制队列长度
【发布时间】:2012-03-12 05:12:58
【问题描述】:

我想限制我的 ZeroMQ 消息队列在 Python 应用程序中消耗的内存量。我知道设置 high-water mark 会限制发送方排队的数量,但有没有办法控制接收方排队的数量? Python ZeroMQ 绑定似乎将其设置为无限制。

我的测试场景:我有两个用于测试的 python 终端。一个是接收者:

Python 2.5.1 (r251:54863, Aug 25 2008, 20:50:04) 
[GCC 4.1.2 20071124 (Red Hat 4.1.2-42)] on linux2
Type "help", "copyright", "credits" or "license" for more information.

>>> import zmq
>>> context = zmq.Context()
>>> socket = context.socket(zmq.PULL)
>>> socket.setsockopt(zmq.RCVBUF, 256)
>>> socket.bind("tcp://127.0.0.1:12345")

另一个是发件人:

Python 2.5.1 (r251:54863, Aug 25 2008, 20:50:04) 
[GCC 4.1.2 20071124 (Red Hat 4.1.2-42)] on linux2
Type "help", "copyright", "credits" or "license" for more information.

>>> import zmq
>>> context=zmq.Context()
>>> socket = context.socket(zmq.PUSH)
>>> socket.setsockopt(zmq.SNDBUF, 2048)
>>> socket.setsockopt(zmq.HWM, 1)
>>> socket.connect("tcp://127.0.0.1:12345")
>>> num = 0
>>> while True:
...  print num
...  socket.send(str(num))
...  num = num + 1
... 

我在接收方运行socket.recv() 几次以确保队列正常工作,但除此之外,让两个终端就坐在那里。发送循环似乎永远不会阻塞,并且接收提示似乎占用了越来越多的内存。

【问题讨论】:

    标签: python sockets zeromq pyzmq


    【解决方案1】:

    与 ZeroMQ 的文档相矛盾的是,high water mark 需要在 PUSH 端和 PULL 端设置。一旦我更改了PULL,它的效果就更好了。新的PULL 代码是:

    Python 2.5.1 (r251:54863, Aug 25 2008, 20:50:04) 
    [GCC 4.1.2 20071124 (Red Hat 4.1.2-42)] on linux2
    Type "help", "copyright", "credits" or "license" for more information.
    
    >>> import zmq
    >>> context=zmq.Context()
    >>> socket = context.socket(zmq.PULL)
    >>> socket.setsockopt(zmq.RCVBUF, 256)
    >>> socket.setsockopt(zmq.HWM, 1)
    >>> socket.bind("tcp://127.0.0.1:12345")
    

    【讨论】:

    • 我使用了这个属性,但是报错:AttributeError: module 'zmq' has no attribute 'HWM'
    【解决方案2】:

    通过zmq.SNDBUFzmq.RCVBUF 选项,您可以set a limit on buffer size


    另外,我在接收方使用zmq.CONFLATE 选项将 ZeroMQ 队列大小限制为一个:

    这是一个 ZMQ PUSH/PULL 的示例:

    发件人(zmq.PUSH):

    def create_pub_socket(ip, port):
        try:
            context = zmq.Context()
            socket = context.socket(zmq.PUSH)
            socket.setsockopt(zmq.SNDHWM, 1)
            zmq_address = "tcp://{}:{}".format(ip, port)
            socket.connect(zmq_address)
            return socket
    
        except zmq.ZMQError as exp:
            print(exp)
            return False
    
    sock = create_push_socket('127.0.0.1', 5558)
    if sock:
        sock.send_json({'a': 1})
    

    Getter (zmq.PULL):

    def listen(self):
        sock = None
        try:
            context = zmq.Context()
            sock = context.socket(zmq.PULL)
            sock.setsockopt(zmq.RCVHWM, 1)
            sock.setsockopt(zmq.CONFLATE, 1)  # last msg only.
            sock.bind("tcp://*:5558")
    
        except zmq.ZMQError:
            logger.captureException()
    
        configs = None
        while configs is None:
            if sock:
                configs = sock.recv_json()
                time.sleep(1e-1)
            else:
                time.sleep(5)
                listen()  # Recursive.
    listen()
    

    【讨论】:

      【解决方案3】:

      实际上,文档是这样说的:

      "当一个 ZMQ_PUSH 套接字由于有 达到所有下游节点的高水位线,或者如果有 根本没有下游节点,然后任何 zmq_send(3) 操作在 套接字应阻塞,直到异常状态结束或至少一个 下游节点可用于发送;消息不是 丢弃。”

      http://api.zeromq.org/2-1:zmq-socket

      这完全表明您可以(并且应该)为下游节点(又名拉)设置高水位线,并且可能暗示将其设置在推送端将没有效果(尽管我怀疑这不是真的,因为有仍然是下游节点可用但消息进来的速度比发送速度快的情况。)

      【讨论】:

        猜你喜欢
        • 2022-10-18
        • 1970-01-01
        • 1970-01-01
        • 2017-10-18
        • 2018-02-14
        • 1970-01-01
        • 1970-01-01
        • 2013-07-25
        • 2012-08-17
        相关资源
        最近更新 更多