1. Обзор Redisson
Что такое Редиссон?——Redisson Wiki
Redisson — это сетка данных в памяти Java, основанная на Redis. Он не только предоставляет ряд распределенных общих объектов Java, но также предоставляет множество распределенных сервисов. К ним относятся (BitSet, Set, Multimap, SortedSet, Map, List, Queue, BlockingQueue, Deque, BlockingDeque, Semaphore, Lock, AtomicLong, CountDownLatch, Publish/Subscribe, фильтр Блума, удаленная служба, кэш Spring, служба Executor, служба Live Object , служба планировщика) Redisson предоставляет самый простой и удобный способ использования Redis. Цель Redisson — способствовать разделению ответственности пользователей Redis, чтобы пользователи могли больше сосредоточиться на обработке бизнес-логики.
Распределенный инструмент, основанный на реализации Redis, с базовыми распределенными объектами и расширенными и абстрактными распределенными службами, обеспечивающий решение большинства распределенных проблем для каждого программиста, пытающегося воссоздать распределенное колесо.
В чем разница между Redisson и Jedis, Lettuce? Не Лэй Фэн и башня Лэй Фэн
Отличие Redisson от двух других как в графическом интерфейсе с мышкой, так и в файле с командной строкой. Redisson — это абстракция более высокого уровня, а Jedis и Lettuce — это оболочки для команд Redis.
- Jedis — это набор инструментов, официально запущенный Redis для подключения клиентов Redis через Java, обеспечивающий поддержку различных команд для Redis.
- Lettuce – это расширяемый потокобезопасный клиент Redis. Платформа связи основана на Netty и поддерживает расширенные функции Redis, такие как сигнальные устройства, кластеры, конвейеры, автоматическое повторное подключение и модель данных Redis. Начиная с Spring Boot 2.x, Lettuce заменил Jedis в качестве предпочтительного клиента для Redis.
- Redisson — это комплексное промежуточное программное обеспечение нового типа, основанное на Redis и коммуникациях на основе Netty.Это лучшая модель для использования Redis в разработке на уровне предприятия.
Jedis инкапсулирует команды Redis, а Lettuce также имеет более богатые API и поддерживает кластеризацию и другие режимы. Но оба они пока что дают вам только основу для работы с базой данных Redis, в то время как Redisson создала зрелое распределенное решение на основе Redis, Lua и Netty и даже набор инструментов, рекомендованный официальными лицами Redis.
2. Распределенная блокировка
Как реализовать распределенную блокировку?
Распределенные блокировки являются жестким требованием для параллельных сервисов.Хотя существуют различные реализации: у ZooKeeper есть последовательные узлы Znode, у баз данных есть блокировки на уровне таблицы и счастливые/пессимистические блокировки, а у Redis есть setNx, но у них одна и та же цель.В конце концов, мы должны вернуться к взаимному исключению.Эта статья знакомит с Redisson, а затем возьмем Redis в качестве примера.
Как написать простую распределенную блокировку Redis?
Возьмите Spring Data Redis в качестве примера, используйте RedisTemplate для работы с Redis (setIfAbsent уже является командой слияния setNx + expire), как показано ниже.
// 加锁
public Boolean tryLock(String key, String value, long timeout, TimeUnit unit) {
return redisTemplate.opsForValue().setIfAbsent(key, value, timeout, unit);
}
// 解锁,防止删错别人的锁,以uuid为value校验是否自己的锁
public void unlock(String lockName, String uuid) {
if(uuid.equals(redisTemplate.opsForValue().get(lockName)){ redisTemplate.opsForValue().del(lockName); }
}
// 结构
if(tryLock){
// todo
}finally{
unlock;
}
Простая версия 1.0 завершена, и умный Сяо Чжан может с первого взгляда увидеть, что это блокировка, но операции get и del не являются атомарными, и когда параллелизм велик, безопасность процесса не может быть гарантирована. Итак, Сяо Чжан предложил использовать скрипт Lua.
Что такое Lua-скриптинг?
Сценарий Lua — это легкий и компактный язык, встроенный в Redis, и его выполнение осуществляется через Redis.eval/evalshaКоманда для запуска, инкапсулируйте операцию в Lua-скрипт, например атомарную операцию, которая выполняется один раз.
Итак, версия 2.0 удаляется через Lua-скрипт.
lockDel.lua выглядит следующим образом
if redis.call('get', KEYS[1]) == ARGV[1]
then
-- 执行删除操作
return redis.call('del', KEYS[1])
else
-- 不成功,返回0
return 0
end
Выполнить команду Lua во время операции удаления
// 解锁脚本
DefaultRedisScript<Object> unlockScript = new DefaultRedisScript();
unlockScript.setScriptSource(new ResourceScriptSource(new ClassPathResource("lockDel.lua")));
// 执行lua脚本解锁
redisTemplate.execute(unlockScript, Collections.singletonList(keyName), value);
2.0 больше похож на блокировку, но, похоже, чего-то не хватает. Сяо Чжан похлопал себя по голове, синхронизация и ReentrantLock очень плавные, потому что они обе являются реентерабельными блокировками, и поток не заблокируется, если он берет блокировку несколько раз. , Нам нужен повторный вход.
Как обеспечить повторный вход?
Повторный вход означает, что одному и тому же потоку разрешено получать одну и ту же блокировку несколько раз, не вызывая взаимоблокировки. Это дает хорошую идею для синхронизированных предвзятых блокировок. Реализация синхронизированного повторного входа осуществляется на уровне JVM. Заголовок объекта JAVA — MARK WORD. Идентификатор потока и счетчик скрыты в потоке, чтобы сделать оценку повторного входа в текущем потоке, чтобы избежать каждого CAS.
Когда поток получает доступ к синхронизированному блоку и получает блокировку, смещенный идентификатор потока будет храниться в записи блокировки в заголовке объекта и кадре стека. Позже потоку не нужно выполнять операции CAS для блокировки и разблокировки при входе и выходе. Просто проверьте, хранит ли слово Mark в заголовке объекта предвзятую блокировку, указывающую на текущий поток. Если проверка прошла успешно, поток получил блокировку. Если тест не пройден, вам нужно снова проверить, установлен ли флаг блокировки смещения в Mark Word в 1: если нет, CAS конкурирует; если он установлен, CAS указывает блокировку смещения головы объекта на текущий поток.
Поддерживается еще один счетчик. Тот же поток будет увеличиваться на 1 при входе и затем уменьшаться на 1 при выходе. Его можно освободить, пока он не достигнет 0.
повторная блокировка
Чтобы сымитировать эту схему, нам нужно модифицировать Lua-скрипт:
1. Необходимо сохранить имя блокировкиlockName, получить замокидентификатор темыи соответствующая веткаВведите количество
2. Блокировка
Каждый раз, когда поток получает блокировку, определяйте, существует ли уже блокировка.
не существует
Установите хэш-ключ на идентификатор потока и инициализируйте значение равным 1.
Установить срок действия
Возвращает истину об успешном получении блокировки
существует
Продолжайте судить, существует ли хеш-ключ идентификатора текущего потока.
Exist, значение ключа потока +1, количество повторных входов увеличено на 1, установлено время истечения
Не существует, отказ блокировки возврата
3. Разблокировать
Каждый раз, когда поток приходит на разблокировку, определяйте, существует ли уже блокировка
существует
Есть ли хеш-ключ идентификатора потока, если да, то уменьшить его на 1, если нет, вернуть ошибку разблокировки
После вычитания 1 определите, равен ли оставшийся счетчик 0. Если он равен 0, это означает, что блокировка больше не нужна.Выполните команду del, чтобы удалить ее.
1. Структура хранения
Чтобы облегчить обслуживание этого объекта, мы используем хеш-структуру для хранения этих полей. Hash Redis похож на HashMap в Java и подходит для хранения объектов.
hset lockname1 threadId 1
установить имя дляlockname1Структура хеша, ключ структуры хэшаthreadId, значение значение равно1
hget lockname1 threadId
Получить значение threadId для lockname1
Структура хранения
lockname 锁名称 key1: threadId 唯一键,线程id value1: count 计数器,记录该线程获取锁的次数
структура в редисе
2. Сложение и вычитание счетчика
Когда один и тот же поток получает одну и ту же блокировку, нам нужно добавить и вычесть счетчик счетчика соответствующего потока.
Чтобы определить, существует ли ключ Redis, вы можете использоватьexists, а чтобы определить, существует ли хэш-ключ, вы можете использоватьhexists
И в Redis также есть команда автоинкремента хэшаhincrby
Каждый раз, когда он увеличивается на 1hincrby lockname1 threadId 1, уменьшено на 1hincrby lockname1 threadId -1
3. Решение о разблокировке
Когда блокировка больше не нужна, каждый раз, когда она разблокируется, счетчик уменьшается на 1, пока не достигнет 0, и выполняется удаление.
Комбинируя вышеуказанную структуру хранения и процесс оценки, блокировка и разблокировка Lua выглядит следующим образом.
блокировка lock.lua
local key = KEYS[1];
local threadId = ARGV[1];
local releaseTime = ARGV[2];
-- lockname不存在
if(redis.call('exists', key) == 0) then
redis.call('hset', key, threadId, '1');
redis.call('expire', key, releaseTime);
return 1;
end;
-- 当前线程已id存在
if(redis.call('hexists', key, threadId) == 1) then
redis.call('hincrby', key, threadId, '1');
redis.call('expire', key, releaseTime);
return 1;
end;
return 0;
разблокировать разблокировать.lua
local key = KEYS[1];
local threadId = ARGV[1];
-- lockname、threadId不存在
if (redis.call('hexists', key, threadId) == 0) then
return nil;
end;
-- 计数器-1
local count = redis.call('hincrby', key, threadId, -1);
-- 删除lock
if (count == 0) then
redis.call('del', key);
return nil;
end;
код
/**
* @description 原生redis实现分布式锁
* @date 2021/2/6 10:51 下午
**/
@Getter
@Setter
public class RedisLock {
private RedisTemplate redisTemplate;
private DefaultRedisScript<Long> lockScript;
private DefaultRedisScript<Object> unlockScript;
public RedisLock(RedisTemplate redisTemplate) {
this.redisTemplate = redisTemplate;
// 加载加锁的脚本
lockScript = new DefaultRedisScript<>();
this.lockScript.setScriptSource(new ResourceScriptSource(new ClassPathResource("lock.lua")));
this.lockScript.setResultType(Long.class);
// 加载释放锁的脚本
unlockScript = new DefaultRedisScript<>();
this.unlockScript.setScriptSource(new ResourceScriptSource(new ClassPathResource("unlock.lua")));
}
/**
* 获取锁
*/
public String tryLock(String lockName, long releaseTime) {
// 存入的线程信息的前缀
String key = UUID.randomUUID().toString();
// 执行脚本
Long result = (Long) redisTemplate.execute(
lockScript,
Collections.singletonList(lockName),
key + Thread.currentThread().getId(),
releaseTime);
if (result != null && result.intValue() == 1) {
return key;
} else {
return null;
}
}
/**
* 解锁
* @param lockName
* @param key
*/
public void unlock(String lockName, String key) {
redisTemplate.execute(unlockScript,
Collections.singletonList(lockName),
key + Thread.currentThread().getId()
);
}
}
На данный момент завершена распределенная блокировка, которая соответствует основным характеристикам взаимного исключения, повторного входа и предотвращения взаимоблокировок.
Строгий Сяо Чжан считает, что, хотя он достаточно стабилен, чтобы быть обычным мьютексом, в бизнесе всегда есть много особых случаев.Например, когда процесс A получает блокировку, поскольку время бизнес-операции слишком велико, блокировка снимается, но бизнес все еще выполняется.В этот момент процесс B может получить блокировку, как обычно, для бизнес-операций, а операции два процесса по-прежнему будут существовать и по-прежнему будут использоваться совместно..
И если человек, ответственный за хранение этой распределенной блокировкиПосле того, как узел Redis не работает, а замок находится в заблокированном состоянии, замок будет заблокирован..
Сяо Чжан не гангстер, потому что складские операции всегда так или иначе особенные.
Поэтому мы надеемся, что в этом случае ReleaseTime блокировки может быть продлен, чтобы отсрочить освобождение блокировки до тех пор, пока ожидаемый результат бизнеса не будет завершен.Операция непрерывного продления времени истечения блокировки для обеспечения завершения бизнеса выполнение - обновление блокировки.
Также распространено разделение чтения и записи.Бизнес, который больше читает и реже пишет, имеет блокировки чтения и записи для повышения производительности.
Расширение на данный момент превысило сложность простого колеса.Сяо Чжану достаточно просто обработки обновлений, чтобы выпить горшок, не говоря уже о производительности (максимальное время ожидания блокировки), элегантности (недействительное приложение блокировки), повторных попытках (сбой механизм повторных попыток) и другие аспекты нуждаются в изучении. Когда Сяо Чжан усердно думал, Сяо Бай рядом с ним подошел и посмотрел на Сяо Чжана. Ему было очень любопытно. Уже 2021 год, так почему бы просто не использовать Redisson?
У Redisson есть замок, который вам нужен.
3. Распределенный замок Redisson
Какова ситуация с использованием так называемой простой распределенной блокировки Redisson?
1. Зависимость
<!-- 原生,本章使用-->
<dependency>
<groupId>org.redisson</groupId>
<artifactId>redisson</artifactId>
<version>3.13.6</version>
</dependency>
<!-- 另一种Spring集成starter,本章未使用 -->
<dependency>
<groupId>org.redisson</groupId>
<artifactId>redisson-spring-boot-starter</artifactId>
<version>3.13.6</version>
</dependency>
2. Конфигурация
@Configuration
public class RedissionConfig {
@Value("${spring.redis.host}")
private String redisHost;
@Value("${spring.redis.password}")
private String password;
private int port = 6379;
@Bean
public RedissonClient getRedisson() {
Config config = new Config();
config.useSingleServer().
setAddress("redis://" + redisHost + ":" + port).
setPassword(password);
config.setCodec(new JsonJacksonCodec());
return Redisson.create(config);
}
}
3. Включите распределенные блокировки
@Resource
private RedissonClient redissonClient;
RLock rLock = redissonClient.getLock(lockName);
try {
boolean isLocked = rLock.tryLock(expireTime, TimeUnit.MILLISECONDS);
if (isLocked) {
// TODO
}
} catch (Exception e) {
rLock.unlock();
}
Кратко и ясно, нужен только один RLock.Поскольку рекомендуется Redisson, давайте заглянем внутрь и посмотрим, как он это реализует.
4. Блокировка
RLock – это основной интерфейс распределенной блокировки Redisson. Он наследует интерфейс Lock параллельного пакета и собственный интерфейс RLockAsync. Возвращаемое значение RLockAsync – RFuture, что является основной логикой асинхронной реализации Redisson и основной позицией Netty.
Как блокируется RLock?
Войдите из RLock, найдите класс RedissonLock, найдитеtryLockЗатем метод передается клерку.tryAcquireOnceAsyncметод, это основной код для блокировки (есть отличия в реализации разных версий, и есть определенные отличия с последней 3.15.х, но основная логика остается неизменной. Вот 3.13.6 как пример)
private RFuture<Boolean> tryAcquireOnceAsync(long waitTime, long leaseTime, TimeUnit unit, long threadId) {
if (leaseTime != -1L) {
return this.tryLockInnerAsync(waitTime, leaseTime, unit, threadId, RedisCommands.EVAL_NULL_BOOLEAN);
} else {
RFuture<Boolean> ttlRemainingFuture = this.tryLockInnerAsync(waitTime, this.commandExecutor.getConnectionManager().getCfg().getLockWatchdogTimeout(), TimeUnit.MILLISECONDS, threadId, RedisCommands.EVAL_NULL_BOOLEAN);
ttlRemainingFuture.onComplete((ttlRemaining, e) -> {
if (e == null) {
if (ttlRemaining) {
this.scheduleExpirationRenewal(threadId);
}
}
});
return ttlRemainingFuture;
}
}
Здесь есть две ветви оценки времени аренды, которые на самом деле установлены ли время истечения при блокировке, и когда время истечения (-1) не установлено, будетwatchDogизобновление блокировки(ниже), задача обновления, которая регистрирует событие блокировки. Давайте сначала посмотрим на время истечения срока действияtryLockInnerAsyncчасть,
evalWriteAsync — это запись для команды eval для выполнения 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, "if (redis.call('exists', KEYS[1]) == 0) then redis.call('hset', KEYS[1], ARGV[2], 1); redis.call('pexpire', KEYS[1], ARGV[1]); return nil; end; 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; return redis.call('pttl', KEYS[1]);", Collections.singletonList(this.getName()), new Object[]{this.internalLockLeaseTime, this.getLockName(threadId)});
}
Вот настоящее лицо, где команда eval выполняет Lua-скрипт, здесь Lua-скрипт расширяется
-- 不存在该key时
if (redis.call('exists', KEYS[1]) == 0) then
-- 新增该锁并且hash中该线程id对应的count置1
redis.call('hincrby', KEYS[1], ARGV[2], 1);
-- 设置过期时间
redis.call('pexpire', KEYS[1], ARGV[1]);
return nil;
end;
-- 存在该key 并且 hash中线程id的key也存在
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;
return redis.call('pttl', KEYS[1]);
Это почти то же самое, что и скрипт, который мы писали ранее для кастомной распределенной блокировки, похоже, так же реализован и редиссон.Анализ конкретных параметров:
// keyName
KEYS[1] = Collections.singletonList(this.getName())
// leaseTime
ARGV[1] = this.internalLockLeaseTime
// uuid+threadId组合的唯一值
ARGV[2] = this.getLockName(threadId)
Всего 3 параметра завершают часть логики:
Определить, есть ли у блокировки уже соответствующая хэш-таблица,
• Если соответствующей хэш-таблицы нет: затем установите ключ записи в хэш-таблице в качестве имени блокировки и значение 1, а затем установите время истечения срока действия хэш-таблицы в виде LeezeTime.
• Имеется соответствующая хеш-таблица: выполнить операцию +1 над значением lockName, то есть подсчитать количество записей, а затем установить время истечения LeezeTime.
• Наконец вернуть оставшееся время ttl этой блокировки
Он ничем не отличается от вышеупомянутого пользовательского замка.
В этом случае шаг разблокировки также должен иметь соответствующую операцию -1, а затем смотреть метод разблокировки, а также искать имя метода, вплоть до
protected RFuture<Boolean> unlockInnerAsync(long threadId) {
return this.commandExecutor.evalWriteAsync(this.getName(), LongCodec.INSTANCE, RedisCommands.EVAL_BOOLEAN, "if (redis.call('exists', KEYS[1]) == 0) then redis.call('publish', KEYS[2], ARGV[1]); return 1; end;if (redis.call('hexists', KEYS[1], ARGV[3]) == 0) then return nil;end; local counter = redis.call('hincrby', KEYS[1], ARGV[3], -1); if (counter > 0) then redis.call('pexpire', KEYS[1], ARGV[2]); return 0; 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.unlockMessage, this.internalLockLeaseTime, this.getLockName(threadId)});
}
Выньте часть Lua
-- 不存在key
if (redis.call('hexists', KEYS[1], ARGV[3]) == 0) then
return nil;
end;
-- 计数器 -1
local counter = redis.call('hincrby', KEYS[1], ARGV[3], -1);
if (counter > 0) then
-- 过期时间重设
redis.call('pexpire', KEYS[1], ARGV[2]);
return 0;
else
-- 删除并发布解锁消息
redis.call('del', KEYS[1]);
redis.call('publish', KEYS[2], ARGV[1]);
return 1;
end;
return nil;
КЛЮЧИ Lua имеют 2Arrays.asList(getName(), getChannelName())
name 锁名称
channelName,用于pubSub发布消息的channel名称
Есть три переменные ARGVLockPubSub.UNLOCK_MESSAGE, internalLockLeaseTime, getLockName(threadId)
LockPubSub.UNLOCK_MESSAGE,channel发送消息的类别,此处解锁为0
internalLockLeaseTime,watchDog配置的超时时间,默认为30s
lockName 这里的lockName指的是uuid和threadId组合的唯一值
Действуйте следующим образом:
1. Если блокировки не существует, вернуть nil;
2. Если блокировка существует, установите для счетчика хеш-ключа потока значение -1,
3. Если счетчик counter>0, сбросить следующее время истечения и вернуть 0, в противном случае удалить блокировку, опубликовать сообщение разблокировки unlockMessage и вернуть 1;
Среди них, когда используется unLock, Redis публикует и подписывается на PubSub для завершения уведомления о сообщении.
Шаг подписки находится в методе блокировки записи блокировки RedissonLock.
long threadId = Thread.currentThread().getId();
Long ttl = this.tryAcquire(-1L, leaseTime, unit, threadId);
if (ttl != null) {
// 订阅
RFuture<RedissonLockEntry> future = this.subscribe(threadId);
if (interruptibly) {
this.commandExecutor.syncSubscriptionInterrupted(future);
} else {
this.commandExecutor.syncSubscription(future);
}
// 省略
Когда блокировка занята другими потоками, отслеживая уведомление об освобождении блокировки (когда другие потоки освобождают блокировку через RedissonLock, уведомление будет отправлено через функцию публикации и подписки pub/sub) и ожидая снятия блокировки. выпущенные другими потоками, также является способом избежать вращения.Общий инструмент повышения эффективности.
1. Разблокировать сообщения
Для того, чтобы узнать, что было уведомлено и что было сделано после уведомления, введите LockPubSub.
Есть только один очевидный метод прослушивания onMessage, подписка и освобождение семафора которого находятся в родительском классе PublishSubscribe, мы фокусируемся только на фактической работе событий прослушивания
protected void onMessage(RedissonLockEntry value, Long message) {
Runnable runnableToExecute;
if (message.equals(unlockMessage)) {
// 从监听器队列取监听线程执行监听回调
runnableToExecute = (Runnable)value.getListeners().poll();
if (runnableToExecute != null) {
runnableToExecute.run();
}
// getLatch()返回的是Semaphore,信号量,此处是释放信号量
// 释放信号量后会唤醒等待的entry.getLatch().tryAcquire去再次尝试申请锁
value.getLatch().release();
} else if (message.equals(readUnlockMessage)) {
while(true) {
runnableToExecute = (Runnable)value.getListeners().poll();
if (runnableToExecute == null) {
value.getLatch().release(value.getLatch().getQueueLength());
break;
}
runnableToExecute.run();
}
}
}
нашел одинСообщение о разблокировке по умолчанию, одно из них — **сообщение о разблокировке блокировки чтения**, ** потому что redisson предоставляет блокировки чтения-записи, а блокировки чтения-записи условия чтения-записи взаимоисключающие с условиями чтения-записи, записи-записи, мы смотрим только на приведенное выше Сообщение разблокировки по умолчанию ветка unlockMessage
Слушатель LockPubSub в конечном итоге делает 2 вещи
-
runnableToExecute.run() выполняет обратный вызов прослушивателя
-
value.getLatch().release(); освободить семафор
Редиссон черезLockPubSubОтслеживайте сообщение разблокировки, выполняйте обратный вызов монитора и освобождайте семафор, чтобы уведомить ожидающий поток о том, что блокировку можно снова захватить.
Затем вернитесь и посмотрите на другую ветку tryAcquireOnceAsync.
private RFuture<Boolean> tryAcquireOnceAsync(long waitTime, long leaseTime, TimeUnit unit, long threadId) {
if (leaseTime != -1L) {
return this.tryLockInnerAsync(waitTime, leaseTime, unit, threadId, RedisCommands.EVAL_NULL_BOOLEAN);
} else {
RFuture<Boolean> ttlRemainingFuture = this.tryLockInnerAsync(waitTime, this.commandExecutor.getConnectionManager().getCfg().getLockWatchdogTimeout(), TimeUnit.MILLISECONDS, threadId, RedisCommands.EVAL_NULL_BOOLEAN);
ttlRemainingFuture.onComplete((ttlRemaining, e) -> {
if (e == null) {
if (ttlRemaining) {
this.scheduleExpirationRenewal(threadId);
}
}
});
return ttlRemainingFuture;
}
}
Видно, что при отсутствии тайм-аута после выполнения операции блокировки выполняется еще и кусок необъяснимой логики
ttlRemainingFuture.onComplete((ttlRemaining, e) -> {
if (e == null) {
if (ttlRemaining) {
this.scheduleExpirationRenewal(threadId);
}
}
})
Это включает в себя модель Netty Future/Promise-Listener (ссылкаАсинхронное программирование в Netty), почти все Redisson общаются таким образом (поэтому Redisson реализован на основе механизма связи Netty), чтобы понять эту логику, можно сначала попробовать ее понять
В будущем Java бизнес-логика представляет собой класс реализации Callable или Runnable. Завершение выполнения call() или run() этого класса означает конец бизнес-логики. В механизме Promise бизнес-логика может быть установлена вручную в бизнес-логике Успех и неудача, удобнее следить за собственной бизнес-логикой.
Поверхностное значение этого кода заключается в том, что после выполнения операции асинхронной блокировки, если блокировка прошла успешно, подтверждается, следует ли выполнять запланированную задачу в зависимости от того, истек ли срок жизни, возвращенный завершением блокировки.
Эта рассчитанная на время задача является ядром watchDog.
2. Обновление блокировки
Просмотр RedissonLock.this.scheduleExpirationRenewal(threadId)
private void scheduleExpirationRenewal(long threadId) {
RedissonLock.ExpirationEntry entry = new RedissonLock.ExpirationEntry();
RedissonLock.ExpirationEntry oldEntry = (RedissonLock.ExpirationEntry)EXPIRATION_RENEWAL_MAP.putIfAbsent(this.getEntryName(), entry);
if (oldEntry != null) {
oldEntry.addThreadId(threadId);
} else {
entry.addThreadId(threadId);
this.renewExpiration();
}
}
private void renewExpiration() {
RedissonLock.ExpirationEntry ee = (RedissonLock.ExpirationEntry)EXPIRATION_RENEWAL_MAP.get(this.getEntryName());
if (ee != null) {
Timeout task = this.commandExecutor.getConnectionManager().newTimeout(new TimerTask() {
public void run(Timeout timeout) throws Exception {
RedissonLock.ExpirationEntry ent = (RedissonLock.ExpirationEntry)RedissonLock.EXPIRATION_RENEWAL_MAP.get(RedissonLock.this.getEntryName());
if (ent != null) {
Long threadId = ent.getFirstThreadId();
if (threadId != null) {
RFuture<Boolean> future = RedissonLock.this.renewExpirationAsync(threadId);
future.onComplete((res, e) -> {
if (e != null) {
RedissonLock.log.error("Can't update lock " + RedissonLock.this.getName() + " expiration", e);
} else {
if (res) {
RedissonLock.this.renewExpiration();
}
}
});
}
}
}
}, this.internalLockLeaseTime / 3L, TimeUnit.MILLISECONDS);
ee.setTimeout(task);
}
}
С точки зрения разделения, этот длинный и постоянно вложенный код фактически выполняет несколько шагов.
• Добавьте задачу обратного вызова netty Timeout, которая выполняется каждые (internalLockLeaseTime / 3) миллисекунды, а метод выполнения — renewExpirationAsync.
• renewExpirationAsync сбрасывает время ожидания блокировки, регистрирует прослушиватель, и обратный вызов прослушивателя выполняет renewExpiration.
Lua renewExpirationAsync выглядит следующим образом
protected RFuture<Boolean> renewExpirationAsync(long threadId) {
return this.commandExecutor.evalWriteAsync(this.getName(), LongCodec.INSTANCE, RedisCommands.EVAL_BOOLEAN, "if (redis.call('hexists', KEYS[1], ARGV[2]) == 1) then redis.call('pexpire', KEYS[1], ARGV[1]); return 1; end; return 0;", Collections.singletonList(this.getName()), new Object[]{this.internalLockLeaseTime, this.getLockName(threadId)});
}
if (redis.call('hexists', KEYS[1], ARGV[2]) == 1) then
redis.call('pexpire', KEYS[1], ARGV[1]);
return 1;
end;
return 0;
Сброс периода тайм-аута.
Какова цель добавления этой логики в Redisson?
Цель состоит в том, чтобы гарантировать, что бизнес не будет затронут в определенных сценариях, таких как проблема, когда время выполнения задачи истекает, но не заканчивается, и блокировка была снята.
Когда поток удерживает блокировку из-за того, что время тайм-аута rentTime не установлено, Redisson по умолчанию настраивает 30 сек., включает сторожевой таймер, обновляет блокировку каждые 10 сек., поддерживает тайм-аут 30 сек. и удаляет блокировку до тех пор, пока задача не будет завершена.
Это Редиссонобновление блокировки, это,WatchDogОсновная идея реализации.
3. Обзор процесса
В общем введении кратко описан процесс:
- Потоки A и B конкурируют за блокировку.После того, как A ее получает, B блокирует
- Когда поток B заблокирован, он не является активным CAS, а подписывается на широковещательное сообщение о блокировке от PubSub.
- Блокировка снимается после завершения операции A, и поток B получает уведомление о подписке.
- B просыпается и начинает хвататься за замок и получать замок
Подробный процесс блокировки и разблокировки резюмируется следующим образом:
5. Честный замок
Повторно входящие блокировки, описанные выше, являются несправедливыми блокировками.Redisson также реализует справедливые блокировки на основе очереди Redis (List) и ZSet.
Каково определение справедливости?
Справедливость заключается в получении блокировок в соответствии с запросом клиента в порядке очереди, то есть по принципу FIFO, поэтому важно последовательное расположение очередей и контейнеров.
FairSync
Проверьте реализацию справедливой блокировки JUC ReentrantLock.
/**
* Sync object for fair locks
*/
static final class FairSync extends Sync {
private static final long serialVersionUID = -3000897897090466540L;
final void lock() {
acquire(1);
}
/**
* Fair version of tryAcquire. Don't grant access unless
* recursive call or no waiters or is first.
*/
protected final boolean tryAcquire(int acquires) {
final Thread current = Thread.currentThread();
int c = getState();
if (c == 0) {
if (!hasQueuedPredecessors() &&
compareAndSetState(0, acquires)) {
setExclusiveOwnerThread(current);
return true;
}
}
else if (current == getExclusiveOwnerThread()) {
int nextc = c + acquires;
if (nextc < 0)
throw new Error("Maximum lock count exceeded");
setState(nextc);
return true;
}
return false;
}
}
AQS предоставил всю реализацию, справедливо ли это или нет, зависит от того, выводит ли класс реализации логику узла по порядку.
AbstractQueuedSynchronizer — это базовая структура для создания блокировок или других компонентов синхронизации.Он завершает создание очереди потоков получения ресурсов через встроенные очереди FIFO.Он не реализует интерфейсы синхронизации, а только определяет несколько методов получения и освобождения состояния синхронизации для пользовательской синхронизации.Компонент использование (выше) поддерживает эксклюзивный и общий доступ, который представляет собой дизайн, основанный на модели шаблонного метода, который обеспечивает почву для справедливости/несправедливости.
Мы используем 2 изображения, чтобы кратко объяснить процесс ожидания AQS (из книги «Искусство параллельного программирования JAVA»)
Одна из них — синхронная очередь (двунаправленная очередь FIFO).Управляйте блок-схемой ссылки на поток, состоянием ожидания и предшествующими и последующими узлами, которые не могут получить статус синхронизации (не удается получить блокировки).
одинОбщий процесс исключительного получения состояния синхронизации, основной процесс вызова метода Acquisition(int arg)
Видно, что процесс получения блокировки
AQS поддерживает очередь синхронизации. Потоки, которым не удается получить состояние, присоединяются к очереди для вращения. Условием для удаления очереди или остановки вращения является то, что узел-предшественник успешно получает состояние синхронизации для головного узла.
И сравните другой класс несправедливой блокировкиNonfairSyncМожно обнаружить, что ключевой код для контроля справедливости и несправедливости лежит в методе hasQueuedPredecessors.
static final class NonfairSync extends Sync {
private static final long serialVersionUID = 7316153563782823691L;
/**
* Performs lock. Try immediate barge, backing up to normal
* acquire on failure.
*/
final void lock() {
if (compareAndSetState(0, 1))
setExclusiveOwnerThread(Thread.currentThread());
else
acquire(1);
}
protected final boolean tryAcquire(int acquires) {
return nonfairTryAcquire(acquires);
}
}
NonfairSync уменьшает условие оценки hasQueuedPredecessors, роль этого метода
Проверьте, есть ли у текущего узла в очереди синхронизации предшествующий узел, и верните true, если есть запрос на получение блокировки раньше, чем текущий поток.
Убедитесь, что первый узел (поток) очереди берется каждый раз для получения блокировки, что является правилом справедливости.
Почему JUC блокирует по умолчанию несправедливо?
Потому что, когда поток запрашивает блокировку, она успешно получена, если она получена для синхронизации состояния. При таком предположении вероятность того, что только что выпущенный поток снова получит состояние синхронизации, очень высока, так что другие потоки могут только ждать в очереди синхронизации. Но преимущество этого заключается в том, что несправедливая блокировка значительно снижает затраты на переключение контекста системного потока.
Видно, что цена честности — это производительность и пропускная способность.
В Redis нет AQS, но есть List и zSet, посмотрим, как Redisson добивается справедливости.
RedissonFairLock
Использование RedissonFairLock по-прежнему очень простое
RLock fairLock = redissonClient.getFairLock(lockName);
fairLock.lock();
RedissonFairLock наследуется от RedissonLock, а также находит метод реализации блокировки до концаtryLockInnerAsync.
Есть 2 длинных раздела Lua, но Debug обнаружил, что запись справедливой блокировки находится после команды == RedisCommands.EVAL_LONG. Этот раздел Lua длиннее и имеет много параметров. Мы сосредоточимся на анализе правил реализации Lua.
параметр
-- lua中的几个参数
KEYS = Arrays.<Object>asList(getName(), threadsQueueName, timeoutSetName)
KEYS[1]: lock_name, 锁名称
KEYS[2]: "redisson_lock_queue:{xxx}" 线程队列
KEYS[3]: "redisson_lock_timeout:{xxx}" 线程id对应的超时集合
ARGV = internalLockLeaseTime, getLockName(threadId), currentTime + threadWaitTime, currentTime
ARGV[1]: "{leaseTime}" 过期时间
ARGV[2]: "{Redisson.UUID}:{threadId}"
ARGV[3] = 当前时间 + 线程等待时间:(10:00:00) + 5000毫秒 = 10:00:05
ARGV[4] = 当前时间(10:00:00) 部署服务器时间,非redis-server服务器时间
Lua-скрипт для реализации справедливой блокировки
-- 1.死循环清除过期key
while true do
-- 获取头节点
local firstThreadId2 = redis.call('lindex', KEYS[2], 0);
-- 首次获取必空跳出循环
if firstThreadId2 == false then
break;
end;
-- 清除过期key
local timeout = tonumber(redis.call('zscore', KEYS[3], firstThreadId2));
if timeout <= tonumber(ARGV[4]) then
redis.call('zrem', KEYS[3], firstThreadId2);
redis.call('lpop', KEYS[2]);
else
break;
end;
end;
-- 2.不存在该锁 && (不存在线程等待队列 || 存在线程等待队列而且第一个节点就是此线程ID),加锁部分主要逻辑
if (redis.call('exists', KEYS[1]) == 0) and
((redis.call('exists', KEYS[2]) == 0) or (redis.call('lindex', KEYS[2], 0) == ARGV[2])) then
-- 弹出队列中线程id元素,删除Zset中该线程id对应的元素
redis.call('lpop', KEYS[2]);
redis.call('zrem', KEYS[3], ARGV[2]);
local keys = redis.call('zrange', KEYS[3], 0, -1);
-- 遍历zSet所有key,将key的超时时间(score) - 当前时间ms
for i = 1, #keys, 1 do
redis.call('zincrby', KEYS[3], -tonumber(ARGV[3]), keys[i]);
end;
-- 加锁设置锁过期时间
redis.call('hset', KEYS[1], ARGV[2], 1);
redis.call('pexpire', KEYS[1], ARGV[1]);
return nil;
end;
-- 3.线程存在,重入判断
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;
-- 4.返回当前线程剩余存活时间
local timeout = redis.call('zscore', KEYS[3], ARGV[2]);
if timeout ~= false then
-- 过期时间timeout的值在下方设置,此处的减法算出的依旧是当前线程的ttl
return timeout - tonumber(ARGV[3]) - tonumber(ARGV[4]);
end;
-- 5.尾节点剩余存活时间
local lastThreadId = redis.call('lindex', KEYS[2], -1);
local ttl;
-- 尾节点不空 && 尾节点非当前线程
if lastThreadId ~= false and lastThreadId ~= ARGV[2] then
-- 计算队尾节点剩余存活时间
ttl = tonumber(redis.call('zscore', KEYS[3], lastThreadId)) - tonumber(ARGV[4]);
else
-- 获取lock_name剩余存活时间
ttl = redis.call('pttl', KEYS[1]);
end;
-- 6.末尾排队
-- zSet 超时时间(score),尾节点ttl + 当前时间 + 5000ms + 当前时间,无则新增,有则更新
-- 线程id放入队列尾部排队,无则插入,有则不再插入
local timeout = ttl + tonumber(ARGV[3]) + tonumber(ARGV[4]);
if redis.call('zadd', KEYS[3], timeout, ARGV[2]) == 1 then
redis.call('rpush', KEYS[2], ARGV[2]);
end;
return ttl;
1. Шаги блокировки Fair Lock
С помощью приведенного выше Lua можно обнаружить, что ключевой структурой операции lua является список (list) и упорядоченный набор (zSet).
где список поддерживает очередь ожидающих потоковredisson_lock_queue:{xxx}, zSet поддерживает упорядоченный набор условий тайм-аута потока.redisson_lock_timeout:{xxx}, хотя lua длиннее, его можно разбить на 6 шагов
- очистка очереди
- Убедитесь, что в очереди есть только ожидающие потоки с неистекшим сроком действия.
- первый замок
- блокировка hset, время истечения срока действия pexpire
- повторное суждение
- То же, что и реентерабельная блокировка lua здесь
- вернуть ttl
- Вычислить ttl хвостового узла
- Начальное значение — оставшееся время истечения блокировки.
- очередь в конце
-
ttl + 2 * currentTime + waitTime — это формула расчета значения по умолчанию для оценки.
2. Моделирование
Если вы смоделируете следующую последовательность, вы поймете весь процесс блокировки redisson fair lock.
Предположим, что t1 10:00:00
t1: когда поток 1 получает блокировку в первый раз
1. Дождаться очереди безголового узла, выпрыгнуть из бесконечного цикла -> 2
2. Блокировка не существует && не установлена очередь ожидания потока
2.1 lpop, zerm и цинкrby — недопустимые операции.Вступает в силу только блокировка, что означает, что это первая блокировка, и после блокировки возвращается nil.
Блокировка выполнена успешно, поток 1 получает блокировку и завершается.
t2: поток 2 попытался получить блокировку (поток 1 не снял блокировку)
1. Дождаться очереди безголового узла, выпрыгнуть из бесконечного цикла -> 2
2. Замок не существует, не установлен -> 3
3. Нереентерабельные потоки -> 4
4.score не имеет значения -> 5
5. Хвостовой узел пуст, установите начальное значение ttl на ttl -> 6 из lock_name
6. Установите оценку времени ожидания zSet в соответствии с ttl + waitTime + currentTime + currentTime и присоединитесь к очереди ожидания, поток 2 является головным узлом.
score = 20S + 5000ms + 10:00:10 + 10:00:10 = 10:00:35 + 10:00:10
t3: поток 3 попытался получить блокировку (поток 1 не снял блокировку)
1. Подождите, пока в очереди появится головной узел
1.1 не истек -> 2
2. Замок не существует, если он не существует -> 3
3. Нереентерабельные потоки -> 4
4.score не имеет значения -> 5
5. Хвостовой узел не пуст && поток хвостового узла равен 2, а не текущему потоку
5.1 Выньте ранее установленный счет и вычтите текущее время: ttl = score - currentTime -> 6
6. Установите оценку времени ожидания zSet в соответствии с ttl + waitTime + currentTime + currentTime и присоединитесь к очереди ожидания.
score = 10S + 5000ms + 10:00:20 + 10:00:20 = 10:00:35 + 10:00:20
Таким образом, три потока, которым необходимо захватить блокировку, формируют очередь, упорядочивают идентификаторы ожидающих потоков в списке и сохраняют время истечения срока действия в zSet (для облегчения расстановки приоритетов). Среди них клиент потока 2 и клиент потока 3, которые возвращают ttl, всегда будут многократно выполнять сегмент Lua с определенным интервалом и пытаться заблокировать, чтобы он имел тот же эффект, что и AQS.
И когда поток 1 снимает блокировку (через Pub/Sub по-прежнему выпускается сообщение о разблокировке, уведомляющее другие потоки о его получении)
10:00:30 Поток 2 пытается получить блокировку (поток 1 снял блокировку)
1. Ожидание появления в очереди головного узла, срок действия которого не истек -> 2
2. Блокировка не существует, и головным узлом очереди ожидания является текущий установленный поток.
2.1 Удалить информацию очереди и информацию zSet текущего потока, время ожидания:
Поток 2 10:00:35 + 10:00:10 - 10:00:30 = 10:00:15
Поток 3 10:00:35 + 10:00:20 - 10:00:30 = 10:00:25
2.2 Поток 2 получает блокировку и сбрасывает время истечения
Блокировка выполнена успешно, поток 2 получает блокировку и завершается.
Структура очереди показана на рисунке
Сценарий освобождения справедливой блокировки аналогичен повторной блокировке с добавлением логики while true для очистки просроченного ключа в начале блокировки, которая не будет здесь описываться.
Из вышеизложенного видно, что игровой процесс справедливой блокировки Redisson аналогичен игровому процессу отложенной очереди.Ядро представляет собой комбинацию Redis List и структуры zSet, но также использует реализацию AQS для справки.Головной узел оценки времени один и тот же (watchDog), что гарантирует, что конкуренция замков является честной и взаимоисключающей. В параллельном сценарии в lua-скрипте оценка zSet решает проблему последовательной вставки и распределяет приоритет. И чтобы предотвратить очистку потока, который выходит из-за исключений, каждый запрос будет оценивать истечение срока действия головного узла и очищать его.Когда будет сделан окончательный выпуск, подписанный поток может быть уведомлен через CHANNEL для получения заблокировать, повторить первые шаги и успешно передать следующей последовательности поток.
6. Резюме
Общая реализация Redisson процесса распределенной разблокировки немного сложна. Автор Руи Гу провел углубленное исследование Netty, JUC и Redis и использовал множество расширенных функций и семантики, которые заслуживают дальнейшего изучения. реализация блокировки Redis для одной машины.Также предусмотрена блокировка Redisson в случае нескольких машин (MultiLock) и официально рекомендуемый красный замок (RedLock), которые будут подробно описаны в следующей главе.
Поэтому, когда вам действительно нужны распределенные блокировки, вы можете сначала найти их в Redisson.