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?
在您现在完整的示例中,只需frontend将转发器更改为PULL套接字,将发布者套接字更改为PUSH,它应该按预期运行.
套接字选择的一般原则与此相关:
最初看起来你可能只想要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:backend是PULL,frontend是PUB
ws.py:SUB连接和订阅forwarder.frontend.
在您的情况下,PUB/SUB在后端失败的相关行为是慢速木匠综合症,"指南"中对此进行了描述.从本质上讲,订阅者花费有限的时间告诉发布者有关订阅的信息,因此如果您在打开PUB套接字后立即发送消息,可能性尚未被告知它还有任何订阅者,所以它只是丢弃消息.