zoukankan      html  css  js  c++  java
  • Redis实现分布式锁与任务队列

    大家都知道在天猫、京东、苏宁等等电商网站上有很多秒杀活动,例如在某一个时刻抢购一个原价1999现在秒杀价只要999的手机时,会迎来一个用户请求的高峰期,会有几十万几百万的并发量,来抢这个手机,在高并发的情形下会对数据库服务器或者是文件服务器应用服务器造成巨大的压力,严重时说不定就宕机了。

    另一个问题是,秒杀的东西都是有量的,例如一款手机只有10台的量秒杀,那么,在高并发的情况下,成千上万条数据更新数据库(例如10台的量被人抢一台就会在数据集某些记录下 减1),那次这个时候的先后顺序是很乱的,很容易出现10台的量,抢到的人就不止10个这种严重的问题。那么,以后所说的问题我们该如何去解决呢? 接下来我所分享的技术就可以拿来处理以上的问题: 分布式锁 和 任务队列。

    1.redis分布式锁

    1)为避免特殊原因导致锁无法释放,在加锁成功后,锁会被赋予一个生存时间(通过lock方法的参数设置或者使用默认值),超出生存时间锁会被自动释放锁的生存时间默认比较短(秒级),因此,若需要长时间加锁,可以通过expire方法延长锁的生存时间为适当时间,比如在循环内。

    2)系统级的锁当进程无论何种原因时出现crash时,操作系统会自己回收锁,所以不会出现资源丢失,但分布式锁不用,若一次性设置很长时间,一旦由于各种原因出现进程crash 或者其他异常导致unlock未被调用时,则该锁在剩下的时间就会变成垃圾锁,导致其他进程或者进程重启后无法进入加锁区域。

       /**
         * 加锁
         * @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对应的id
        foreach ($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对应的id
        foreach ($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脚本,让服务器定时运行某些程序,实现入队出队等操作,这里我就不在将其与实际应用结合起来去实现了,大家理解好这两大功能的实现思路即可。

  • 相关阅读:
    我的屌丝giser成长记-研二篇
    我的屌丝giser成长记-研一篇(下)
    C#连接Oracle数据库的方法(Oracle.DataAccess.Client也叫ODP.net)
    C# 日期格式化的中的(/)正斜杠的问题(与操作系统设置有关)
    C#,SOAP1.1与1.2的发布与禁用(SOAP 1.2 in .NET Framework 2.0)
    C#使用WebService 常见问题处理
    sql查询数据库中所有表的记录条数,以及占用磁盘空间大小。
    eclipse中的XML文件无法快捷键注释问题
    对比两个表中,字段名不一样的SQL
    oracle 恢复备份
  • 原文地址:https://www.cnblogs.com/jackzhuo/p/13652621.html
Copyright © 2011-2022 走看看