【问题标题】:How do I run pyzmq and a webserver in one ioloop?如何在一个 ioloop 中运行 pyzmq 和网络服务器?
【发布时间】:2017-06-09 19:03:44
【问题描述】:

我想编写一个单线程程序,该程序使用 Tornado 托管网络服务器,并在 ZMQ 套接字上接收消息(使用 PyZMQ Tornado 事件循环:http://learning-0mq-with-pyzmq.readthedocs.org/en/latest/pyzmq/multisocket/tornadoeventloop.html),但我不确定如何构建它。我应该使用

from zmq.eventloop import ioloop

from tornado.ioloop import IOLoop

或两者兼而有之?

【问题讨论】:

    标签: python tornado zeromq pyzmq


    【解决方案1】:

    在所有 Tornado 导入之前,您需要导入 zmq.eventloop.ioloop 并调用 zmq.eventloop.ioloop.install 函数。然后你可以导入 Tornado ioloop 并使用它。

    见: http://zeromq.github.io/pyzmq/eventloop.html

    【讨论】:

      【解决方案2】:

      这里是 Tornado HTTP 服务器的示例,带有 zeroMQ PUB SUB 套接字。

      #!/usr/bin/env python
      import json
      import tornado
      import tornado.web
      import zmq
      from tornado import httpserver
      from zmq.eventloop import ioloop
      from zmq.eventloop.zmqstream import ZMQStream
      
      ioloop.install()
      tornado.ioloop = ioloop
      import sys
      
      
      def ping_remote():
          """callback to keep the connection with remote server alive while we wait
          Network routers between raspberry pie and cloud server will close the socket
          if there is no data exchanged for long time.
          """
          pub_inst.send_json_data(msg="Ping", req_id="##")
          sys.stdout.write('.')
          sys.stdout.flush()
      
      
      pending_requests = {}
      
      
      class ZMQSub(object):
          def __init__(self, callback):
              self.callback = callback
              context = zmq.Context()
              socket = context.socket(zmq.SUB)
              # socket.connect('tcp://127.0.0.1:5559')
              socket.bind('tcp://*:8081')
              self.stream = ZMQStream(socket)
              self.stream.on_recv(self.callback)
              socket.setsockopt(zmq.SUBSCRIBE, "")
      
          def shutdown_zmq_sub(self):
              self.stream.close()
      
      
      class ZMQPub(object):
          def __init__(self):
              context = zmq.Context()
              socket = context.socket(zmq.PUB)
              socket.bind('tcp://*:8082')
              self.publish_stream = ZMQStream(socket)
      
          def send_json_data(self, msg, req_id):
              topic = str(req_id)
              self.publish_stream.send_multipart([topic, msg])
      
          def shutdown_zmq_sub(self):
              self.publish_stream.close()
      
      
      def SensorCb(msg):
          # decode message from raspberry pie and the channel ID.
          key, msg = (i for i in msg)
          if not key == "##":
              msg = json.loads(msg)
      
              if key in pending_requests.keys():
                  req_inst = pending_requests[key]
                  req_inst.write(msg)
                  req_inst.finish()
                  del pending_requests[key]
              else:
                  print "no such request"
                  print pending_requests
          else:
              print "received ping"
      
      
      class Handler(tornado.web.RequestHandler):
          def __init__(self, *args, **kwargs):
              super(Handler, self).__init__(*args, **kwargs)
      
              # get the unique req id
              self.req_id = str(self.application.req_id) + "#"
              self.application.req_id += 1
      
              # set headers
              self.set_header("Access-Control-Allow-Origin", "*")
              self.set_header("Access-Control-Allow-Headers", "x-requested-with")
              self.set_header('Access-Control-Allow-Methods', 'POST, GET, OPTIONS, PUT')
      
          @tornado.web.asynchronous
          def get(self):
              print self.request
              if self.req_id not in pending_requests.keys():
                  pending_requests[self.req_id] = self
              else:
                  print "WTF"
              pub_inst.send_json_data(msg=json.dumps({"op": "ServiceCall"}), req_id=self.req_id)
      
      
      if __name__ == "__main__":
          pub_inst = ZMQPub()
          sub_inst = ZMQSub(callback=SensorCb)
          application = tornado.web.Application(
              [(r'/get_sensor_data', Handler), (r'/(.*)')])
          application.req_id = 0
          server = httpserver.HTTPServer(application, )
          port = 8080
          server.listen(port)
          print "Sensor server ready on port: ", port
          ping = ioloop.PeriodicCallback(ping_remote, 3000)
          ping.start()
          tornado.ioloop.IOLoop.instance().start()
      

      【讨论】:

        猜你喜欢
        • 2014-02-28
        • 1970-01-01
        • 2017-07-07
        • 2019-05-20
        • 1970-01-01
        • 2018-06-06
        • 2010-12-08
        • 2023-03-07
        • 2010-09-27
        相关资源
        最近更新 更多