skip to content

How do you implement Redis pub/sub messaging in Spring using RedisMessageListenerContainer?

level: seniorimportance: should knowfreq 45%

answer

  1. convertAndSend to publish; container to receive
  2. ChannelTopic vs PatternTopic (PSUBSCRIBE)
  3. MessageListener.onMessage vs MessageListenerAdapter POJO
  4. dedicated subscription connection + task executor
  5. fire-and-forget: no ack, no persistence, no replay

basics

~10 s

Publish with redisTemplate.convertAndSend(channel, message). To receive, register a RedisMessageListenerContainer bean, give it the connection factory, and addMessageListener(listener, topic) where the listener implements MessageListener (often via a MessageListenerAdapter that delegates to a POJO).

solid answer

~40 s

Spring Data Redis wraps Redis SUBSCRIBE/PUBLISH. The receiving side centers on RedisMessageListenerContainer: you give it the RedisConnectionFactory and register listeners against topics — ChannelTopic for exact channel names or PatternTopic for glob patterns. A listener is either a MessageListener (onMessage(Message, byte[] pattern)) or a POJO wrapped in a MessageListenerAdapter that reflectively invokes a handler method and can deserialize the payload. The container owns a dedicated thread/connection for the blocking subscription and dispatches received messages on a task executor. Publishing is just redisTemplate.convertAndSend("channel", payload) (or the low-level connection.publish). It's fire-and-forget: no acknowledgement, no persistence, no replay — subscribers must be connected at publish time or the message is lost. Good for cache invalidation broadcasts and lightweight signaling, not durable work queues.

code

java · 28 lines
java
@Configuration
public class PubSubConfig {

    @Bean
    RedisMessageListenerContainer container(RedisConnectionFactory cf,
                                            MessageListenerAdapter adapter) {
        RedisMessageListenerContainer c = new RedisMessageListenerContainer();
        c.setConnectionFactory(cf);
        c.addMessageListener(adapter, new ChannelTopic("cache.invalidate"));
        return c;
    }

    @Bean
    MessageListenerAdapter adapter(CacheInvalidator invalidator) {
        // reflectively calls invalidator.handleMessage(String)
        MessageListenerAdapter a = new MessageListenerAdapter(invalidator, "handleMessage");
        a.setSerializer(new StringRedisSerializer());
        return a;
    }
}

@Component
class CacheInvalidator {
    public void handleMessage(String key) { /* evict key locally */ }
}

// Publisher side:
// redisTemplate.convertAndSend("cache.invalidate", "books::978-1");

go deeper

for a junior

Know convertAndSend publishes and RedisMessageListenerContainer + a MessageListener receives.

for a middle

Wire ChannelTopic/PatternTopic and a MessageListenerAdapter over a POJO with the right serializer.

for a senior

Explain the dedicated subscription connection, async dispatch, error handling, and the fire-and-forget delivery model.

for a principal

Choose pub/sub vs Streams vs a broker per delivery guarantees; design cache-invalidation fan-out and reason about ordering, back-pressure, and reconnect gaps.

**Redis pub/sub model.** Redis natively supports `PUBLISH channel message`, `SUBSCRIBE channel`, and `PSUBSCRIBE pattern`. It is a *broadcast* bus: every currently-subscribed client on a channel gets a copy. There is **no durability** — messages are not stored, not acknowledged, and not replayed. A subscriber that is down at publish time never sees the message. (Redis Streams — `XADD`/consumer groups — is the durable alternative and a different API.) **Publishing.** The high-level way is `RedisTemplate.convertAndSend(String channel, Object message)`, which serializes `message` with the template's value serializer and issues `PUBLISH`. Low-level: `connection.publish(byte[] channel, byte[] message)`. **Receiving — the container.** `RedisMessageListenerContainer` is a lifecycle bean (`SmartLifecycle`) that manages the subscription connection(s). Because `SUBSCRIBE` blocks the connection, the container holds a **dedicated** connection for subscriptions (separate from the pool used for normal commands) and runs the blocking receive on its own thread. Incoming messages are handed to a `taskExecutor` (a `SimpleAsyncTaskExecutor` by default) that invokes your listeners — so listener code runs off the subscription thread and multiple messages can be processed concurrently. Configure it minimally with `setConnectionFactory(cf)` and one or more `addMessageListener(listener, topic)` calls. **Topics.** `new ChannelTopic("orders")` subscribes to the exact channel; `new PatternTopic("orders.*")` uses `PSUBSCRIBE` glob matching. Pattern subscriptions are more expensive server-side. **Listener shapes.** 1. Implement `org.springframework.data.redis.connection.MessageListener` directly: `void onMessage(Message message, byte[] pattern)`. `message.getBody()` and `message.getChannel()` are raw `byte[]` — you deserialize yourself. 2. Use `MessageListenerAdapter`: wrap a plain POJO whose method (default `handleMessage`) receives the already-deserialized payload. Set its serializer so the adapter converts bytes to your type; set the delegate method name via constructor or `setDefaultListenerMethod`. **Threading & ordering.** Because dispatch is async on a task executor, **ordering across messages is not guaranteed** and a single logical channel may have its messages handled concurrently. If you need strict per-key ordering you must serialize handling yourself. Also, a slow listener can back up processing; keep handlers fast or hand off to a queue. **Error handling.** Set `setErrorHandler(...)` — uncaught exceptions in listeners otherwise just get logged. If the Redis connection drops, the container transparently re-subscribes on reconnect, but messages published during the gap are lost (pub/sub has no replay). **Serializer mismatch gotcha.** The bytes on the wire depend on the *publisher's* serializer. If you publish with `GenericJackson2JsonRedisSerializer` you must decode with a matching one on the listener/adapter, or you'll get garbage/`byte[]`. **When to use.** Great for broadcast/fan-out signals: distributed cache invalidation ('evict user 42 everywhere'), real-time notifications, config-change pings. **Not** for guaranteed delivery, work distribution, or anything needing acknowledgement/replay — reach for Redis Streams, or a real broker (Kafka/RabbitMQ) there.

  • What happens to messages published while a subscriber is disconnected?
    They are lost. Redis pub/sub does not persist or replay messages; only clients subscribed at publish time receive them. Use Redis Streams or a durable broker if you need delivery guarantees.
  • Why does the container need its own connection?
    SUBSCRIBE puts the connection into subscribe mode and blocks it from issuing normal commands, so the container dedicates a connection/thread to the subscription instead of borrowing from the command pool.

saying these in an interview costs you the question

  • Claiming Redis pub/sub guarantees delivery or persists messages (it does neither)
  • Expecting acknowledgements or consumer-group semantics (that's Redis Streams, not pub/sub)
  • Assuming listener invocations preserve message order (dispatch is async on a task executor)
  • Forgetting the listener must deserialize with the same serializer the publisher used

context