一、什么是 redis pub/sub
redis pub/sub(publish/subscribe) 是 redis 内置的发布/订阅消息机制。发布者向一个"频道(channel)"发送消息,所有订阅了该频道的客户端都能实时收到消息。
类比理解
广播电台模型: 发布者 = 电台主播 频道 = fm 103.7 订阅者 = 正在收听103.7的听众 主播说了一句话 → 所有正在听的人都能听到 如果你的收音机关着(不在线)→ 这句话你就永远听不到了(不持久化)
核心特点
| 特点 | 说明 |
|---|---|
| 实时性 | 发布后订阅者立即收到 |
| 不持久化 | 消息发完就没了,不会存储 |
| 无确认机制 | 不知道订阅者是否成功处理 |
| 不可回溯 | 不在线时发的消息,上线后拿不到 |
| 轻量级 | 无需额外中间件,redis 自带 |
| 无队列 | 不是消息队列,是实时广播 |
二、与消息队列(mq)的核心区别
redis pub/sub(广播模型):
发布者 ──► 频道 ──► 订阅者a(在线✓ → 收到)
──► 订阅者b(在线✓ → 收到)
──► 订阅者c(离线✗ → 丢失)
rabbitmq(队列模型):
生产者 ──► 队列 ──► 消费者a(在线✓ → 消费)
消费者a离线?消息在队列中等着,
等a上线了再消费(持久化)| 维度 | redis pub/sub | rabbitmq/kafka |
|---|---|---|
| 消息持久化 | ❌ 不持久化 | ✅ 持久化到磁盘 |
| 离线消费 | ❌ 离线就丢 | ✅ 上线后继续消费 |
| 消息确认 | ❌ 无ack | ✅ 有ack机制 |
| 消息堆积 | ❌ 不支持 | ✅ 可以堆积 |
| 重试机制 | ❌ 无 | ✅ 支持重试/死信 |
| 消费者分组 | ❌ 所有订阅者都收 | ✅ 同组只一个消费 |
| 性能 | 极高(内存操作) | 高(需要磁盘io) |
| 部署依赖 | 只需redis | 需要额外中间件 |
| 适用场景 | 实时通知、缓存同步 | 业务解耦、可靠投递 |
三、redis pub/sub 工作原理
┌──────────────────────────────────────────────────────────┐
│ redis server │
│ │
│ channel: "order.created" │
│ │ │
│ ├── subscriber 1 (连接中) ← 收到消息 │
│ ├── subscriber 2 (连接中) ← 收到消息 │
│ └── subscriber 3 (断开了) ← 收不到,消息丢失 │
│ │
│ channel: "stock.warning" │
│ │ │
│ └── subscriber 4 (连接中) ← 收到消息 │
│ │
└──────────────────────────────────────────────────────────┘
publisher ──publish "order.created" "{...}"──► redis
redis ──推送消息──► 所有订阅了 "order.created" 的客户端关键: 订阅者必须保持与 redis 的长连接。连接断了就收不到消息。
四、redis 命令
基础命令
# 订阅频道(客户端1)
subscribe order.created stock.warning
# 发布消息(客户端2)
publish order.created '{"orderid":"so001","userid":123}'
# 订阅模式匹配(通配符)
psubscribe order.* # 匹配 order.created、order.paid、order.shipped 等
psubscribe stock.* # 匹配 stock.warning、stock.low 等
命令演示
# 终端1:订阅者
127.0.0.1:6379> subscribe order.created
reading messages... (press ctrl-c to quit)
1) "subscribe"
2) "order.created"
3) (integer) 1
# 等待消息...
# 终端2:发布者
127.0.0.1:6379> publish order.created '{"orderid":"so001"}'
(integer) 1 # 返回值 = 收到消息的订阅者数量
# 终端1 立即收到:
1) "message"
2) "order.created"
3) "{\"orderid\":\"so001\"}"
五、java 代码示例
1. 使用 spring data redis
配置:
@configuration
public class redispubsubconfig {
@bean
public redismessagelistenercontainer redismessagelistenercontainer(
redisconnectionfactory connectionfactory) {
redismessagelistenercontainer container = new redismessagelistenercontainer();
container.setconnectionfactory(connectionfactory);
// 注册监听器
container.addmessagelistener(ordercreatedlistener(),
new patterntopic("order.*"));
container.addmessagelistener(stockwarninglistener(),
new channeltopic("stock.warning"));
return container;
}
@bean
public messagelisteneradapter ordercreatedlistener() {
return new messagelisteneradapter(new ordercreatedsubscriber(), "onmessage");
}
@bean
public messagelisteneradapter stockwarninglistener() {
return new messagelisteneradapter(new stockwarningsubscriber(), "onmessage");
}
}发布者:
@service
@slf4j
public class rediseventpublisher {
@resource
private stringredistemplate stringredistemplate;
/**
* 发布事件到redis频道.
*/
public void publish(string channel, object event) {
string message = json.tojsonstring(event);
log.info("redis发布消息, channel:{}, message:{}", channel, message);
stringredistemplate.convertandsend(channel, message);
}
}订阅者:
/**
* 订单创建事件订阅者.
*/
@slf4j
public class ordercreatedsubscriber implements messagelistener {
@override
public void onmessage(message message, byte[] pattern) {
string channel = new string(message.getchannel());
string body = new string(message.getbody());
log.info("收到redis消息, channel:{}, body:{}", channel, body);
try {
ordercreatedevent event = json.parseobject(body, ordercreatedevent.class);
// 处理业务逻辑
handleordercreated(event);
} catch (exception e) {
log.warn("处理redis消息失败, channel:{}", channel, e);
}
}
private void handleordercreated(ordercreatedevent event) {
log.info("处理订单创建事件, orderid:{}", event.getorderid());
// 业务逻辑...
}
}2. 使用注解方式(spring boot)
@configuration
public class redispubsubconfig {
@bean
public redismessagelistenercontainer container(
redisconnectionfactory factory,
ordereventsubscriber ordersubscriber,
cacheinvalidatesubscriber cachesubscriber) {
redismessagelistenercontainer container = new redismessagelistenercontainer();
container.setconnectionfactory(factory);
// 订阅多个频道
container.addmessagelistener(ordersubscriber, new channeltopic("order.created"));
container.addmessagelistener(ordersubscriber, new channeltopic("order.paid"));
container.addmessagelistener(cachesubscriber, new patterntopic("cache.invalidate.*"));
return container;
}
}
@component
@slf4j
public class ordereventsubscriber implements messagelistener {
@override
public void onmessage(message message, byte[] pattern) {
string channel = new string(message.getchannel());
string body = new string(message.getbody());
log.info("订单事件, channel:{}, body:{}", channel, body);
}
}
@component
@slf4j
public class cacheinvalidatesubscriber implements messagelistener {
@resource
private cachemanager cachemanager;
@override
public void onmessage(message message, byte[] pattern) {
string channel = new string(message.getchannel());
string cachekey = new string(message.getbody());
log.info("缓存失效通知, channel:{}, key:{}", channel, cachekey);
// 清除本地缓存
cachemanager.getcache("localcache").evict(cachekey);
}
}3. 完整业务示例:多实例缓存同步
/**
* 场景:应用部署了3个实例,一个实例更新了缓存,需要通知其他实例清除本地缓存.
*/
@service
@slf4j
public class cachesyncservice {
@resource
private stringredistemplate redistemplate;
@resource
private caffeinecachemanager localcachemanager;
private static final string cache_sync_channel = "cache.sync";
/**
* 数据更新后,发布缓存失效通知.
*/
public void notifycacheinvalidate(string cachename, string key) {
cachesyncmessage msg = new cachesyncmessage();
msg.setcachename(cachename);
msg.setkey(key);
msg.setsourceinstance(getinstanceid()); // 标记来源,避免自己处理自己发的消息
msg.settimestamp(system.currenttimemillis());
redistemplate.convertandsend(cache_sync_channel, json.tojsonstring(msg));
log.info("发布缓存同步通知, cache:{}, key:{}", cachename, key);
}
/**
* 收到其他实例的缓存失效通知.
*/
public void oncachesyncmessage(string messagebody) {
cachesyncmessage msg = json.parseobject(messagebody, cachesyncmessage.class);
// 忽略自己发的消息
if (getinstanceid().equals(msg.getsourceinstance())) {
return;
}
// 清除本地缓存
cache cache = localcachemanager.getcache(msg.getcachename());
if (cache != null) {
cache.evict(msg.getkey());
log.info("本地缓存已清除, cache:{}, key:{}, from:{}",
msg.getcachename(), msg.getkey(), msg.getsourceinstance());
}
}
private string getinstanceid() {
// 每个实例的唯一标识(如 ip:port 或 uuid)
return system.getproperty("instance.id", uuid.randomuuid().tostring());
}
}
@data
class cachesyncmessage {
private string cachename;
private string key;
private string sourceinstance;
private long timestamp;
}六、模式匹配订阅(pattern subscribe)
redis pub/sub 支持通配符订阅:
# 精确订阅 subscribe order.created # 只收 order.created # 模式订阅 psubscribe order.* # 收 order.created、order.paid、order.shipped... psubscribe *.warning # 收 stock.warning、price.warning... psubscribe * # 收所有频道的消息
代码中的模式订阅:
// 精确频道
container.addmessagelistener(listener, new channeltopic("order.created"));
// 模式匹配
container.addmessagelistener(listener, new patterntopic("order.*"));
container.addmessagelistener(listener, new patterntopic("cache.invalidate.*"));七、适用场景
✅ 适合用 redis pub/sub 的场景
| 场景 | 说明 |
|---|---|
| 多实例缓存同步 | 一个实例更新缓存,通知其他实例清除本地缓存 |
| 实时通知/聊天 | websocket 配合 pub/sub 实现多节点消息分发 |
| 配置变更广播 | 配置中心修改后通知所有应用实例热加载 |
| 在线状态通知 | 用户上线/下线通知给好友列表 |
| 实时监控仪表盘 | 指标数据实时推送到前端 |
共同特征: 消息丢了不要紧,要求实时性,不需要历史回溯。
❌ 不适合的场景
| 场景 | 原因 | 应该用什么 |
|---|---|---|
| 订单处理 | 消息不能丢 | rabbitmq/kafka |
| 异步任务 | 需要重试和确认 | rabbitmq |
| 日志采集 | 需要持久化和回溯 | kafka |
| 业务解耦 | 需要可靠投递 | rabbitmq |
| 延迟任务 | 需要延迟投递 | rabbitmq(延迟队列)/redis(sorted set) |
八、redis pub/sub 的局限性
1. 消息丢失
时间线: t1: 订阅者a在线,订阅者b离线 t2: 发布者发布消息 t3: 订阅者a收到 ✓,订阅者b收不到 ✗ t4: 订阅者b上线 t5: 订阅者b永远拿不到t2时刻的消息(已经丢了)
2. 无消息积压能力
如果订阅者处理速度跟不上发布速度,消息会堆积在 redis 的输出缓冲区中,当缓冲区超过限制时 redis 会断开订阅者的连接。
# redis配置中的保护机制 client-output-buffer-limit pubsub 32mb 8mb 60 # 含义:pubsub客户端输出缓冲区超过32mb,或持续60秒超过8mb,强制断开
3. 无消费确认
发布者不知道消息是否被成功处理:
发布者:publish channel "message" → 返回值只是"有几个订阅者收到了"
不是"有几个订阅者成功处理了"4. 集群模式下的限制
redis cluster 中,pub/sub 消息会广播到集群的所有节点,带来额外的网络开销。
九、redis pub/sub 的升级方案:redis stream
redis 5.0 引入了 stream,弥补了 pub/sub 的不足:
| 维度 | pub/sub | stream |
|---|---|---|
| 持久化 | ❌ | ✅ 消息持久化存储 |
| 消费确认 | ❌ | ✅ ack机制 |
| 消费者组 | ❌ | ✅ consumer group |
| 历史回溯 | ❌ | ✅ 可从任意位置消费 |
| 消息积压 | ❌ 会被丢弃 | ✅ 按需存储 |
# stream 基本用法 # 发布消息 xadd order.stream * orderid so001 userid 123 # 消费消息(消费者组) xreadgroup group notification-group consumer1 count 1 block 5000 streams order.stream > # 确认消费 xack order.stream notification-group 1234567890-0
如果你需要"redis + 可靠性",用 stream 而不是 pub/sub。
十、总结
redis pub/sub 是什么:
redis内置的实时广播机制
发布者发消息到频道,所有在线订阅者立即收到
核心特征:
实时性极高(毫秒级)
不持久化(发完就没了)
不可靠(离线就丢)
极轻量(无需额外中间件)
适合场景:
缓存同步、实时通知、配置广播
共同点 = "丢了无所谓,要的是实时"
不适合场景:
业务消息、订单处理、异步任务
共同点 = "消息不能丢"
选择建议:
需要可靠 → rabbitmq/kafka
需要实时 + 可靠 → redis stream
只需要实时广播 + 丢了无所谓 → redis pub/sub
到此这篇关于redis pub/sub中原理、场景与代码示例的文章就介绍到这了,更多相关redis pub/sub 内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!
发表评论