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.