你真的了解Redis的发布订阅?(含Java版实现源码)

Redis发布订阅使用场景及JAVA代码实现(含源码)html

导语

Redis是咱们很经常使用的一款nosql数据库产品,咱们一般会用Redis来配合关系型数据库一块儿使用,弥补关系型数据库的不足。java

其中,Redis的发布订阅功能也是它的一大亮点。虽然它不是一款专门作发布订阅的产品,但其自带的发布订阅功能已经知足咱们平常需求。web

那Redis的发布订阅功能的原理和它均可以用在哪些场景呢?今天咱们就来探讨一下这个问题。redis

什么是发布订阅

所谓发布订阅,就是消息发布者发布消息及消息订阅者接收消息,两者经过某种媒介关联起来。这相似之前的『订报』,当咱们订阅了某种报纸后(好比财经报),每当报纸有新的期刊出版后,就会有邮递员给咱们送过来。即,只有定了这种报纸才会收到出版社发布的这种新报纸。sql

Redis的发布订阅功能也是相似,首先要有消息的发布者,其次要有消息的订阅者。有了消息发布者和订阅者以后,还缺乏什么?数据库

那就是上述的『某种报纸』,并非出版社出版的每一种报纸(如人民日报,财经报,体育报)都给你送过来,而是明确你要定哪种,你定了哪种才给你送哪种。编程

回到Redis的发布订阅上,上述的『某种报纸』就抽象为频道channel,客户端订阅了某channel后,当发布者经过此channel发布消息时,全部订阅者就会收到该频道发布的消息。缓存

发布和订阅机制

当一个客户端经过 PUBLISH 命令向订阅者发送信息的时候,咱们称这个客户端为发布者(publisher)。

而当一个客户端使用 SUBSCRIBE 或者 PSUBSCRIBE命令接收信息的时候,咱们称这个客户端为订阅者(subscriber)。

为了解耦发布者(publisher)和订阅者(subscriber)之间的关系,Redis 使用了 channel (频道)做为二者的中介 —— 发布者将信息直接发布给 channel ,而 channel 负责将信息发送给适当的订阅者,发布者和订阅者之间没有相互关系,也不知道对方的存在。app

如上图所示, Redis client ARedis client B 订阅了 channel-> Financial newspapers,当 Redis client C经过 channel-> Financial newspapers 发布消息 Stocks are up today! 时, Redis client ARedis client B 就会收到该消息。

原理

Redis是使用C实现的,经过分析 Redis 源码里的 pubsub.c 文件,了解发布和订阅机制的底层实现,籍此加深对 Redis 的理解。异步

Redis 经过 PUBLISH 、SUBSCRIBE 和 PSUBSCRIBE 等命令实现发布和订阅功能。

经过 SUBSCRIBE 命令订阅某频道后,redis-server 里维护了一个字典,字典的键就是一个个 channel ,而字典的值则是一个链表,链表中保存了全部订阅这个 channel 的客户端。SUBSCRIBE 命令的关键,就是将客户端添加到给定 channel 的订阅链表中。

经过 PUBLISH 命令向订阅者发送消息,redis-server 会使用给定的频道做为键,在它所维护的 channel 字典中查找记录了订阅这个频道的全部客户端的链表,遍历这个链表,将消息发布给全部订阅者。

详细参考:Redis 发布/订阅机制原理分析

业务场景

明确了Redis发布订阅的原理和基本流程后,咱们来看一下Redis的发布订阅到底具体能作什么。

一、异步消息通知

好比渠道在调支付平台的时候,咱们能够用回调的方式给支付平台一个咱们的回调接口来通知咱们支付状态,还能够利用Redis的发布订阅来实现。好比咱们发起支付的同时订阅频道pay_notice_ + wk (假如咱们的渠道标识是wk,不能让其余渠道也订阅这个频道),当支付平台处理完成后,支付平台往该频道发布消息,告诉频道的订阅者该订单的支付信息及状态。收到消息后,根据消息内容更新订单信息及后续操做。

当不少人都调用支付平台时,支付时都去订阅同一个频道会有问题。好比用户A支付完订阅频道pay_notice_wk,在支付平台未处理完时,用户B支付完也订阅了pay_notice_wk,当A收到通知后,接着B的支付通知也发布了,这时渠道收不到第二次消息发布。由于同一个频道收到消息后,订阅自动取消,也就是订阅是一次性的。

因此咱们订阅的订单支付状态的频道就得惟一,一个订单一个频道,咱们能够在频道上加上订单号pay_notice_wk+orderNo保证频道惟一。这样咱们能够把频道号在支付时当作参数一并传过去,支付平台处理完就能够用此频道发布消息给咱们了。(实际大多使用接口回调通知的方式,由于用Redis发布订阅限制条件苛刻,系统间必须共用一套Redis)

二、任务通知

好比经过跑批系统通知应用系统作一些事(跑批系统没法拿到用户数据,且应用系统又不能作定时任务的状况下)。如天天凌晨3点提早加载一些用户的用户数据到Redis,应用系统不能作定时任务,能够经过系统公共的Redis来由跑批系统发布任务给应用系统,应用系统收到指令,去作相应的操做。

这里须要注意的是在线上集群部署的状况下,全部服务实例都会收到通知,都要作一样的操做吗?彻底不必。能够用Redis实现锁机制,其中一台实例拿到锁后执行任务。另外若是任务比较耗时,能够不用锁,能够考虑一下任务分片执行。固然这不在本文的讨论范畴,这里不在赘述。

三、参数刷新加载

众所周知,咱们用Redis无非就是将系统中不怎么变的、查询又比较频繁的数据缓存起来,例如咱们系统首页的轮播图啊,页面的动态连接啊,一些系统参数啊,公共数据啊都加载到Redis,而后有个后台管理系统去配置修改这些数据。

打个比方咱们首页的轮播图要再增长一个图,那咱们就在后管系统加上,加上就完事了吗?固然没有,由于Redis里仍是老数据。那你会说不是有过时时间吗?是的,但有的过时时间设置的较长如24小时而且咱们想当即生效怎么办?这时候咱们就能够利用Redis的发布订阅机制来实现数据的实时刷新。当咱们修改完数据后,点击刷新按钮,经过发布订阅机制,订阅者接收到消息后调用从新加载的方法便可。

代码实现

发布订阅的理论以及使用场景你们都已经有了大体了解了,可是怎么用代码实现发布订阅呢?在这里给你们分享一下实现方式。

咱们以第三种使用场景为例,先来看一下总体实现类图吧。

解释一下,这里咱们首先定义一个统一接口ICacheUpdate,只有一个update方法,咱们令Service层实现这个方法,执行具体的更新操做。咱们再来看RedisMsgPubSub,它继承redis.clients.jedis.JedisPubSub,主要重写其onMessage()方法(订阅的频道有消息到来时会触发这个方法),咱们在这个方法里调用RedisMsgPubSubupdate方法执行更新操做。当咱们有多个Service实现ICacheUpdate时,咱们就很是迫切地须要一个管理器来集中管理这些Service,而且当触发onMessage方法时要告诉onMessage方法具体调用哪一个ICacheUpdate的实现类,因此咱们有了PubSubManager。而且咱们单独开启一个线程来维护发布订阅,因此管理器继承了Thread类。

具体代码:

统一接口

ICacheUpdate.java

public interface ICacheUpdate {
    public void update();
}
复制代码

Service层

实现ICacheUpdate的update方法,执行具体的更新操做

InfoService.java

public class InfoService implements ICacheUpdate {
	private static Logger logger = LoggerFactory.getLogger(InfoService.class);
	@Autowired
	private RedisCache redisCache;
	@Autowired
	private InfoMapper infoMapper;
	/** * 按信息类型分类查询信息 * @return */
	public Map<String, List<Map<String, Object>>> selectAllInfo(){
		Map<String, List<Map<String, Object>>> resultMap = new HashMap<String, List<Map<String, Object>>>();
		List<String> infoTypeList = infoMapper.selectInfoType();//信息表中全部涉及的信息类型
		logger.info("-------按信息类型查找公共信息开始----"+infoTypeList);
		if(infoTypeList!=null && infoTypeList.size()>0) {
			for (String infoType : infoTypeList) {
				List<Map<String, Object>> result = infoMapper.selectByInfoType(infoType);
				resultMap.put(infoType, result);
			}
		}
		return resultMap;
	}
	@Override
	public void update() {
		//缓存首页信息
		logger.info("InfoService selectAllInfo 刷新缓存");
		Map<String, List<Map<String, Object>>> resultMap = this.selectAllInfo();
		Set<String> keySet = resultMap.keySet();
		for(String key:keySet){
			List<Map<String, Object>> value = resultMap.get(key);
			redisCache.putObject(GlobalSt.PUBLIC_INFO_ALL+key, value);
		}
	}
}
复制代码

Redis发布订阅的扩展类

做用:

一、统一管理ICacheUpdate,把全部实现ICacheUpdate接口的类添加到updates容器

二、重写onMessage方法,订阅到消息后进行刷新缓存的操做

RedisMsgPubSub.java

/** * Redis发布订阅的扩展类 * 做用:一、统一管理ICacheUpdate,把全部实现ICacheUpdate接口的类添加到updates容器 * 二、重写onMessage方法,订阅到消息后进行刷新缓存的操做 */
public class RedisMsgPubSub extends JedisPubSub {
    private static Logger logger = LoggerFactory.getLogger(RedisMsgPubSub.class);
    private Map<String , ICacheUpdate> updates = new HashMap<String , ICacheUpdate>();
    //一、由updates统一管理ICacheUpdate
    public boolean addListener(String key , ICacheUpdate update) {
        if(update == null) 
            return false;
	updates.put(key, update);
	return true;
    }
    /** * 二、重写onMessage方法,订阅到消息后进行刷新缓存的操做 * 订阅频道收到的消息 */
    @Override  
    public void onMessage(String channel, String message) {
        logger.info("RedisMsgPubSub onMessage channel:{},message :{}" ,channel, message);
        ICacheUpdate updater = null;
        if(StringUtil.isNotEmpty(message)) 
            updater = updates.get(message);
        if(updater!=null)
            updater.update();
    }
    //other code...
}
复制代码

发布订阅的管理器

执行的操做:

一、将全部须要刷新加载的Service类(实现ICacheUpdate接口)添加到RedisMsgPubSub的updates中

二、启动线程订阅pubsub_config频道,收到消息后的五秒后再次订阅(避免订阅到一次消息后结束订阅)

PubSubManager.java

public class PubSubManager extends Thread{
    private static Logger logger = LoggerFactory.getLogger(PubSubManager.class);

    public static Jedis jedis;
    RedisMsgPubSub msgPubSub = new RedisMsgPubSub();
    //频道
    public static final String PUNSUB_CONFIG = "pubsub_config";
    //1.将全部须要刷新加载的Service类(实现ICacheUpdate接口)添加到RedisMsgPubSub的updates中
    public boolean addListener(String key, ICacheUpdate listener){
        return msgPubSub.addListener(key,listener);
    }
    @Override
    public void run(){
        while (true){
            try {
                JedisPool jedisPool = SpringTools.getBean("jedisPool", JedisPool.class);
                if(jedisPool!=null){
                    jedis = jedisPool.getResource();
                    if(jedis!=null){
                        //2.启动线程订阅pubsub_config频道 阻塞
                        jedis.subscribe(msgPubSub,PUNSUB_CONFIG);
                    }
                }
            } catch (Exception e) {
                logger.error("redis connect error!");
            } finally {
                if(jedis!=null)
                    jedis.close();
            }
            try {
                //3.收到消息后的五秒后再次订阅(避免订阅到一次消息后结束订阅)
                Thread.sleep(5000);
            } catch (InterruptedException e) {
                logger.error("InterruptedException in redis sleep!");
            }
        }
    }
}
复制代码

到此,Redis的发布订阅大体已经实现。咱们何时启用呢?咱们能够选择在启动项目时完成订阅和基础数据的加载,因此咱们经过实现javax.servlet.SevletContextListener来完成这一操做。而后将监听器添加到web.xml

CacheInitListener.java

/** * 加载系统参数 */
public class CacheInitListener implements ServletContextListener{
    private static Logger logger = LoggerFactory.getLogger(CacheInitListener.class);

    @Override
    public void contextDestroyed(ServletContextEvent arg0) {
    }

    @Override
    public void contextInitialized(ServletContextEvent arg0) {
        logger.info("---CacheListener初始化开始---");
        init();
        logger.info("---CacheListener初始化结束---");
    }

    public void init() {
        try {
            //得到管理器
            PubSubManager pubSubManager = SpringTools.getBean("pubSubManager", PubSubManager.class);

            InfoService infoService = SpringTools.getBean("infoService", InfoService.class);
            //添加到管理器
            pubSubManager.addListener("infoService", infoService);
            //other service...

            //启动线程执行订阅操做
            pubSubManager.start();
            //初始化加载
            loadParamToRedis();
        } catch (Exception e) {
            logger.info(e.getMessage(), e);
        }
    }

    private void loadParamToRedis() {
        InfoService infoService = SpringTools.getBean("infoService", InfoService.class);
        infoService.update();
        //other service...
    }
}
复制代码

web.xml

<listener>
	<listener-class>com.xxx.listener.CacheInitListener</listener-class>
</listener>
复制代码

【end】

文章首发于公众号@编程大道

相关文章
相关标签/搜索