Реализация и принцип распределенной блокировки Redis

Java

1. Введение

В программе мы хотим обеспечить видимость и атомарность переменной, мы можем использовать volatile (чтение/запись любой отдельной volatile переменной является атомарным, но составные операции, подобные volatile++, не являются атомарными), synchronized, оптимистическая блокировка, пессимистическая блокировка и т.д. для контроля. Это можно сделать в одном приложении.Сейчас,с развитием времени,большинство проектов попрощались с эрой одной машины и вступили в эру микросервисов.В этом случае многие сервисы должны быть кластеризованы,и приложение должно быть развернуто на нескольких машинах.Для балансировки нагрузки невозможно использовать вышеупомянутый механизм для обеспечения видимости и атомарности переменных в параллельных условиях (как показано на рисунке ниже), что приводит к множеству распределенных механизмов. (например, распределенные транзакции, распределенные блокировки и т. д.), основная функция которых заключается в обеспечении согласованности данных:

Как показано на рисунке выше, если предположить, что переменная a — это оставшиеся запасы, значение равно 1. В это время три пользователя приходят, чтобы разместить заказ, и ровно три запроса распределяются по трем различным узлам обслуживания. Проверьте оставшийся инвентарь и обнаружите, что есть еще 1, а затем вычтите его, что приведет к отрицательному инвентарю, и есть два пользователя, которые не отправили товар, что обычно называется перепроданностью. Такая ситуация недопустима, пользователь будет драться с бизнесом, бизнес поссорится с вашим руководителем, а потом вы соберете портфель и пойдете домой!

В этом сценарии нам нужен способ решить эту проблему,Это проблема, которую должны решить распределенные блокировки..

2 Реализация и характеристики распределенных замков

2.1 Реализация распределенных блокировок

Локальные блокировки могут поддерживаться самим языком.Для реализации распределенных блокировок необходимо полагаться на промежуточное ПО, базы данных, Redis, Zookeeper и т. д. В основном существуют следующие методы реализации:
1) Memcached: используйте команду добавления Memcached. Эта команда является атомарной операцией, и добавление может быть успешным только в том случае, если ключ не существует, что означает, что поток получает блокировку.
2) Redis: аналогично Memcached, с использованием команды Redis setnx. Эта команда также является атомарной операцией, и набор может быть выполнен успешно только в том случае, если ключ не существует.
3) Zookeeper: используйте последовательные временные узлы Zookeeper для реализации распределенных блокировок и очередей ожидания. Первоначальная цель дизайна Zookeeper — реализовать распределенную службу блокировки.
4) Chubby: крупнозернистый распределенный сервис блокировки, реализованный Google.Нижний уровень использует алгоритм консенсуса Paxos.

2.2 Характеристики распределенных замков

1) В среде распределенной системы метод может выполняться только одним потоком одной машины одновременно.
2) Высокодоступная блокировка захвата и блокировка освобождения.
3) Высокопроизводительное получение и снятие блокировки.
4) Он имеет реентерабельные функции.
5) Он имеет механизм отказа блокировки для предотвращения взаимоблокировки.
6) Он имеет функцию неблокирующей блокировки, то есть, если блокировка не получена, он напрямую возвращает отказ в получении блокировки.

3 Redisson реализует распределенную блокировку Redis и принцип реализации

3.1 Добавить зависимости

<dependency>
    <groupId>org.redisson</groupId>
    <artifactId>redisson-spring-boot-starter</artifactId>
    <version>3.12.4</version>
</dependency>

3.2 Тестовый вид

Если количество запасов равно 100, оно будет вызвано один раз минус 1. Если оно меньше или равно 0, будет возвращено значение false, что означает, что заказ не выполнен.

@Component
public class RedissonLock {
    private static Integer inventory = 100;
    /**
     * 测试
     *
     * @return true:下单成功 false:下单失败
     */
    public Boolean redisLockTest(){
        // 获取锁实例
        RLock inventoryLock = RedissonService.getRLock("inventory-number");
        try {
            // 加锁
            inventoryLock.lock();
            if (inventory <= 0){
                return false;
            }
            inventory--;
            System.out.println("线程名称:" + Thread.currentThread().getName() + "剩余数量:" + RedissonLock.inventory);
        }catch (Exception e){
            e.printStackTrace();
        }finally {
            // 释放锁
            inventoryLock.unlock();
        }
        return true;
    }
}

Измерение давления с помощью jmeter:

Группа потоков 100 выполняется в течение 20 секунд:

Ответ утверждает true для true и false для отказа:

результат:

3.3 Случаи получения замков

RLock inventoryLock = RedissonService.getRLock("inventory-number"); Этот раздел является экземпляром получения блокировки. Inventory-number - указанное имя блокировки. После ввода метода getLock(String name) вы можете увидеть, что экземпляр получение блокировки построено в методе RedissonLock, инициализирующем некоторые свойства.

public RLock getLock(String name) {
        return new RedissonLock(this.connectionManager.getCommandExecutor(), name);
}

Взгляните на конструктор RedissonLock:

public RedissonLock(CommandAsyncExecutor commandExecutor, String name) {
        super(commandExecutor, name);
        //命令执行器
        this.commandExecutor = commandExecutor;
        //UUID字符串(MasterSlaveConnectionManager类的构造函数 传入UUID)
        this.id = commandExecutor.getConnectionManager().getId();
        //内部锁过期时间(防止死锁,默认时间为30s)
        this.internalLockLeaseTime = commandExecutor.getConnectionManager().getCfg().getLockWatchdogTimeout();
        //uuid+传进来的锁名称
        this.entryName = this.id + ":" + name;
        //redis消息体
        this.pubSub = commandExecutor.getConnectionManager().getSubscribeService().getLockPubSub();
    }

Время истечения внутренней блокировки (по умолчанию 30 с, если бизнес-код не был выполнен по истечении этого времени, время истечения срока действия будет автоматически продлено):

3.4 Блокировка

inventoryLock.lock(); Этот код означает блокировку, шаг за шагом в исходный код, чтобы увидеть, сначала посмотрите следующий метод lock():

public void lock() {
        try {
            this.lock(-1L, (TimeUnit)null, false);
        } catch (InterruptedException var2) {
            throw new IllegalStateException();
        }
    }

Вы можете видеть, что здесь установлены некоторые значения по умолчанию, а затем продолжить вызов метода lock() с параметрами, который также находится здесь для завершения логики блокировки.Исходный код выглядит следующим образом:

private void lock(long leaseTime, TimeUnit unit, boolean interruptibly) throws InterruptedException {
        // 线程ID
        long threadId = Thread.currentThread().getId();
        // 尝试获取锁
        Long ttl = this.tryAcquire(leaseTime, unit, threadId);
        // 如果过期时间等于null,则表示获取到锁,直接返回,不等于null继续往下执行
        if (ttl != null) {
            // 如果获取锁失败,则订阅到对应这个锁的channel
            RFuture<RedissonLockEntry> future = this.subscribe(threadId);
            if (interruptibly) {
                // 可中断订阅
                this.commandExecutor.syncSubscriptionInterrupted(future);
            } else {
                // 不可中断订阅
                this.commandExecutor.syncSubscription(future);
            }
            try {
                // 不断循环
                while(true) {
                    // 再次尝试获取锁
                    ttl = this.tryAcquire(leaseTime, unit, threadId);
                    // ttl(过期时间)为空,说明成功获取锁,返回
                    if (ttl == null) {
                        return;
                    }
                    // ttl(过期时间)大于0 则等待ttl时间后继续尝试获取
                    if (ttl >= 0L) {
                        try {
                            ((RedissonLockEntry)future.getNow()).getLatch().tryAcquire(ttl, TimeUnit.MILLISECONDS);
                        } catch (InterruptedException var13) {
                            if (interruptibly) {
                                throw var13;
                            }
                            ((RedissonLockEntry)future.getNow()).getLatch().tryAcquire(ttl, TimeUnit.MILLISECONDS);
                        }
                    } else if (interruptibly) {
                        ((RedissonLockEntry)future.getNow()).getLatch().acquire();
                    } else {
                        ((RedissonLockEntry)future.getNow()).getLatch().acquireUninterruptibly();
                    }
                }
            } finally {
                // 取消对channel的订阅
                this.unsubscribe(future, threadId);
            }
        }
    }

Давайте взглянем на метод tryAcquire для получения блокировки:

private Long tryAcquire(long leaseTime, TimeUnit unit, long threadId) {
        return (Long)this.get(this.tryAcquireAsync(leaseTime, unit, threadId));
    }

Заходим и видим метод tryAcquireAsync:

private <T> RFuture<Long> tryAcquireAsync(long leaseTime, TimeUnit unit, long threadId) {
        // 有设置过期时间
        if (leaseTime != -1L) {
            return this.tryLockInnerAsync(leaseTime, unit, threadId, RedisCommands.EVAL_LONG);
        } else {
            // 没有设置过期时间
            RFuture<Long> ttlRemainingFuture = this.tryLockInnerAsync(this.commandExecutor.getConnectionManager().getCfg().getLockWatchdogTimeout(), TimeUnit.MILLISECONDS, threadId, RedisCommands.EVAL_LONG);
            ttlRemainingFuture.onComplete((ttlRemaining, e) -> {
                if (e == null) {
                    if (ttlRemaining == null) {
                        this.scheduleExpirationRenewal(threadId);
                    }
                }
            });
            return ttlRemainingFuture;
        }
    }
  • Метод tryLockInnerAsync фактически выполняет логику получения блокировки, которая является частью кода сценария LUA. Здесь используется хэш-структура данных.
<T> RFuture<T> tryLockInnerAsync(long leaseTime, TimeUnit unit, long threadId, RedisStrictCommand<T> command) {
        this.internalLockLeaseTime = unit.toMillis(leaseTime);
        return this.commandExecutor.evalWriteAsync(this.getName(), LongCodec.INSTANCE, command, 
        // 如果锁不存在,则通过hset设置它的值,并设置过期时间
        "if (redis.call('exists', KEYS[1]) == 0) then redis.call('hincrby', KEYS[1], ARGV[2], 1); redis.call('pexpire', KEYS[1], ARGV[1]); return nil; end; 
        // 如果锁已存在,并且锁的是当前线程,则通过hincrby给数值递增1(这里显示了redis分布式锁的可重入性)
        if (redis.call('hexists', KEYS[1], ARGV[2]) == 1) then redis.call('hincrby', KEYS[1], ARGV[2], 1); redis.call('pexpire', KEYS[1], ARGV[1]); return nil; end; 
        // 如果锁已存在,但并非本线程,则返回过期时间ttl
        return redis.call('pttl', KEYS[1]);", Collections.singletonList(this.getName()), new Object[]{this.internalLockLeaseTime, this.getLockName(threadId)});
    }

KEYS[1] представляет ключ, который вы блокируете, например: RLock inventoryLock = RedissonService.getRLock("inventory-number"); здесь ключом блокировки, который вы установили для блокировки, является "inventory-number".
ARGV[1] представляет время жизни ключа блокировки по умолчанию, как показано на снимке экрана выше, время по умолчанию составляет 30 секунд.
ARGV[2] представляет идентификатор заблокированного клиента, подобный следующему: 8743c9c0-0795-4907-87fd-6c719a6b4586:1

Приведенный выше код LUA не выглядит очень сложным, и есть три суждения:

Судя по тому, что защелки не существует, если блокировки не существует, установите значение и срок действия, и блокировка пройдет успешно.
Судя по гексистам, если блокировка уже существует, а текущий поток заблокирован, это оказывается блокировкой с повторным входом и блокировка прошла успешно.Значение +1 ARGV[2], которое изначально было 1, теперь равно 2. Конечно , когда он отпущен И отпустите дважды.
Если блокировка уже существует, но блокировка не является текущим потоком, это доказывает, что другой поток удерживает блокировку. Возвращает время истечения текущей блокировки, если блокировка не удалась

3.5 Разблокировать

inventoryLock.unlock(); Этот код означает разблокировку, как и прежде, шаг за шагом переходим к исходному коду и видим следующий метод unlock():

public void unlock() {
        try {
            this.get(this.unlockAsync(Thread.currentThread().getId()));
        } catch (RedisException var2) {
            if (var2.getCause() instanceof IllegalMonitorStateException) {
                throw (IllegalMonitorStateException)var2.getCause();
            } else {
                throw var2;
            }
        }
    }

Войдите в unlockAsync(), чтобы увидеть, это метод разблокировки:

public RFuture<Void> unlockAsync(long threadId) {
        RPromise<Void> result = new RedissonPromise();
        // 释放锁的方法
        RFuture<Boolean> future = this.unlockInnerAsync(threadId);
        // 添加监听器 解锁opStatus:返回值
        future.onComplete((opStatus, e) -> {
            this.cancelExpirationRenewal(threadId);
            if (e != null) {
                result.tryFailure(e);
            //如果返回null,则证明解锁的线程和当前锁不是同一个线程,抛出异常
            } else if (opStatus == null) {
                IllegalMonitorStateException cause = new IllegalMonitorStateException("attempt to unlock lock, not locked by current thread by node id: " + this.id + " thread-id: " + threadId);
                result.tryFailure(cause);
            } else {
                // 解锁成功
                result.trySuccess((Object)null);
            }
        });
        return result;
    }

Зайдите и посмотрите способ снятия блокировки: unlockInnerAsync():

protected RFuture<Boolean> unlockInnerAsync(long threadId) {
        return this.commandExecutor.evalWriteAsync(this.getName(), LongCodec.INSTANCE, RedisCommands.EVAL_BOOLEAN, 
        // 如果释放锁的线程和已存在锁的线程不是同一个线程,返回null
        "if (redis.call('hexists', KEYS[1], ARGV[3]) == 0) then return nil;end; 
        // 如果是同一个线程,就通过hincrby减1的方式,释放一次锁
        local counter = redis.call('hincrby', KEYS[1], ARGV[3], -1);
        // 若剩余次数大于0 ,则刷新过期时间
        if (counter > 0) then redis.call('pexpire', KEYS[1], ARGV[2]); return 0; 
        // 其他就证明锁已经释放,删除key并发布锁释放的消息
        else redis.call('del', KEYS[1]); redis.call('publish', KEYS[2], ARGV[1]); return 1; end; 
        return nil;", 
        Arrays.asList(this.getName(), this.getChannelName()), new Object[]{LockPubSub.UNLOCK_MESSAGE, this.internalLockLeaseTime, this.getLockName(threadId)});
    }

Приведенный выше код представляет собой логику для снятия блокировки. Точно так же у него также есть три суждения:

Если разблокированный поток и текущий поток блокировки не совпадают, разблокировка завершается неудачно и выдается исключение.
Если разблокированный поток и текущий поток блокировки совпадают, блокировка снимается один раз путем вычитания 1 из hincrby. Если оставшееся количество раз все еще больше 0, это подтверждается повторной блокировкой, и время истечения срока действия снова обновляется.
Блокировка больше не существует, сообщение о снятии блокировки публикуется через публикацию, и разблокировка выполнена успешно.

Это здесь.Лучше пройти его руками.Вы поймете это,если сами попробуете.Можете ли вы поставить большой палец вверх боссам,когда увидите это?Я хочу увидеть такой ужасающий мир второго уровня, Спасибо!

Если вам это нужно, вы можете обратить внимание на мой официальный аккаунт, и я буду обновлять технические статьи, связанные с Java в режиме реального времени.Также в официальном аккаунте есть некоторые практические материалы, такие как видеоурок системы Java seckill, обучение материалы темной лошадки 2019 (версия IDEA) и резюме вопросов интервью BAT (полная классификация), компьютеры MAC, обычно используемые установочные пакеты (некоторые из них покупаются на Taobao, уже PJ).