【问题标题】:Is non-blocking Redis pubsub possible?非阻塞 Redis pubsub 可能吗?
【发布时间】:2011-12-13 20:46:30
【问题描述】:

我想用redis的pubsub传输一些消息,但不想被listen屏蔽,如下代码:

import redis
rc = redis.Redis()

ps = rc.pubsub()
ps.subscribe(['foo', 'bar'])

rc.publish('foo', 'hello world')

for item in ps.listen():
    if item['type'] == 'message':
        print item['channel']
        print item['data']

最后一个 for 部分将被阻止。我只想检查给定频道是否有数据,我该如何完成?有checklike 方法吗?

【问题讨论】:

  • 您是否有不想被监听阻止的原因? Redis 连接非常便宜,通常会生成多个连接。
  • 在 Python 中使用 Redis、ZMQ、Tornado 实现异步 PubSub - github.com/abhinavsingh/async_pubsub
  • 使用 pubsub 对象的 .get_message() 方法而不是 .listen() (下面有一个示例)。 [发布此问题时,Python Redis 驱动程序可能不支持该方法]。

标签: python redis redis-py


【解决方案1】:

如果您正在考虑非阻塞、异步处理,您可能正在使用(或应该使用)异步框架/服务器。

更新: 距离最初的答案已经 5 年了,与此同时 Python 得到了native async IO support。现在有AIORedis, an async IO Redis client。

【讨论】:

  • 这是应该勾选的正确答案。我不知道为什么人们会重新发明轮子,redis 已经有一个异步客户端,在这样的客户端存在的情况下并不真正需要生成一个新线程。
  • @securecurve:公平地说,我已经在标记答案一年多之后添加了该答案。但是,txRedis 和 brükva(Tornado-Redis 从中分叉)都已经 3 岁了,所以这并不是一个真正的借口。
  • 这个答案与问题无关。正如在接受的答案中指出的那样,redis 将消息推送给正在监听的客户端。因此无法请求消息。
  • @Glaslos:有办法不阻止收听新消息。 正是问的问题。
  • 我不同意,redis pub/sub 通道上的监听调用总是会阻塞,即使你在线程/greenlet 中运行它并在阻塞时切换上下文。
【解决方案2】:

已接受的答案已过时,因为 redis-py 建议您使用非阻塞 get_message()。但它也提供了一种轻松使用线程的方法。

https://pypi.python.org/pypi/redis

阅读消息有三种不同的策略。

在幕后,get_message() 使用系统的“选择”模块快速轮询连接的套接字。如果有数据可供读取,get_message() 将读取它,格式化消息并将其返回或将其传递给消息处理程序。如果没有要读取的数据,get_message() 将立即返回 None。这使得集成到应用程序内的现有事件循环中变得微不足道。

 while True:
     message = p.get_message()
     if message:
         # do something with the message
     time.sleep(0.001)  # be nice to the system :)

旧版本的 redis-py 只能使用 pubsub.listen() 读取消息。 listen() 是一个阻塞直到消息可用的生成器。如果您的应用程序除了接收和处理从 redis 收到的消息之外不需要做任何其他事情,listen() 是一种启动运行的简单方法。

 for message in p.listen():
     # do something with the message

第三个选项在单独的线程中运行事件循环。 pubsub.run_in_thread() 创建一个新线程并启动事件循环。线程对象返回给 run_in_thread() 的调用者。调用者可以使用 thread.stop() 方法来关闭事件循环和线程。在幕后,这只是一个围绕 get_message() 的包装器,它在单独的线程中运行,本质上为您创建了一个微小的非阻塞事件循环。 run_in_thread() 采用可选的 sleep_time 参数。如果指定,事件循环将使用循环的每次迭代中的值调用 time.sleep()。

注意:由于我们在单独的线程中运行,因此无法处理注册消息处理程序无法自动处理的消息。因此,如果您订阅了没有附加消息处理程序的模式或通道,redis-py 会阻止您调用 run_in_thread()。

p.subscribe(**{'my-channel': my_handler})
thread = p.run_in_thread(sleep_time=0.001)
# the event loop is now running in the background processing messages
# when it's time to shut it down...
thread.stop()

所以要回答您的问题,只需在您想知道消息是否到达时检查 get_message。

【讨论】:

  • 这里的消息阅读在一个循环中。为什么订阅者不会多次收到相同的消息(因为它在循环中)?发布者是否会跟踪哪个订阅者已经收到了消息?
【解决方案3】:

新版 redis-py 支持异步发布订阅,详情请查看https://github.com/andymccurdy/redis-py。 以下是文档本身的示例:

while True:
    message = p.get_message()
    if message:
        # do something with the message
    time.sleep(0.001)  # be nice to the system :)

【讨论】:

    【解决方案4】:

    我认为这是不可能的。 Channel 没有任何“当前数据”,您订阅了一个频道并开始接收其他客户端在该频道上推送的消息,因此它是一个阻塞 API。此外,如果您查看 Redis Commands documentation 的 pub/sub 会更清楚。

    【讨论】:

    • 我认为这个答案与另一个答案结合起来非常完整。他可以把它放到一个线程中。如果他不想在 chanel 有活动时立即采取行动,那么他可以将其存储在一个字典中,并拥有自己的检查方法,该方法使用锁定互斥锁查看字典
    【解决方案5】:

    这是一个线程阻塞侦听器的工作示例。

    import sys
    import cmd
    import redis
    import threading
    
    
    def monitor():
        r = redis.Redis(YOURHOST, YOURPORT, YOURPASSWORD, db=0)
    
        channel = sys.argv[1]
        p = r.pubsub()
        p.subscribe(channel)
    
        print 'monitoring channel', channel
        for m in p.listen():
            print m['data']
    
    
    class my_cmd(cmd.Cmd):
        """Simple command processor example."""
    
        def do_start(self, line):
            my_thread.start()
    
        def do_EOF(self, line):
            return True
    
    
    if __name__ == '__main__':
        if len(sys.argv) == 1:
            print "missing argument! please provide the channel name."
        else:
            my_thread = threading.Thread(target=monitor)
            my_thread.setDaemon(True)
    
            my_cmd().cmdloop()
    

    【讨论】:

    【解决方案6】:

    这是一个没有线程的非阻塞解决方案:

    fd = ps.connection._sock.fileno();
    rlist,, = select.select([fd], [], [], 0) # or replace 0 with None to block
    if rlist:
        for rfd in rlist:
            if fd == rfd:
                message = ps.get_message()
    

    ps.get_message() 本身就足够了,但是我使用这种方法,以便我可以等待多个 fd 而不仅仅是 redis 连接。

    【讨论】:

      【解决方案7】:

      要达到无阻塞代码,您必须执行另一种范例代码。这并不难,使用一个新线程来监听所有的变化,让主线程去做其他事情。

      此外,您将需要一些机制来在主线程和 redis 订阅者线程之间交换数据。

      【讨论】:

        【解决方案8】:

        最有效的方法是基于 greenlet 而不是基于线程。作为一个基于 greenlet 的并发框架,gevent 在 Python 世界中已经相当成熟。因此,gevent 与 redis-py 的集成会很棒。这正是github上本期讨论的内容:

        https://github.com/andymccurdy/redis-py/issues/310

        【讨论】:

          【解决方案9】:

          你可以使用gevent、gevent猴子补丁来构建一个非阻塞的redis pubsub应用。

          【讨论】:

            【解决方案10】:

            Redis 的 pub/sub 向频道上订阅(收听)的客户端发送消息。如果您不听,您将错过消息(因此阻塞呼叫)。如果你想让它不阻塞,我建议改用队列(redis 也很擅长)。如果您必须使用 pub/sub,您可以按照建议使用 gevent 来拥有一个异步阻塞侦听器,将消息推送到队列并使用单独的使用者以非阻塞方式处理来自该队列的消息。

            【讨论】:

              【解决方案11】:

              这很简单。我们检查消息是否存在,并继续进行订阅,直到处理完所有消息。

              import redis
              
              r = redis.Redis(decode_responses=True)
              subscription = r.pubsub()
              subscription.psubscribe('channel')
              
              r.publish('channel', 'foo')
              r.publish('channel', 'bar')
              r.publish('channel', 'baz')
              
              message = subscription.get_message()
              while message is not None:
                if message['data'] != 1:
                    # Do something with message
                    print(message)
                # Get next message
                message = subscription.get_message()
              

              【讨论】:

                猜你喜欢
                • 2013-04-10
                • 1970-01-01
                • 1970-01-01
                • 1970-01-01
                • 2015-11-27
                • 1970-01-01
                • 1970-01-01
                • 1970-01-01
                • 2018-06-18
                相关资源
                最近更新 更多