讲解一下Redis中的订阅发布模式

发布与订阅模型在许多编程语言中都有实现,也就是我们经常说的设计模式中的一种–观察者模式。在一些应用场合,例如发送方并不是以固定频率发送消息,如果接收方频繁去咨询发送方,这种操作无疑是很麻烦并且不友好的。

为什么做订阅分布?

随着业务复杂, 业务的项目依赖关系增强, 使用消息队列帮助系统降低耦合度.

  • 订阅分布本身也是一种生产者消费者模式, 订阅者是消费者, 发布者是生产者.
  • 订阅发布模式, 发布者发布消息后, 只要有订阅方, 则多个订阅方会收到同样的消息
  • 生产者消费者模式, 生产者往队列里放入消息, 由多个消费者对一条消息进行抢占.
  • 订阅分布模式可以将一些不着急完成的工作放到其他进程或者线程中进行离线处理.

Redis中的订阅发布

Redis中的订阅发布模式, 当没有订阅者时, 消息会被直接丢弃(Redis不会持久化保存消息)

Redis生产者消费者

生产者使用Redis中的list数据结构进行实现, 将待处理的消息塞入到消息队列中.

class Producer(object):

def __init__(self, host="localhost", port=6379):
self._conn = redis.StrictRedis(host=host, port=port)
self.key = "test_key"
self.value = "test_value_{id}"

def produce(self):
for id in xrange(5):
msg = self.value.format(id=id)
self._conn.lpush(self.key, msg)

消费者使用redis中brpop进行实现, brpop会从list头部消息, 并能够设置超时等待时间.

class Consumer(object):

def __init__(self, host="localhost", port=6379):
self._conn = redis.StrictRedis(host=host, port=port)
self.key = "test_key"

def consume(self, timeout=0):
# timeout=0 表示会无线阻塞, 直到获得消息while True:
msg = self._conn.brpop(self.key, timeout=timeout)
process(msg)


def process(msg):
print msg

if __name__ == '__main__':
consumer = Consumer()
consumer.consume()
# 输出结果
('test_key''test_value_1')
('test_key''test_value_2')
('test_key''test_value_3')
('test_key''test_value_4')
('test_key''test_value_5')

Redis中订阅发布

在Redis Pubsub中, 一个频道(channel)相当于一个消息队列

class Publisher(object):

def __init__(self, host, port):
self._conn = redis.StrictRedis(host=host, port=port)
self.channel = "test_channel"
self.value = "test_value_{id}"

def pub(self):
for id in xrange(5):
msg = self.value.format(id=id)
self._conn.publish(self.channel, msg)

其中get_message使用了select IO多路复用来检查socket连接是否是否可读.

class Subscriber(object):

def __init__(self, host="localhost", port=6379):
self._conn = redis.StrictRedis(host=host, port=port)
self._pubsub = self._conn.pubsub() # 生成pubsub对象
self.channel = "test_channel"
self._pubsub.subscribe(self.channel)

def sub(self):
while True:
msg = self._pubsub.get_message()
if msg and isinstance(msg.get("data"), basestring):
process(msg.get("data"))

def close(self):
self._pubsub.close()

# 输出结果
test_value_1
test_value_2
test_value_3
test_value_4
test_value_

Java Jedis踩过的坑

在Jedis中订阅方处理是采用同步的方式, 看源码中PubSub模块的process函数

do-while循环中, 会等到当前消息处理完毕才能够处理下一条消息, 这样会导致当入队列消息量过大的时候, redis链接被强制关闭.

解决方案: 将整个处理函数改为异步的方式.

文章来源网络,作者:管理,如若转载,请注明出处:https://shuyeidc.com/wp/223403.html<

(0)
管理的头像管理
上一篇2025-04-15 22:57
下一篇 2025-04-15 22:58

相关推荐

  • 站群服务器如何批量管理更高效,有哪些管理技巧?

    站群服务器批量管理想提效,自动化是唯一出路,通过统一配置管理工具与面板系统,结合服务商提供的底层基础设施支持,能将运维效率提升数倍,批量管理的核心痛点与解决思路多台站群服务器分散管理,最常见的问题就是重复劳动,每次软件更新、配置修改、安全加固,都需要逐台登录操作,不仅耗时,还容易漏掉某台机器,更头疼的是,一旦某……

    2026-07-27
    0
  • 服务器磁盘IO过高如何优化?,磁盘IO过高的原因有哪些?

    服务器磁盘IO过高,核心优化路径是“先定位、再分流、后升级”,你需要通过系统工具精确判断究竟是应用程序、日志策略还是硬件瓶颈导致,然后针对性地从代码、缓存、存储架构和硬件选型四个层面下手,其中选择持有持牌自营机房和增值电信业务经营许可证的服务商,能从根本上保障底层IO稳定性,定位IO瓶颈:动手优化的第一步盲目优……

    2026-07-27
    0
  • 跨境网站访问延迟高怎么解决,网站访问慢的原因是什么?

    跨境网站访问延迟高的核心解决思路在于多维度优化网络路径,包括使用全球CDN加速、选择靠近目标区域的优质IDC机房、调整传输协议以及精简应用层资源,其中服务商的基础设施质量直接决定优化上限,为什么跨境访问延迟高?三大核心因素物理距离与光速限制数据包在海底光缆中的传输速度受限于介质,从中国到美国西海岸的物理往返时间……

    2026-07-27
    0
  • 站群服务器到底是什么意思,怎么选择比较好

    站群服务器就是一台拥有多个独立IP地址、专门用于托管和管理多个网站的高性能服务器,其核心价值在于通过独立IP降低网站间的关联风险,并提升搜索引擎优化效果,站群服务器的工作原理与适用场景站群服务器本质上是将一台物理服务器通过虚拟化或直接配置的方式,分配给多个独立IP地址,每个IP对应一个独立的网站,这些网站共享服……

    2026-07-27
    0
  • 高防服务器误封正常流量如何调整,怎么解决?

    高防服务器误封正常流量,核心调整思路是从“一刀切”转向“精细化”——通过分析业务特征,调整防护阈值、配置白名单和启用智能学习模式,让防护系统学会区分真假流量,为什么会误封正常流量误封主要源于防护策略的通用化,高防服务器通常默认启用严格防护规则,当流量特征与攻击特征库部分匹配时,就会被拦截,据行业安全白皮书指出……

    2026-07-27
    0

发表回复

您的邮箱地址不会被公开。必填项已用 * 标注