redis分布式锁
ᅩᩩٚ፡ِ҅౮ԟబ҅ஙמᔱ̓ӣॡৼයӰ̈́ىဳᬯӻᘶᗑᝌӬ؎ኞጱૡٍՈ̶
GitHub https://github.com/JavaFamily ૪ත୯҅ํӞᕚय़ܯᶎᦶਠෆᘍᅩ̵ᩒාզ݊౯
ጱᔮڜᒍ̶
ڹ
ӤӞᒍᜓ౯کԧचԭzkړୗᲁጱਫሿ҅ᬯᒍᜓ੪᧔ӞӥचԭRedisጱړୗᲁਫሿމ̶
zkਫሿړୗᲁጱփᭆᳪғzkړୗᲁ
ࣁতکRedisړୗᲁԏڹ҅౯మ᪙य़ਹᘱᅩRedisጱचᏐᎣᦩ̶
᧔ӞӥRedisጱӷӻեғ
setnx ฎSET if Not eXists(ইຎӧਂࣁ҅ڞ SET)ጱᓌ̶ٟ
አဩইࢶ҅ইຎӧਂࣁset౮ۑᬬࢧintጱ1҅ᬯӻkeyਂࣁԧᬬࢧ0̶
ਖ਼꧊ value ىᘶک key ҅ଚਖ਼ key ጱኞਂᳵᦡԅ seconds (զᑁԅܔ֖)̶
ইຎ key ૪ᕪਂࣁ҅ setex եਖ਼ᥟٟ෯꧊̶
ํੜվ֎ᙗਧտወజӡӞset value ౮ۑ set time०ᨳ҅ᮎӧ੪؞ԧԍ҅ᬯࠡRedisਥᗑమکԧ̶
setex ฎӞӻܻৼ(atomic)֢҅ىᘶ꧊ᦡᗝኞਂᳵӷӻ֢ۖտࣁݶӞᳵٖਠ౮̶
SETNX key value
SETEX key seconds value
౯ᦡᗝԧ10ᑁጱ०පᳵ҅ttlեݢզັ፡ׯᦇ҅ᨮጱ᧔ก૪ᕪک๗ԧ̶
᪙य़ਹᦖᬯӷӻݷԞฎํܻࢩጱ҅ࢩԅ՜ժฎRedisਫሿړୗᲁጱىᲫ̶
ྋ
তڹᬮฎ፡፡࣋ว:
౯ׁᆐฎڠୌԧஉग़ӻᕚᑕ݄ಕٺପਂinventory҅ӧڊक़ጱପਂಕٺᶲଧݒԧ๋҅ᕣጱᕮຎԞฎӧ
ጱ̶
ܔےsynchronizedᘏLockᬯԶଉᥢ֢౯੪ӧ᧔ԧঅމ҅ᕮຎᙗਧฎጱ̶
౯ضਫሿӞӻᓌܔጱRedisᲁ҅ᆐݸ౯ժٚਫሿړୗᲁ҅ݢᚆๅොय़ਹጱቘᥴ̶
ᬮᦕӤᶎ౯᧔ᬦጱեԍ҅ਫሿӞӻܔጱٌਫྲᓌܔ֦҅ժضᘍӞӥ҅ڦஃӥ፡̶
setnx
ݢզ፡ک҅ᒫӞӻ౮ۑԧ҅ဌ᯽නᲁ҅ݸᶎጱ᮷०ᨳԧ҅ᛗᶲଧᳯ᷌ᳯ᷌ฎᥴ٬ԧ҅ݝᥝےᲁ҅ᖽන
ݸᶎጱ೭ک҅᯽නইྌሾ҅੪ᚆכᦤೲᆙᶲଧಗᤈ̶
֕ฎ֦ժԞݎሿᳯ᷌ԧ҅ᬮฎӞጱ҅ᒫӞӻ՚set౮ۑԧ҅֕ฎᑱᆐ೯ԧ҅ᮎᲁ੪Ӟፗࣁᮎ෫ဩک᯽
න҅ݸᶎጱᕚᑕԞᬱӧکᲁ݈҅ྒᲁԧ̶
ಅզ....
setex
Ꭳ᭲౯ԏڹ᧔ᬯӻեጱܻࢩԧމ҅ᦡᗝӞӻᬦ๗ᳵ҅੪ᓒᕚᑕ1೯ԧ҅Ԟտࣁ०පᳵکԧ҅ᛔۖ
᯽න̶
౯ᬯ᯾੪አکԧnxpxጱᕮݳ݇හ҅੪ฎset꧊ଚӬےԧᬦ๗ᳵ҅ᬯ᯾౯ᬮᦡᗝԧӞӻᬦ๗ᳵ҅੪
ฎᬯᳵٖইຎᒫԫӻဌ೭کᒫӞӻጱᲁ҅੪ᭅڊᴥलԧ҅ࢩԅݢᚆฎਮಁᒒෙᬳԧ̶
ےᲁ
ෆ֛ےᲁጱ᭦ᬋྲᓌܔ҅य़ਹचӤ᮷ᚆ፡҅ӧᬦ౯೭ک୮ڹᳵ݄ٺতᳵጱ֢ఽᥧํᅩ
ᒨ҅ System.currentTimeMillis()ၾᘙஉय़ጱ̶
/**
*
*/
public boolean lock(String id) {
Long start = System.currentTimeMillis();
try {
for (; ; ) {
//SETեᬬࢧOK ҅ڞᦤก឴ݐᲁ౮ۑ
String lock = jedis.set(LOCK_KEY, id, params);
if ("OK".equals(lock)) {
return true;
}
//ވڞሾᒵஇ҅ࣁtimeoutᳵٖՖ๚឴ݐکᲁ҅ڞ឴ݐ०ᨳ
long l = System.currentTimeMillis() - start;
System.currentTimeMillisၾᘙय़҅ྯӻᕚᑕᬰ᮷ᬯ҅౯ԏڹٟդᎱ҅੪տࣁ๐ۓސۖጱײ҅
Ӟӻᕚᑕӧෙ݄೭҅᧣አොፗള឴ݐ꧊੪অԧ҅ӧᬦԞӧฎ๋սᥴ҅෭๗ᔄᬮฎํஉग़অොဩጱ̶
if (l >= timeout) {
return false;
}
try {
Thread.sleep(100);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
} finally {
jedis.close();
}
}
@Service
public class TimeServcie {
private static long time;
static {
new Thread(new Runnable(){
@Override
public void run() {
while (true){
try {
Thread.sleep(5);
} catch (InterruptedException e) {
e.printStackTrace();
}
long cur = System.currentTimeMillis();
setTime(cur);
}
}
}).start();
}
public static long getTime() {
return time;
}
public static void setTime(long time) {
TimeServcie.time = time;
}
ᥴᲁ
ᥴᲁጱ᭦ᬋๅےᓌܔ҅੪ฎӞྦྷLuaጱ೪ᤰ҅Key؉ԧڢᴻ̶
֦ժݎሿဌ҅౯Ӥᶎےᲁᥴᲁ᮷አԧUUID҅ᬯ੪ฎԅԧכᦤ҅᧡ےᲁԧ᧡ᥴᲁ҅ᥝฎ֦ڢധԧ౯ጱ
ᲁ҅ᮎӧԤॺԧࡶ̶
LUAฎܻৼጱ҅Ԟྲᓌܔ҅੪ฎڣෙӞӥKey౯ժ݇හฎވፘᒵ҅ฎጱᦾ੪ڢᴻ҅ᬬࢧ౮ۑ1҅0
੪ฎ०ᨳ̶
ḵᦤ
౯ժݢզአ౯ժٟጱRedisᲁᦶᦶපຎ҅ݢզ፡ک᮷ೲᆙᶲଧ݄ಗᤈԧ
}
/**
*
*/
public boolean unlock(String id) {
String script =
"if redis.call('get',KEYS[1]) == ARGV[1] then" +
" return redis.call('del',KEYS[1]) " +
"else" +
" return 0 " +
"end";
try {
String result = jedis.eval(script,
Collections.singletonList(LOCK_KEY),
Collections.singletonList(id)).toString();
return "1".equals(result) ? true : false;
} finally {
jedis.close();
}
}
ᘍ
य़ਹฎӧฎᥧਠᗦԧ҅֕ฎӤᶎጱᲁ҅ํӧቭዟጱ҅౯ဌᘍஉग़ᅩ֦҅ᦜݢզᘍӞӥ҅რᎱ
౯᮷რک౯ጱGItHubԧ̶
ᘒӬ҅ᲁӞᛱ᮷ฎᵱᥝݢ᯿فᤈጱ҅Ӥᶎጱᕚᑕ᮷ฎಗᤈਠԧ੪᯽නԧ҅෫ဩེٚᬰفԧ҅ᬰ݄Ԟฎ᯿
ෛےᲁԧ҅ԭӞӻᲁጱᦡᦇ᧔ᙗਧӧฎஉݳቘጱ̶
౯ӧᓒಋٟ҅ࢩԅ᮷ํሿ౮ጱ҅ڦՈଆ౯ժٟঅԧ̶
redisson
redissonጱᲁ҅੪ਫሿԧݢ᯿فԧ҅֕ฎ՜ጱრᎱྲฤ႑ᵙ̶
ֵአ᩸உᓌܔ҅ࢩԅ՜ժବ੶᮷ᤰঅԧ֦҅ᬳളӤ֦ጱRedisਮಁᒒ҅՜ଆ֦؉ԧ౯ӤᶎٟጱӞ
ڔ҅ᆐݸๅਠᗦ̶
ᓌܔ፡፡՜ጱֵአމ҅᪙ྋଉֵአLockဌࠨ܄ڦ̶
ThreadPoolExecutor threadPoolExecutor =
new ThreadPoolExecutor(inventory, inventory, 10L, SECONDS,
linkedBlockingQueue);
long start = System.currentTimeMillis();
Config config = new Config();
config.useSingleServer().setAddress("redis://127.0.0.1:6379");
final RedissonClient client = Redisson.create(config);
Ӥᶎݢզ፡ک౯አکԧgetLockٌ҅ਫ੪ฎ឴ݐӞӻᲁጱਫ̶ֺ
RedissionLock Ԟဌ؉ࠨ҅੪ฎᆧఀጱڡত۸̶
ےᲁ
ํဌํݎሿஉग़᪙Lockஉग़ፘ֒ጱࣈොޫҘ
ᦶےᲁ҅೭ک୮ڹᕚᑕ҅ᆐݸ౯१᧔ጱttlԞ፡کԧ҅ฎӧฎӞڔ᮷ฎᮎԍᆧఀҘ
final RLock lock = client.getLock("lock1");
for (int i = 0; i <= NUM; i++) {
threadPoolExecutor.execute(new Runnable() {
public void run() {
lock.lock();
inventory--;
System.out.println(inventory);
lock.unlock();
}
});
}
long end = System.currentTimeMillis();
System.out.println("ಗᤈᕚᑕහ:" + NUM + " ᘙ:" + (end - start) + " ପ
ਂහԅ:" + inventory);
public RLock getLock(String name) {
return new RedissonLock(connectionManager.getCommandExecutor(), name);
}
public RedissonLock(CommandAsyncExecutor commandExecutor, String name) {
super(commandExecutor, name);
//եಗᤈ
this.commandExecutor = commandExecutor;
//UUIDਁᒧԀ
this.id = commandExecutor.getConnectionManager().getId();
//ٖ᮱ᲁᬦ๗ᳵ
this.internalLockLeaseTime = commandExecutor.
getConnectionManager().getCfg().getLockWatchdogTimeout();
this.entryName = id + ":" + name;
}
឴ݐᲁ
឴ݐᲁጱײ҅Ԟྲᓌܔ֦҅ݢզ፡ک҅՜Ԟฎӧෙڬෛᬦ๗ᳵ҅᪙౯Ӥᶎӧෙ݄೭୮ڹᳵ໊҅
ḵᬦ๗ฎӞӻ᭲ቘ҅ݝฎ౯ྲᔋᔧ̶
public void lockInterruptibly(long leaseTime, TimeUnit unit) throws
InterruptedException {
//୮ڹᕚᑕID
long threadId = Thread.currentThread().getId();
//ᦶ឴ݐᲁ
Long ttl = tryAcquire(leaseTime, unit, threadId);
// ইຎttlԅᑮ҅ڞᦤก឴ݐᲁ౮ۑ
if (ttl == null) {
return;
}
//ইຎ឴ݐᲁ०ᨳ҅ڞᦈᴅکଫᬯӻᲁጱchannel
RFuture
commandExecutor.syncSubscription(future);
try {
while (true) {
//ེٚᦶ឴ݐᲁ
ttl = tryAcquire(leaseTime, unit, threadId);
//ttlԅᑮ҅᧔ก౮ۑ឴ݐᲁ҅ᬬࢧ
if (ttl == null) {
break;
}
//ttlय़ԭ0 ڞᒵஇttlᳵݸᖀᖅᦶ឴ݐ
if (ttl >= 0) {
getEntry(threadId).getLatch().tryAcquire(ttl,
TimeUnit.MILLISECONDS);
} else {
getEntry(threadId).getLatch().acquire();
}
}
} finally {
//ݐၾchannelጱᦈᴅ
unsubscribe(future, threadId);
}
//get(lockAsync(leaseTime, unit));
}
ବ੶ےᲁ᭦ᬋ
֦ݢᚆտమᬯԍग़֢҅ࣁӞ᩸ӧฎܻৼӧᬮฎํᳯ᷌ԍҘ
य़֬ժᙗਧమکޚ҅ಅզᬮฎLUA҅՜ֵአԧHashጱහഝᕮ̶
Ԇᥝฎڣෙᲁฎވਂࣁ҅ਂࣁ੪ᦡᗝᬦ๗ᳵ҅ইຎᲁ૪ᕪਂࣁԧ҅ᮎྲӞӥᕚᑕ҅ᕚᑕฎӞӻᮎ੪
ᦤกݢզ᯿ف҅ᲁࣁԧ҅֕ฎӧฎ୮ڹᕚᑕ҅ᦤกڦՈᬮဌ᯽න҅ᮎ੪ۃ֟ᳵᬬࢧ҅ےᲁ०ᨳ̶
ฎӧฎํᅩᕰ҅ग़ቘᥴӞ̶᭭
private
final long threadId) {
//ইຎଃํᬦ๗ᳵ҅ڞೲᆙฦ᭗ොୗ឴ݐᲁ
if (leaseTime != -1) {
return tryLockInnerAsync(leaseTime, unit, threadId,
RedisCommands.EVAL_LONG);
}
//ضೲᆙ30ᑁጱᬦ๗ᳵಗᤈ឴ݐᲁጱොဩ
RFuture
commandExecutor.getConnectionManager().getCfg().getLockWatchdogTimeout(),
TimeUnit.MILLISECONDS, threadId, RedisCommands.EVAL_LONG);
//ইຎᬮ೮ํᬯӻᲁ҅ڞސਧձۓӧෙڬෛᧆᲁጱᬦ๗ᳵ
ttlRemainingFuture.addListener(new FutureListener
@Override
public void operationComplete(Future
{
if (!future.isSuccess()) {
return;
}
Long ttlRemaining = future.getNow();
// lock acquired
if (ttlRemaining == null) {
scheduleExpirationRenewal(threadId);
}
}
});
return ttlRemainingFuture;
}
ᥴᲁ
ᲁጱ᯽නԆᥝฎpublish᯽නᲁጱמ௳҅ᆐݸ؉໊ḵ҅Ӟտڣෙฎވ୮ڹᕚᑕ҅౮ۑ੪᯽නᲁ҅ᬮํ
ӻhincrby᭓ٺጱ֢҅ᲁጱ꧊य़ԭ0᧔กฎݢ᯿فᲁ҅ᮎ੪ڬෛᬦ๗ᳵ̶
ইຎ꧊ੜԭ0ԧ҅ᮎڢധKey᯽නᲁ̶
ฎӧฎ݈AQSஉ؟ԧҘ
AQS੪ฎ᭗ᬦӞӻvolatileץ᷶status݄፡ᲁጱᇫா҅Ԟտ፡හ꧊ڣෙฎވฎݢ᯿فጱ̶
ಅզ౯᧔դᎱጱᦡᦇ๋҅ݸ੪ӡڻ୭Ӟ҅᮷ฎӞጱ̶
long threadId, RedisStrictCommand
//ᬦ๗ᳵ
internalLockLeaseTime = unit.toMillis(leaseTime);
return commandExecutor.evalWriteAsync(getName(),
LongCodec.INSTANCE, command,
//ইຎᲁӧਂࣁ҅ڞ᭗ᬦhsetᦡᗝਙጱ꧊҅ଚᦡᗝᬦ๗ᳵ
"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; " +
//ইຎᲁ૪ਂࣁ҅ଚӬᲁጱฎ୮ڹᕚᑕ҅ڞ᭗ᬦhincrbyᕳහ꧊᭓ी1
"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.