Publish / Subscribe#

Redis pub/sub lets one client broadcast a message and every client subscribed to that channel receive it. Nothing is stored: a message is delivered to whoever is subscribed at that moment, and is gone.

That property is the thing to keep in mind throughout this page. If a subscriber is not connected when a message is published, it does not get it later. For a stream you can replay, see Redis Streams instead.

Connecting#

[1]:
import redis

r = redis.Redis(decode_responses=True)
r.ping()
[1]:
True

A PubSub object is created from the client. It holds its own connection to the server.

[2]:
pubsub = r.pubsub()

Subscribing and reading messages#

subscribe() takes one or more channel names. The server confirms each subscription with a message of its own, which is why the first thing get_message() returns is of type subscribe rather than data.

[3]:
pubsub.subscribe('channel-one')
pubsub.get_message()
[3]:
{'type': 'subscribe', 'pattern': None, 'channel': 'channel-one', 'data': 1}

Now a message published on that channel is delivered. get_message() returns None when there is nothing waiting, so a timeout is usually what you want rather than a busy loop.

[4]:
r.publish('channel-one', 'hello')
pubsub.get_message(timeout=1)
[4]:
{'type': 'message', 'pattern': None, 'channel': 'channel-one', 'data': 'hello'}

The data field holds the payload. The channel field matters once you are subscribed to more than one.

Ignoring the subscribe confirmations#

Those subscribe and unsubscribe confirmations are rarely interesting to application code. ignore_subscribe_messages=True filters them out, so everything get_message() returns is real data.

[5]:
quiet = r.pubsub(ignore_subscribe_messages=True)
quiet.subscribe('channel-two')
print(quiet.get_message(timeout=1))  # the confirmation, filtered out

r.publish('channel-two', 'only data from here')
quiet.get_message(timeout=1)
None
[5]:
{'type': 'message',
 'pattern': None,
 'channel': 'channel-two',
 'data': 'only data from here'}

Handlers instead of return values#

A channel can be given a callable instead. Messages on that channel are passed to it, and get_message() returns None for them rather than handing them back.

[6]:
received = []

handled = r.pubsub()
handled.subscribe(**{'channel-three': received.append})
handled.get_message(timeout=1)  # the subscribe confirmation

r.publish('channel-three', 'goes to the handler')
handled.get_message(timeout=1)  # returns None: the handler took it

received
[6]:
[{'type': 'message',
  'pattern': None,
  'channel': 'channel-three',
  'data': 'goes to the handler'}]

Patterns#

psubscribe() matches channel names with glob-style patterns, so one subscription can cover a family of channels. Messages arrive with type pmessage and carry both the pattern that matched and the actual channel.

[7]:
patterned = r.pubsub(ignore_subscribe_messages=True)
patterned.psubscribe('news:*')
patterned.get_message(timeout=1)

r.publish('news:sport', 'a goal')
patterned.get_message(timeout=1)
[7]:
{'type': 'pmessage',
 'pattern': 'news:*',
 'channel': 'news:sport',
 'data': 'a goal'}

Reading in a loop with listen()#

listen() is a generator that blocks until the next message arrives. It is the natural shape for a worker whose whole job is to consume a channel.

Subscribe before you iterate. listen() runs while the connection has at least one subscription, so a loop started before any subscribe() call ends immediately instead of waiting for one - it yields nothing, blocks on nothing and raises nothing. The symptom is a consumer that silently never receives, which is hard to trace back to ordering.

[8]:
listener = r.pubsub(ignore_subscribe_messages=True)
listener.subscribe('channel-four')  # first, not after

r.publish('channel-four', 'one')
r.publish('channel-four', 'two')

seen = []
for message in listener.listen():
    seen.append(message['data'])
    if len(seen) == 2:
        break

seen
[8]:
['one', 'two']

Reading in the background#

run_in_thread() runs that loop on a thread for you. Every subscribed channel needs a handler, since there is no caller to return a message to.

It returns the thread, and stop() ends it. A thread left running holds its connection open for the life of the process.

[9]:
import time

collected = []

background = r.pubsub()
background.subscribe(**{'channel-five': collected.append})
thread = background.run_in_thread(sleep_time=0.01)

r.publish('channel-five', 'handled on the thread')
time.sleep(0.1)

thread.stop()
thread.join(timeout=1)

[message['data'] for message in collected]
[9]:
['handled on the thread']

Unsubscribing and closing#

unsubscribe() with no arguments drops every channel and punsubscribe() does the same for patterns, but neither is required before close(): closing resets the subscription state on its way out.

A PubSub that simply goes out of scope is not a leak. __del__ resets it, which disconnects the connection and releases it back to the pool. What leaks is one that stays referenced and never closed – held in a global, kept alive by a running run_in_thread worker, or caught in a reference cycle.

A with block settles the question, since it closes however the block ends:

[10]:
with r.pubsub() as scoped:
    scoped.subscribe('channel-six')
    scoped.get_message(timeout=1)  # the subscribe confirmation
    r.publish('channel-six', 'inside the block')
    print(scoped.get_message(timeout=1))

print('connection released on the way out')
{'type': 'message', 'pattern': None, 'channel': 'channel-six', 'data': 'inside the block'}
connection released on the way out

The ones opened earlier on this page are closed the same way.

[11]:
# close() is all that is needed. `background` is already closed: its worker
# thread calls close() when stop() ends the loop.
for p in (pubsub, quiet, handled, patterned, listener, background):
    p.close()

print('closed')
closed

What pub/sub does not do#

It does not store anything. A subscriber that is disconnected when a message is published never sees it, and there is no replay after a reconnect. If losing a message is not survivable, use Redis Streams.

It does not acknowledge. The publisher learns how many subscribers the message was delivered to - that is the integer publish() returns - and nothing about whether any of them processed it.