用Redis实现分布式锁 与 实现任务队列【转载】
这一次总结和分享用Redis实现分布式锁 与 实现任务队列 这两大强大的功能。先扯点个人观点,之前我看了一篇博文说博客园的文章大部分都是分享代码,博文里强调说分享思路比分享代码更重要(貌似大概是这个意思,若有误请谅解),但我觉得,分享思路固然重要,但有了思路,却没有实现的代码,那会让人觉得很浮夸的,在工作中的程序猿都知道,你去实现一个功能模块,一段代码,虽然你有了思路,但是实现的过程也是很耗时的,特别是代码调试,还有各种测试等等。所以我认为,思路+代码,才是一篇好博文的主要核心。
直接进入主题。
一、前言
双十一刚过不久,大家都知道在天猫、京东、苏宁等等电商网站上有很多秒杀活动,例如在某一个时刻抢购一个原价1999现在秒杀价只要999的手机时,会迎来一个用户请求的高峰期,可能会有几十万几百万的并发量,来抢这个手机,在高并发的情形下会对数据库服务器或者是文件服务器应用服务器造成巨大的压力,严重时说不定就宕机了,另一个问题是,秒杀的东西都是有量的,例如一款手机只有10台的量秒杀,那么,在高并发的情况下,成千上万条数据更新数据库(例如10台的量被人抢一台就会在数据集某些记录下 减1),那次这个时候的先后顺序是很乱的,很容易出现10台的量,抢到的人就不止10个这种严重的问题。那么,以后所说的问题我们该如何去解决呢? 接下来我所分享的技术就可以拿来处理以上的问题: 分布式锁 和 任务队列。
二、实现思路
1.Redis实现分布式锁思路
思路很简单,主要用到的redis函数是setnx(),这个应该是实现分布式锁最主要的函数。首先是将某一任务标识名(这里用Lock:order作为标识名的例子)作为键存到redis里,并为其设个过期时间,如果是还有Lock:order请求过来,先是通过setnx()看看是否能将Lock:order插入到redis里,可以的话就返回true,不可以就返回false。当然,在我的代码里会比这个思路复杂一些,我会在分析代码时进一步说明。
2.Redis实现任务队列
这里的实现会用到上面的Redis分布式的锁机制,主要是用到了Redis里的有序集合这一数据结构。例如入队时,通过zset的add()函数进行入队,而出对时,可以用到zset的getScore()函数。另外还可以弹出顶部的几个任务。
以上就是实现 分布式锁 和 任务队列 的简单思路,如果你看完有点模棱两可,那请看接下来的代码实现。
三、代码分析
(一)先来分析Redis分布式锁的代码实现
(1)为避免特殊原因导致锁无法释放,在加锁成功后,锁会被赋予一个生存时间(通过lock方法的参数设置或者使用默认值),超出生存时间锁会被自动释放锁的生存时间默认比较短(秒级),因此,若需要长时间加锁,可以通过expire方法延长锁的生存时间为适当时间,比如在循环内。
(2)系统级的锁当进程无论何种原因时出现crash时,操作系统会自己回收锁,所以不会出现资源丢失,但分布式锁不用,若一次性设置很长时间,一旦由于各种原因出现进程crash 或者其他异常导致unlock未被调用时,则该锁在剩下的时间就会变成垃圾锁,导致其他进程或者进程重启后无法进入加锁区域。
先看加锁的实现代码:这里需要主要两个参数,一个是$timeout,这个是循环获取锁的等待时间,在这个时间内会一直尝试获取锁知道超时,如果为0,则表示获取锁失败后直接返回而不再等待;另一个重要参数的$expire,这个参数指当前锁的最大生存时间,以秒为单位的,它必须大于0,如果超过生存时间锁仍未被释放,则系统会自动强制释放。这个参数的最要作用请看上面的(1)里的解释。
这里先取得当前时间,然后再获取到锁失败时的等待超时的时刻(是个时间戳),再获取到锁的最大生存时刻是多少。这里redis的key用这种格式:"Lock:锁的标识名",这里就开始进入循环了,先是插入数据到redis里,使用setnx()函数,这函数的意思是,如果该键不存在则插入数据,将最大生存时刻作为值存储,假如插入成功,则对该键进行失效时间的设置,并将该键放在$lockedName数组里,返回true,也就是上锁成功;如果该键存在,则不会插入操作了,这里有一步严谨的操作,那就是取得当前键的剩余时间,假如这个时间小于0,表示key上没有设置生存时间(key是不会不存在的,因为前面setnx会自动创建)如果出现这种状况,那就是进程的某个实例setnx成功后 crash 导致紧跟着的expire没有被调用,这时可以直接设置expire并把锁纳为己用。如果没设置锁失败的等待时间 或者 已超过最大等待时间了,那就退出循环,反之则 隔 $waitIntervalUs 后继续 请求。 这就是加锁的整一个代码分析。
/*** 加锁* @param [type] $name 锁的标识名* @param integer $timeout 循环获取锁的等待超时时间,在此时间内会一直尝试获取锁直到超时,为0表示失败后直接返回不等待* @param integer $expire 当前锁的最大生存时间(秒),必须大于0,如果超过生存时间锁仍未被释放,则系统会自动强制释放* @param integer $waitIntervalUs 获取锁失败后挂起再试的时间间隔(微秒)* @return [type] [description]*/public function lock($name, $timeout = 0, $expire = 15, $waitIntervalUs = 100000) {if ($name == null) return false;//取得当前时间$now = time();//获取锁失败时的等待超时时刻$timeoutAt = $now + $timeout;//锁的最大生存时刻$expireAt = $now + $expire;$redisKey = "Lock:{$name}";while (true) {//将rediskey的最大生存时刻存到redis里,过了这个时刻该锁会被自动释放$result = $this->redisString->setnx($redisKey, $expireAt);if ($result != false) {//设置key的失效时间$this->redisString->expire($redisKey, $expireAt);//将锁标志放到lockedNames数组里$this->lockedNames[$name] = $expireAt;return true;}//以秒为单位,返回给定key的剩余生存时间$ttl = $this->redisString->ttl($redisKey);//ttl小于0 表示key上没有设置生存时间(key是不会不存在的,因为前面setnx会自动创建)//如果出现这种状况,那就是进程的某个实例setnx成功后 crash 导致紧跟着的expire没有被调用//这时可以直接设置expire并把锁纳为己用if ($ttl < 0) {$this->redisString->set($redisKey, $expireAt);$this->lockedNames[$name] = $expireAt;return true;}/*****循环请求锁部分*****///如果没设置锁失败的等待时间 或者 已超过最大等待时间了,那就退出if ($timeout <= 0 || $timeoutAt < microtime(true)) break;//隔 $waitIntervalUs 后继续 请求usleep($waitIntervalUs);}return false;}
接着看解锁的代码分析:解锁就简单多了,传入参数就是锁标识,先是判断是否存在该锁,存在的话,就从redis里面通过deleteKey()函数删除掉锁标识即可。
/*** 解锁* @param [type] $name [description]* @return [type] [description]*/public function unlock($name) {//先判断是否存在此锁if ($this->isLocking($name)) {//删除锁if ($this->redisString->deleteKey("Lock:$name")) {//清掉lockedNames里的锁标志unset($this->lockedNames[$name]);return true;}}return false;}
在贴上删除掉所有锁的方法,其实都一个样,多了个循环遍历而已。
/*** 释放当前所有获得的锁* @return [type] [description]*/public function unlockAll() {//此标志是用来标志是否释放所有锁成功$allSuccess = true;foreach ($this->lockedNames as $name => $expireAt) {if (false === $this->unlock($name)) {$allSuccess = false; }}return $allSuccess;}
以上就是用Redis实现分布式锁的整一套思路和代码实现的总结和分享,这里我附上正一个实现类的代码,代码里我基本上对每一行进行了注释,方便大家快速看懂并且能模拟应用。想要深入了解的请看整个类的代码:
/***在redis上实现分布式锁*/ class RedisLock {private $redisString;private $lockedNames = [];public function __construct($param = NULL) {$this->redisString = RedisFactory::get($param)->string;}/*** 加锁* @param [type] $name 锁的标识名* @param integer $timeout 循环获取锁的等待超时时间,在此时间内会一直尝试获取锁直到超时,为0表示失败后直接返回不等待* @param integer $expire 当前锁的最大生存时间(秒),必须大于0,如果超过生存时间锁仍未被释放,则系统会自动强制释放* @param integer $waitIntervalUs 获取锁失败后挂起再试的时间间隔(微秒)* @return [type] [description]*/public function lock($name, $timeout = 0, $expire = 15, $waitIntervalUs = 100000) {if ($name == null) return false;//取得当前时间$now = time();//获取锁失败时的等待超时时刻$timeoutAt = $now + $timeout;//锁的最大生存时刻$expireAt = $now + $expire;$redisKey = "Lock:{$name}";while (true) {//将rediskey的最大生存时刻存到redis里,过了这个时刻该锁会被自动释放$result = $this->redisString->setnx($redisKey, $expireAt);if ($result != false) {//设置key的失效时间$this->redisString->expire($redisKey, $expireAt);//将锁标志放到lockedNames数组里$this->lockedNames[$name] = $expireAt;return true;}//以秒为单位,返回给定key的剩余生存时间$ttl = $this->redisString->ttl($redisKey);//ttl小于0 表示key上没有设置生存时间(key是不会不存在的,因为前面setnx会自动创建)//如果出现这种状况,那就是进程的某个实例setnx成功后 crash 导致紧跟着的expire没有被调用//这时可以直接设置expire并把锁纳为己用if ($ttl < 0) {$this->redisString->set($redisKey, $expireAt);$this->lockedNames[$name] = $expireAt;return true;}/*****循环请求锁部分*****///如果没设置锁失败的等待时间 或者 已超过最大等待时间了,那就退出if ($timeout <= 0 || $timeoutAt < microtime(true)) break;//隔 $waitIntervalUs 后继续 请求usleep($waitIntervalUs);}return false;}/*** 解锁* @param [type] $name [description]* @return [type] [description]*/public function unlock($name) {//先判断是否存在此锁if ($this->isLocking($name)) {//删除锁if ($this->redisString->deleteKey("Lock:$name")) {//清掉lockedNames里的锁标志unset($this->lockedNames[$name]);return true;}}return false;}/*** 释放当前所有获得的锁* @return [type] [description]*/public function unlockAll() {//此标志是用来标志是否释放所有锁成功$allSuccess = true;foreach ($this->lockedNames as $name => $expireAt) {if (false === $this->unlock($name)) {$allSuccess = false; }}return $allSuccess;}/*** 给当前所增加指定生存时间,必须大于0* @param [type] $name [description]* @return [type] [description]*/public function expire($name, $expire) {//先判断是否存在该锁if ($this->isLocking($name)) {//所指定的生存时间必须大于0$expire = max($expire, 1);//增加锁生存时间if ($this->redisString->expire("Lock:$name", $expire)) {return true;}}return false;}/*** 判断当前是否拥有指定名字的所* @param [type] $name [description]* @return boolean [description]*/public function isLocking($name) {//先看lonkedName[$name]是否存在该锁标志名if (isset($this->lockedNames[$name])) {//从redis返回该锁的生存时间return (string)$this->lockedNames[$name] = (string)$this->redisString->get("Lock:$name");}return false;}}Redis实现分布式锁
(二)用Redis实现任务队列的代码分析
(1)任务队列,用于将业务逻辑中可以异步处理的操作放入队列中,在其他线程中处理后出队
(2)队列中使用了分布式锁和其他逻辑,保证入队和出队的一致性
(3)这个队列和普通队列不一样,入队时的id是用来区分重复入队的,队列里面只会有一条记录,同一个id后入的覆盖前入的,而不是追加, 如果需求要求重复入队当做不用的任务,请使用不同的id区分
先看入队的代码分析:首先当然是对参数的合法性检测,接着就用到上面加锁机制的内容了,就是开始加锁,入队时我这里选择当前时间戳作为score,接着就是入队了,使用的是zset数据结构的add()方法,入队完成后,就对该任务解锁,即完成了一个入队的操作。
/*** 入队一个 Task* @param [type] $name 队列名称* @param [type] $id 任务id(或者其数组)* @param integer $timeout 入队超时时间(秒)* @param integer $afterInterval [description]* @return [type] [description]*/public function enqueue($name, $id, $timeout = 10, $afterInterval = 0) {//合法性检测if (empty($name) || empty($id) || $timeout <= 0) return false;//加锁if (!$this->_redis->lock->lock("Queue:{$name}", $timeout)) {Logger::get('queue')->error("enqueue faild becouse of lock failure: name = $name, id = $id");return false;}//入队时以当前时间戳作为 score$score = microtime(true) + $afterInterval;//入队foreach ((array)$id as $item) {//先判断下是否已经存在该id了if (false === $this->_redis->zset->getScore("Queue:$name", $item)) {$this->_redis->zset->add("Queue:$name", $score, $item);}}//解锁$this->_redis->lock->unlock("Queue:$name");return true;}
接着来看一下出队的代码分析:出队一个Task,需要指定它的$id 和 $score,如果$score与队列中的匹配则出队,否则认为该Task已被重新入队过,当前操作按失败处理。首先和对参数进行合法性检测,接着又用到加锁的功能了,然后及时出队了,先使用getScore()从Redis里获取到该id的score,然后将传入的$score和Redis里存储的score进行对比,如果两者相等就进行出队操作,也就是使用zset里的delete()方法删掉该任务id,最后当前就是解锁了。这就是出队的代码分析。
/*** 出队一个Task,需要指定$id 和 $score* 如果$score 与队列中的匹配则出队,否则认为该Task已被重新入队过,当前操作按失败处理* * @param [type] $name 队列名称 * @param [type] $id 任务标识* @param [type] $score 任务对应score,从队列中获取任务时会返回一个score,只有$score和队列中的值匹配时Task才会被出队* @param integer $timeout 超时时间(秒)* @return [type] Task是否成功,返回false可能是redis操作失败,也有可能是$score与队列中的值不匹配(这表示该Task自从获取到本地之后被其他线程入队过)*/public function dequeue($name, $id, $score, $timeout = 10) {//合法性检测if (empty($name) || empty($id) || empty($score)) return false;//加锁if (!$this->_redis->lock->lock("Queue:$name", $timeout)) {Logger:get('queue')->error("dequeue faild becouse of lock lailure:name=$name, id = $id");return false;}//出队//先取出redis的score$serverScore = $this->_redis->zset->getScore("Queue:$name", $id);$result = false;//先判断传进来的score和redis的score是否是一样if ($serverScore == $score) {//删掉该$id$result = (float)$this->_redis->zset->delete("Queue:$name", $id);if ($result == false) {Logger::get('queue')->error("dequeue faild because of redis delete failure: name =$name, id = $id");}}//解锁$this->_redis->lock->unlock("Queue:$name");return $result;}
学过数据结构这门课的朋友都应该知道,队列操作还有弹出顶部某个值的方法等等,这里处理入队出队操作,我还实现了 获取队列顶部若干个Task 并将其出队的方法,想了解的朋友可以看这段代码,假如看不太明白就留言,这里我不再对其进行分析了。
/*** 获取队列顶部若干个Task 并将其出队* @param [type] $name 队列名称* @param integer $count 数量* @param integer $timeout 超时时间* @return [type] 返回数组[0=>['id'=> , 'score'=> ], 1=>['id'=> , 'score'=> ], 2=>['id'=> , 'score'=> ]]*/public function pop($name, $count = 1, $timeout = 10) {//合法性检测if (empty($name) || $count <= 0) return []; //加锁if (!$this->_redis->lock->lock("Queue:$name")) {Log::get('queue')->error("pop faild because of pop failure: name = $name, count = $count");return false;}//取出若干的Task$result = [];$array = $this->_redis->zset->getByScore("Queue:$name", false, microtime(true), true, false, [0, $count]);//将其放在$result数组里 并 删除掉redis对应的idforeach ($array as $id => $score) {$result[] = ['id'=>$id, 'score'=>$score];$this->_redis->zset->delete("Queue:$name", $id);}//解锁$this->_redis->lock->unlock("Queue:$name");return $count == 1 ? (empty($result) ? false : $result[0]) : $result;}
以上就是用Redis实现任务队列的整一套思路和代码实现的总结和分享,这里我附上正一个实现类的代码,代码里我基本上对每一行进行了注释,方便大家快速看懂并且能模拟应用。想要深入了解的请看整个类的代码:
/*** 任务队列* */ class RedisQueue {private $_redis;public function __construct($param = null) {$this->_redis = RedisFactory::get($param);}/*** 入队一个 Task* @param [type] $name 队列名称* @param [type] $id 任务id(或者其数组)* @param integer $timeout 入队超时时间(秒)* @param integer $afterInterval [description]* @return [type] [description]*/public function enqueue($name, $id, $timeout = 10, $afterInterval = 0) {//合法性检测if (empty($name) || empty($id) || $timeout <= 0) return false;//加锁if (!$this->_redis->lock->lock("Queue:{$name}", $timeout)) {Logger::get('queue')->error("enqueue faild becouse of lock failure: name = $name, id = $id");return false;}//入队时以当前时间戳作为 score$score = microtime(true) + $afterInterval;//入队foreach ((array)$id as $item) {//先判断下是否已经存在该id了if (false === $this->_redis->zset->getScore("Queue:$name", $item)) {$this->_redis->zset->add("Queue:$name", $score, $item);}}//解锁$this->_redis->lock->unlock("Queue:$name");return true;}/*** 出队一个Task,需要指定$id 和 $score* 如果$score 与队列中的匹配则出队,否则认为该Task已被重新入队过,当前操作按失败处理* * @param [type] $name 队列名称 * @param [type] $id 任务标识* @param [type] $score 任务对应score,从队列中获取任务时会返回一个score,只有$score和队列中的值匹配时Task才会被出队* @param integer $timeout 超时时间(秒)* @return [type] Task是否成功,返回false可能是redis操作失败,也有可能是$score与队列中的值不匹配(这表示该Task自从获取到本地之后被其他线程入队过)*/public function dequeue($name, $id, $score, $timeout = 10) {//合法性检测if (empty($name) || empty($id) || empty($score)) return false;//加锁if (!$this->_redis->lock->lock("Queue:$name", $timeout)) {Logger:get('queue')->error("dequeue faild becouse of lock lailure:name=$name, id = $id");return false;}//出队//先取出redis的score$serverScore = $this->_redis->zset->getScore("Queue:$name", $id);$result = false;//先判断传进来的score和redis的score是否是一样if ($serverScore == $score) {//删掉该$id$result = (float)$this->_redis->zset->delete("Queue:$name", $id);if ($result == false) {Logger::get('queue')->error("dequeue faild because of redis delete failure: name =$name, id = $id");}}//解锁$this->_redis->lock->unlock("Queue:$name");return $result;}/*** 获取队列顶部若干个Task 并将其出队* @param [type] $name 队列名称* @param integer $count 数量* @param integer $timeout 超时时间* @return [type] 返回数组[0=>['id'=> , 'score'=> ], 1=>['id'=> , 'score'=> ], 2=>['id'=> , 'score'=> ]]*/public function pop($name, $count = 1, $timeout = 10) {//合法性检测if (empty($name) || $count <= 0) return []; //加锁if (!$this->_redis->lock->lock("Queue:$name")) {Logger::get('queue')->error("pop faild because of pop failure: name = $name, count = $count");return false;}//取出若干的Task$result = [];$array = $this->_redis->zset->getByScore("Queue:$name", false, microtime(true), true, false, [0, $count]);//将其放在$result数组里 并 删除掉redis对应的idforeach ($array as $id => $score) {$result[] = ['id'=>$id, 'score'=>$score];$this->_redis->zset->delete("Queue:$name", $id);}//解锁$this->_redis->lock->unlock("Queue:$name");return $count == 1 ? (empty($result) ? false : $result[0]) : $result;}/*** 获取队列顶部的若干个Task* @param [type] $name 队列名称* @param integer $count 数量* @return [type] 返回数组[0=>['id'=> , 'score'=> ], 1=>['id'=> , 'score'=> ], 2=>['id'=> , 'score'=> ]]*/public function top($name, $count = 1) {//合法性检测if (empty($name) || $count < 1) return [];//取错若干个Task$result = [];$array = $this->_redis->zset->getByScore("Queue:$name", false, microtime(true), true, false, [0, $count]);//将Task存放在数组里foreach ($array as $id => $score) {$result[] = ['id'=>$id, 'score'=>$score];}//返回数组 return $count == 1 ? (empty($result) ? false : $result[0]) : $result; } }Redis实现任务队列
到此,这两大块功能基本讲解完毕,对于任务队列,你可以写一个shell脚本,让服务器定时运行某些程序,实现入队出队等操作,这里我就不在将其与实际应用结合起来去实现了,大家理解好这两大功能的实现思路即可,由于代码用的是PHP语言来写的,如果你理解了实现思路,你完全可以使用java或者是.net等等其他语言去实现这两个功能。这两大功能的应用场景十分多,特别是秒杀,另一个就是春运抢火车票,这两个是最鲜明的例子了。当然还有很多地方用到,这里我不再一一列举。
好了,本次总结和分享到此完毕。最后我附上 分布式锁和任务队列这两个类:
/***在redis上实现分布式锁*/ class RedisLock {private $redisString;private $lockedNames = [];public function __construct($param = NULL) {$this->redisString = RedisFactory::get($param)->string;}/*** 加锁* @param [type] $name 锁的标识名* @param integer $timeout 循环获取锁的等待超时时间,在此时间内会一直尝试获取锁直到超时,为0表示失败后直接返回不等待* @param integer $expire 当前锁的最大生存时间(秒),必须大于0,如果超过生存时间锁仍未被释放,则系统会自动强制释放* @param integer $waitIntervalUs 获取锁失败后挂起再试的时间间隔(微秒)* @return [type] [description]*/public function lock($name, $timeout = 0, $expire = 15, $waitIntervalUs = 100000) {if ($name == null) return false;//取得当前时间$now = time();//获取锁失败时的等待超时时刻$timeoutAt = $now + $timeout;//锁的最大生存时刻$expireAt = $now + $expire;$redisKey = "Lock:{$name}";while (true) {//将rediskey的最大生存时刻存到redis里,过了这个时刻该锁会被自动释放$result = $this->redisString->setnx($redisKey, $expireAt);if ($result != false) {//设置key的失效时间$this->redisString->expire($redisKey, $expireAt);//将锁标志放到lockedNames数组里$this->lockedNames[$name] = $expireAt;return true;}//以秒为单位,返回给定key的剩余生存时间$ttl = $this->redisString->ttl($redisKey);//ttl小于0 表示key上没有设置生存时间(key是不会不存在的,因为前面setnx会自动创建)//如果出现这种状况,那就是进程的某个实例setnx成功后 crash 导致紧跟着的expire没有被调用//这时可以直接设置expire并把锁纳为己用if ($ttl < 0) {$this->redisString->set($redisKey, $expireAt);$this->lockedNames[$name] = $expireAt;return true;}/*****循环请求锁部分*****///如果没设置锁失败的等待时间 或者 已超过最大等待时间了,那就退出if ($timeout <= 0 || $timeoutAt < microtime(true)) break;//隔 $waitIntervalUs 后继续 请求usleep($waitIntervalUs);}return false;}/*** 解锁* @param [type] $name [description]* @return [type] [description]*/public function unlock($name) {//先判断是否存在此锁if ($this->isLocking($name)) {//删除锁if ($this->redisString->deleteKey("Lock:$name")) {//清掉lockedNames里的锁标志unset($this->lockedNames[$name]);return true;}}return false;}/*** 释放当前所有获得的锁* @return [type] [description]*/public function unlockAll() {//此标志是用来标志是否释放所有锁成功$allSuccess = true;foreach ($this->lockedNames as $name => $expireAt) {if (false === $this->unlock($name)) {$allSuccess = false; }}return $allSuccess;}/*** 给当前所增加指定生存时间,必须大于0* @param [type] $name [description]* @return [type] [description]*/public function expire($name, $expire) {//先判断是否存在该锁if ($this->isLocking($name)) {//所指定的生存时间必须大于0$expire = max($expire, 1);//增加锁生存时间if ($this->redisString->expire("Lock:$name", $expire)) {return true;}}return false;}/*** 判断当前是否拥有指定名字的所* @param [type] $name [description]* @return boolean [description]*/public function isLocking($name) {//先看lonkedName[$name]是否存在该锁标志名if (isset($this->lockedNames[$name])) {//从redis返回该锁的生存时间return (string)$this->lockedNames[$name] = (string)$this->redisString->get("Lock:$name");}return false;}}/*** 任务队列*/ class RedisQueue {private $_redis;public function __construct($param = null) {$this->_redis = RedisFactory::get($param);}/*** 入队一个 Task* @param [type] $name 队列名称* @param [type] $id 任务id(或者其数组)* @param integer $timeout 入队超时时间(秒)* @param integer $afterInterval [description]* @return [type] [description]*/public function enqueue($name, $id, $timeout = 10, $afterInterval = 0) {//合法性检测if (empty($name) || empty($id) || $timeout <= 0) return false;//加锁if (!$this->_redis->lock->lock("Queue:{$name}", $timeout)) {Logger::get('queue')->error("enqueue faild becouse of lock failure: name = $name, id = $id");return false;}//入队时以当前时间戳作为 score$score = microtime(true) + $afterInterval;//入队foreach ((array)$id as $item) {//先判断下是否已经存在该id了if (false === $this->_redis->zset->getScore("Queue:$name", $item)) {$this->_redis->zset->add("Queue:$name", $score, $item);}}//解锁$this->_redis->lock->unlock("Queue:$name");return true;}/*** 出队一个Task,需要指定$id 和 $score* 如果$score 与队列中的匹配则出队,否则认为该Task已被重新入队过,当前操作按失败处理* * @param [type] $name 队列名称 * @param [type] $id 任务标识* @param [type] $score 任务对应score,从队列中获取任务时会返回一个score,只有$score和队列中的值匹配时Task才会被出队* @param integer $timeout 超时时间(秒)* @return [type] Task是否成功,返回false可能是redis操作失败,也有可能是$score与队列中的值不匹配(这表示该Task自从获取到本地之后被其他线程入队过)*/public function dequeue($name, $id, $score, $timeout = 10) {//合法性检测if (empty($name) || empty($id) || empty($score)) return false;//加锁if (!$this->_redis->lock->lock("Queue:$name", $timeout)) {Logger:get('queue')->error("dequeue faild becouse of lock lailure:name=$name, id = $id");return false;}//出队//先取出redis的score$serverScore = $this->_redis->zset->getScore("Queue:$name", $id);$result = false;//先判断传进来的score和redis的score是否是一样if ($serverScore == $score) {//删掉该$id$result = (float)$this->_redis->zset->delete("Queue:$name", $id);if ($result == false) {Logger::get('queue')->error("dequeue faild because of redis delete failure: name =$name, id = $id");}}//解锁$this->_redis->lock->unlock("Queue:$name");return $result;}/*** 获取队列顶部若干个Task 并将其出队* @param [type] $name 队列名称* @param integer $count 数量* @param integer $timeout 超时时间* @return [type] 返回数组[0=>['id'=> , 'score'=> ], 1=>['id'=> , 'score'=> ], 2=>['id'=> , 'score'=> ]]*/public function pop($name, $count = 1, $timeout = 10) {//合法性检测if (empty($name) || $count <= 0) return []; //加锁if (!$this->_redis->lock->lock("Queue:$name")) {Logger::get('queue')->error("pop faild because of pop failure: name = $name, count = $count");return false;}//取出若干的Task$result = [];$array = $this->_redis->zset->getByScore("Queue:$name", false, microtime(true), true, false, [0, $count]);//将其放在$result数组里 并 删除掉redis对应的idforeach ($array as $id => $score) {$result[] = ['id'=>$id, 'score'=>$score];$this->_redis->zset->delete("Queue:$name", $id);}//解锁$this->_redis->lock->unlock("Queue:$name");return $count == 1 ? (empty($result) ? false : $result[0]) : $result;}/*** 获取队列顶部的若干个Task* @param [type] $name 队列名称* @param integer $count 数量* @return [type] 返回数组[0=>['id'=> , 'score'=> ], 1=>['id'=> , 'score'=> ], 2=>['id'=> , 'score'=> ]]*/public function top($name, $count = 1) {//合法性检测if (empty($name) || $count < 1) return [];//取错若干个Task$result = [];$array = $this->_redis->zset->getByScore("Queue:$name", false, microtime(true), true, false, [0, $count]);//将Task存放在数组里foreach ($array as $id => $score) {$result[] = ['id'=>$id, 'score'=>$score];}//返回数组 return $count == 1 ? (empty($result) ? false : $result[0]) : $result; } }Redis分布式锁和任务队列代码
转载于:https://www.cnblogs.com/liliuguang/p/8807459.html
用Redis实现分布式锁 与 实现任务队列【转载】相关推荐
- 基于 Redis 实现分布式锁思考
以下文章来源方志朋的博客,回复"666"获面试宝典 来源:blog.csdn.net/xuan_lu/article/details/111600302 分布式锁 基于redis实 ...
- Redis实现分布式锁的深入探究
点击上方"方志朋",选择"设为星标" 回复"666"获取新整理的面试文章 一.分布式锁简介 锁 是一种用来解决多个执行线程 访问共享资源 错 ...
- nx set 怎么实现的原子性_基于Redis的分布式锁实现
前言 本篇文章主要介绍基于Redis的分布式锁实现到底是怎么一回事,其中参考了许多大佬写的文章,算是对分布式锁做一个总结 分布式锁概览 在多线程的环境下,为了保证一个代码块在同一时间只能由一个线程访问 ...
- Zookeeper和Redis实现分布式锁,附我的可靠性分析
作者:今天你敲代码了吗 链接:https://www.jianshu.com/p/b6953745e341 在分布式系统中,为保证同一时间只有一个客户端可以对共享资源进行操作,需要对共享资源加锁来实现 ...
- Redis——由分布式锁造成的重大事故
作者:浪漫先生 原文:juejin.im/post/6854573212831842311 前言 基于Redis使用分布式锁在当今已经不是什么新鲜事了.本篇文章主要是基于我们实际项目中因为redis分 ...
- 基于Redis的分布式锁和Redlock算法
来自:后端技术指南针 1 前言 今天开始来和大家一起学习一下Redis实际应用篇,会写几个Redis的常见应用. 在我看来Redis最为典型的应用就是作为分布式缓存系统,其他的一些应用本质上并不是杀手 ...
- 《Redis官方文档》用Redis构建分布式锁
<Redis官方文档>用Redis构建分布式锁 用Redis构建分布式锁 在不同进程需要互斥地访问共享资源时,分布式锁是一种非常有用的技术手段. 有很多三方库和文章描述如何用Redis实现 ...
- 《Redis官方文档》用Redis构建分布式锁(悲观锁)
2019独角兽企业重金招聘Python工程师标准>>> **用Redis构建分布式锁 ** 在不同进程需要互斥地访问共享资源时,分布式锁是一种非常有用的技术手段. 有很多三方库和文章 ...
- redis 实现分布式锁
为什么80%的码农都做不了架构师?>>> redis 实现分布式锁 伪代码 lock(){if(jedis.setNx("key",timestamp)){ ...
- redis系列:基于redis的分布式锁
一.介绍 这篇博文讲介绍如何一步步构建一个基于Redis的分布式锁.会从最原始的版本开始,然后根据问题进行调整,最后完成一个较为合理的分布式锁. 本篇文章会将分布式锁的实现分为两部分,一个是单机环境, ...
最新文章
- DOM中严格区分大小写
- C++中return语句的用法
- aop简介-基于jdk的动态代理
- JS浏览器加载一个页面的过程
- Java编程提高性能的26个方法
- 设置TextField内文字距左边框的距离
- 15 PP配置-生产计划-主数据-定义特殊采购类型
- vuedraggable嵌套块拖拽_Vue.Draggable拖拽效果
- lazyload 加载
- 运行python脚本时出现no module named cv2怎么解决
- QoS中流量监管和流量整形详解
- 使用MSDN学习ASP.NET的工作流程
- 微电子电路——一位全加器
- pass 软件_PASS软件非劣效Logrank检验的h1参数如何设置?
- 读书笔记3——《用户故事与敏捷方法》
- centos查询 硬盘序列号查询_关于使用java执行shell脚本获取centos的硬盘序列号和mac地址...
- LQ-1600K打印机色带传动故障分析
- STM32F205 PWM配置
- 公司IT管理制度——案例分享
- 基于片内Flash的提示音播放程序
热门文章
- Flink CDC 系列 - 同步 MySQL 分库分表,构建 Iceberg 实时数据湖
- Android自定义View【实战教程】5⃣️---Canvas详解及代码绘制安卓机器人
- matlab中求解非线性方程组的函数,利用solve函数求解非线性方程组的问题
- python中文分词与词云画像_使用Python绘制肖像词云
- python 开发公众号sdk_「公众号开发」基于Serverless架构Python实现公众号图文搜索...
- python3 django 中文乱码_python3 wsgi服务和响应数据中文乱码问题
- linux命令(47):Linux下对文件进行按行排序,去除重复行
- server.htaccess 具体解释以及 .htaccess 參数说明
- 使用at任务定点执行
- 团队作业三——项目思考