使用现有龙卷​​风ioloop附加ZMQStream

Abh*_*ngh 4 python tornado zeromq

我有一个应用程序,其中每个websocket连接(在tornado打开回调内)创建一个zmq.SUB现有zmq.FORWARDER设备的套接字.想法是从zmq接收数据作为回调,然后可以通过websocket连接将其转发到前端客户端.

https://gist.github.com/abhinavsingh/6378134

ws.py

import zmq
from zmq.eventloop import ioloop
from zmq.eventloop.zmqstream import ZMQStream
ioloop.install()

from tornado.websocket import WebSocketHandler
from tornado.web import Application
from tornado.ioloop import IOLoop
ioloop = IOLoop.instance()

class ZMQPubSub(object):

    def __init__(self, callback):
        self.callback = callback

    def connect(self):
        self.context = zmq.Context()
        self.socket = self.context.socket(zmq.SUB)
        self.socket.connect('tcp://127.0.0.1:5560')
        self.stream = ZMQStream(self.socket)
        self.stream.on_recv(self.callback)

    def subscribe(self, channel_id):
        self.socket.setsockopt(zmq.SUBSCRIBE, channel_id)

class MyWebSocket(WebSocketHandler):

    def open(self):
        self.pubsub = ZMQPubSub(self.on_data)
        self.pubsub.connect()
        self.pubsub.subscribe("session_id")
        print 'ws opened'

    def on_message(self, message):
        print message

    def on_close(self):
        print 'ws closed'

    def on_data(self, data):
        print data

def main():
    application = Application([(r'/channel', MyWebSocket)])
    application.listen(10001)
    print 'starting ws on port 10001'
    ioloop.start()

if __name__ == '__main__':
    main()
Run Code Online (Sandbox Code Playgroud)

forwarder.py

import zmq

def main():
    try:
        context = zmq.Context(1)

        frontend = context.socket(zmq.SUB)
        frontend.bind('tcp://*:5559')
        frontend.setsockopt(zmq.SUBSCRIBE, '')

        backend = context.socket(zmq.PUB)
        backend.bind('tcp://*:5560')

        print 'starting zmq forwarder'
        zmq.device(zmq.FORWARDER, frontend, backend)
    except KeyboardInterrupt:
        pass
    except Exception as e:
        logger.exception(e)
    finally:
        frontend.close()
        backend.close()
        context.term()

if __name__ == '__main__':
    main()
Run Code Online (Sandbox Code Playgroud)

publish.py

import zmq

if __name__ == '__main__':
    context = zmq.Context()
    socket = context.socket(zmq.PUB)
    socket.connect('tcp://127.0.0.1:5559')
    socket.send('session_id helloworld')
    print 'sent data for channel session_id'
Run Code Online (Sandbox Code Playgroud)

但是,我的ZMQPubSub班级似乎根本没有收到任何数据.

我进一步尝试并意识到我需要ioloop.IOLoop.instance().start()在注册on_recv回调之后调用ZMQPubSub.但是,这只会阻止执行.

我也试过将main.ioloop实例传递给ZMQStream构造函数,但也没有帮助.

有没有办法可以绑定ZMQStream到现有main.ioloop实例而不阻塞流程MyWebSocket.open

min*_*nrk 5

在您现在完整的示例中,只需frontend将转发器更改为PULL套接字,将发布者套接字更改为PUSH,它应该按预期运行.

套接字选择的一般原则与此相关:

  • 当你想向准备好接收它的每个人发送消息时,可以使用PUB/SUB(可能没有人)
  • 如果要向一个对等体发送消息,等待它们准备就绪,请使用PUSH/PULL

最初看起来你可能只想要PUB-SUB,但是一旦你开始查看每个套接字对,你会发现它们非常不同.该frontend-websocket连接是绝对PUB-SUB -你可以有零到多个接收者,只是想和你发来的邮件通过发送消息给大家谁恰好是可用的.但是后端方面是不同的 - 只有一个接收器,它肯定希望来自发布者的每条消息.

所以你有它 - 后端应该是PULL和前端PUB.你所有的插座:

PUSH -> [PULL-PUB] -> SUB
Run Code Online (Sandbox Code Playgroud)

publisher.py:socket是PUSH,连接到backenddevice.py

forwarder.py:backendPULL,frontendPUB ws.py:SUB连接和订阅forwarder.frontend.

在您的情况下,PUB/SUB在后端失败的相关行为是慢速木匠综合症,"指南"中对此进行了描述.从本质上讲,订阅者花费有限的时间告诉发布者有关订阅的信息,因此如果您在打开PUB套接字后立即发送消息,可能性尚未被告知它还有任何订阅者,所以它只是丢弃消息.