【问题标题】:What is the proper way to handle Redis connection in Tornado ? (Async - Pub/Sub)在 Tornado 中处理 Redis 连接的正确方法是什么? (异步 - 发布/订阅)
【发布时间】:2012-01-18 00:08:02
【问题描述】:

我将 Redis 与我的 Tornado 应用程序与 asyc 客户端 Brukva 一起使用,当我查看 Brukva 站点上的示例应用程序时,他们正在 websocket 中的“init”方法上建立新连接

class MessagesCatcher(tornado.websocket.WebSocketHandler):
    def __init__(self, *args, **kwargs):
        super(MessagesCatcher, self).__init__(*args, **kwargs)
        self.client = brukva.Client()
        self.client.connect()
        self.client.subscribe('test_channel')

    def open(self):
        self.client.listen(self.on_message)

    def on_message(self, result):
        self.write_message(str(result.body))

    def close(self):
        self.client.unsubscribe('test_channel')
        self.client.disconnect()

在 websocket 的情况下很好,但是如何在常见的 Tornado RequestHandler post 方法中处理它说长轮询操作 (发布-订阅模型)。我正在更新处理程序的每个发布方法中建立新的客户端连接这是正确的方法吗?当我在 redis 控制台检查时,我发现客户端在每次新的发布操作中都在增加。

这是我的代码示例。

c = brukva.Client(host = '127.0.0.1')
c.connect()

class MessageNewHandler(BaseHandler):
    @tornado.web.authenticated
    def post(self):

        self.listing_id = self.get_argument("listing_id")
        message = {
            "id": str(uuid.uuid4()),
            "from": str(self.get_secure_cookie("username")),
            "body": str(self.get_argument("body")),
        }
        message["html"] = self.render_string("message.html", message=message)

        if self.get_argument("next", None):
            self.redirect(self.get_argument("next"))
        else:
            c.publish(self.listing_id, message)
            logging.info("Writing message : " + json.dumps(message))
            self.write(json.dumps(message))

    class MessageUpdatesHandler(BaseHandler):
        @tornado.web.authenticated
        @tornado.web.asynchronous
        def post(self):
            self.listing_id = self.get_argument("listing_id", None)
            self.client = brukva.Client()
            self.client.connect()
            self.client.subscribe(self.listing_id)
            self.client.listen(self.on_new_messages)

        def on_new_messages(self, messages):
            # Closed client connection
            if self.request.connection.stream.closed():
                return
            logging.info("Getting update : " + json.dumps(messages.body))
            self.finish(json.dumps(messages.body))
            self.client.unsubscribe(self.listing_id)


        def on_connection_close(self):
            # unsubscribe user from channel
            self.client.unsubscribe(self.listing_id)
            self.client.disconnect()

感谢您提供一些类似案例的示例代码。

【问题讨论】:

标签: python redis tornado publish-subscribe


【解决方案1】:

有点晚了,但我一直在使用tornado-redis。它适用于 tornado 的 ioloop 和 tornado.gen 模块

安装 tornadoredis

可以从pip安装

pip install tornadoredis

或使用设置工具

easy_install tornadoredis

但你真的不应该那样做。您还可以克隆存储库并提取它。然后运行

python setup.py build
python setup.py install

连接到 redis

以下代码位于您的 main.py 或等效文件中

redis_conn = tornadoredis.Client('hostname', 'port')
redis_conn.connect()

redis.connect 只被调用一次。这是一个阻塞调用,因此应该在启动主 ioloop 之前调用它。所有处理程序之间共享相同的连接对象。

您可以将其添加到您的应用程序设置中,例如

settings = {
    redis = redis_conn
}
app = tornado.web.Application([('/.*', Handler),],
                              **settings)

使用 tornadoredis

连接可以在处理程序中作为self.settings['redis'] 使用,也可以作为 BaseHandler 类的属性添加。您的请求处理程序子类化该类并访问该属性。

class BaseHandler(tornado.web.RequestHandler):

    @property
    def redis():
        return self.settings['redis']

为了与 redis 通信,使用了 tornado.web.asynchronoustornado.gen.engine 装饰器

class SomeHandler(BaseHandler):

    @tornado.web.asynchronous
    @tornado.gen.engine
    def get(self):
        foo = yield gen.Task(self.redis.get, 'foo')
        self.render('sometemplate.html', {'foo': foo}

额外信息

更多示例和其他功能(如连接池和管道)可以在 github 存储库中找到。

【讨论】:

    【解决方案2】:

    您应该在您的应用中汇集连接。因为似乎 brukva 不自动支持这一点(redis-py 支持这一点,但本质上是阻塞的,所以它不适合龙卷风),你需要编写自己的连接池。

    不过,模式非常简单。类似这样的东西(这不是真正的操作代码):

    class BrukvaPool():
    
        __conns = {}
    
    
        def get(host, port,db):
            ''' Get a client for host, port, db '''
    
            key = "%s:%s:%s" % (host, port, db)
    
            conns = self.__conns.get(key, [])
            if conns:
                ret = conns.pop()
                return ret
            else:
               ## Init brukva client here and connect it
    
        def release(client):
            ''' release a client at the end of a request '''
            key = "%s:%s:%s" % (client.connection.host, client.connection.port, client.connection.db)
            self.__conns.setdefault(key, []).append(client)
    

    这可能有点棘手,但这是主要思想。

    【讨论】:

      猜你喜欢
      • 2013-02-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-12-31
      • 2015-10-26
      • 2012-03-29
      • 1970-01-01
      • 2017-05-30
      相关资源
      最近更新 更多