分布式锁-心跳模式

请添加图片描述

什么是分布式锁?

​ 分布式锁是一种在分布式系统中实现同步和控制的机制。在多个进程或者节点之间共享数据或资源时,为了避免数据竞争等问题出现,需要对其进行同步和控制。分布式锁是一种实现这种同步和控制的机制。

​ 分布式锁的工作原理是通过在系统中创建共享锁,保证同一时间只有一个进程或者节点可以对锁定资源进行访问或更新,从而实现同步和控制。

​ 分布式锁可以实现在分布式系统中对全局共享资源或者一部分资源进行访问控制和同步。如果没有使用分布式锁,多个进程或者节点同时对同一个资源进行访问或修改是非常危险的,可能引发各种竞态条件。因此,分布式锁在分布式系统中起到了非常重要的作用。

分布式锁的实现方式

Java中的分布式锁可以通过以下几种方式实现:

  1. 基于数据库实现分布式锁:

    ​ 将锁信息存储在共享的数据库中,在需要加锁的时候先尝试插入一条记录,若插入成功,则成功获取到锁,否则已经有其他进程持有锁,需要等待该锁被释放后再尝试获取。

  2. 基于缓存实现分布式锁:

    ​ 将锁信息存储在分布式缓存中,如Redis,同时利用Redis的原子性操作实现加锁和解锁,即通过SETNX命令尝试将键值对设置成锁状态,若设置成功则成功获取锁,否则已有其他应用持有锁。

  3. 基于Zookeeper实现分布式锁:

    ​ 通过Zookeeper的临时节点机制实现分布式锁。当需要加锁时,创建具有唯一名称和标识的临时znode节点,如果创建成功则获取锁,如果创建失败则表示锁已经被其他进程持有,需要等待该锁被释放后再尝试获取。

	需要注意的是,在分布式系统中使用锁需要考虑一些特殊情况,如节点故障、并发冲突、死锁等。因此,在实现分布式锁时需要考虑应用场景,根据需求选择适当的实现方式,并解决潜在的问题。

通过Reids来实现分布式锁

心跳模式

​ Redis分布式锁的心跳模式是一种保证锁有效性的机制。在并发环境下,锁的有效性可能会受到网络闪断等因素的影响,导致锁不被及时释放,从而造成锁的失效。Redis分布式锁的心跳模式通过定期发送心跳信息,保持锁的有效性,从而防止锁的失效。

​ 具体来说,当客户端持有正在使用的锁时,会启动一个定时器,每隔一段时间发送一个心跳请求更新锁的过期时间,确保锁一直保持有效。如果控制锁的线程崩溃或者断电,由于超时时间已经设定,Redis会在到期时间之后自动将锁释放掉。

​ 在Redis中,使用setex命令创建一个锁,并保证设置过期时间的同时,才能创建成功。获取锁使用setnx命令,采用自旋锁的机制,直到获取到锁或者超时才退出等待。释放锁的时候,使用Lua脚本来保证原子性。

​ 心跳模式的应用使Redis分布式锁更加健壮和实用,避免了锁失效的问题,确保了分布式锁的正确性。但是心跳模式会带来更多的网络开销和额外的代码实现。

具体代码实现看以下例子:

配置类:


import jakarta.annotation.PostConstruct;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;

@ConditionalOnProperty(
    name = {"distributed-lock-conf.enabled"},
    havingValue = "true"
)
@Component
@ConfigurationProperties(
    prefix = "distributed-lock-conf"
)
public class DistributedLockConf {
    private static final Logger logger = LoggerFactory.getLogger(DistributedLockConf.class);
    private Boolean enabled = false;
    private Integer tryTimes = 3;
    private Boolean throwException = true;
    private Integer secondsToExpire = 10;
    private RedisConf redisConf = new RedisConf();
    private Integer heartbeatSeconds = 2;

    public DistributedLockConf() {
    }

    public Boolean getEnabled() {
        return this.enabled;
    }

    public DistributedLockConf setEnabled(Boolean enabled) {
        this.enabled = enabled;
        return this;
    }

    public Integer getTryTimes() {
        return this.tryTimes;
    }

    public DistributedLockConf setTryTimes(Integer tryTimes) {
        this.tryTimes = tryTimes;
        return this;
    }

    public Boolean getThrowException() {
        return this.throwException;
    }

    public DistributedLockConf setThrowException(Boolean throwException) {
        this.throwException = throwException;
        return this;
    }

    public RedisConf getRedisConf() {
        return this.redisConf;
    }

    public DistributedLockConf setRedisConf(RedisConf redisConf) {
        this.redisConf = redisConf;
        return this;
    }

    public Integer getSecondsToExpire() {
        return this.secondsToExpire;
    }

    public DistributedLockConf setSecondsToExpire(Integer secondsToExpire) {
        this.secondsToExpire = secondsToExpire;
        return this;
    }

    public Integer getHeartbeatSeconds() {
        return this.heartbeatSeconds;
    }

    public DistributedLockConf setHeartbeatSeconds(Integer heartbeatSeconds) {
        this.heartbeatSeconds = heartbeatSeconds;
        return this;
    }

    @PostConstruct
    private void init() {
        logger.info("DistributedLockConf init ...");
    }
}

分布式锁生成类:

其中有一些工具类不在进行逐个引入了

import com.google.common.util.concurrent.ThreadFactoryBuilder;
import jakarta.annotation.PostConstruct;
import java.util.Iterator;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import org.apache.commons.lang3.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Component;

@ConditionalOnProperty(
    name = {"distributed-lock-conf.enabled"},
    havingValue = "true"
)
@Component
public class DistributedLockBuilder {
    private static final Logger logger = LoggerFactory.getLogger(DistributedLockBuilder.class);
    private static DistributedLockConf distributedLockConf;
    private static RedisClient redisClient;
    private static final ConcurrentHashMap<String, DistributedLock> allLock = new ConcurrentHashMap();
    private ScheduledExecutorService heartbeatExecutor;
    private ExecutorService heartbeatDoExecutor;
    private ExecutorService releaseDoExecutor;

    public DistributedLockBuilder(DistributedLockConf distributedLockConf) {
        DistributedLockBuilder.distributedLockConf = distributedLockConf;
    }
/**
用于初始化分布式锁的相关配置和线程池。具体步骤如下: 
 1. 使用@PostConstruct注解表示该方法会在类初始化完成后自动调用。 
2. 从配置文件中获取Redis的相关配置信息,如果配置了jedisFromBean,则通过SpringContextUtil.getBean方法获取对应的RedisClient实例,否则新建一个RedisClient实例。 
3. 初始化心跳线程池,使用Executors.newSingleThreadScheduledExecutor()方法创建一个单线程池,该线程池会定时执行心跳方法heartbeat()。 
4. 初始化心跳执行线程池和释放锁执行线程池,使用Executors.newFixedThreadPool()方法创建固定大小的线程池,用于执行心跳和释放锁的任务。线程池的大小为10,使用ThreadFactoryBuilder()方法设置线程池的名称格式。
*/
    @PostConstruct
    private void init() {
        String jedisFromBean = distributedLockConf.getRedisConf().getJedisFromBean();
        if (StringUtils.isNotBlank(jedisFromBean)) {
            redisClient = (RedisClient)SpringContextUtil.getBean(jedisFromBean);
        } else {
            redisClient = new RedisClient(distributedLockConf.getRedisConf());
        }

        this.heartbeatExecutor = Executors.newSingleThreadScheduledExecutor();
        this.heartbeatExecutor.scheduleWithFixedDelay(() -> {
            try {
                this.heartbeat();
            } catch (Exception var2) {
                logger.error(var2.toString(), var2);
            }

        }, 0L, 1L, TimeUnit.SECONDS);
        this.heartbeatDoExecutor = Executors.newFixedThreadPool(10, (new ThreadFactoryBuilder()).setNameFormat("distributed-lock-hb-%d").build());
        this.releaseDoExecutor = Executors.newFixedThreadPool(10, (new ThreadFactoryBuilder()).setNameFormat("distributed-lock-release-%d").build());
    }

    public static DistributedLock lock(String redisKey, String owner) {
        return lock(redisKey, owner, Boolean.TRUE.equals(distributedLockConf.getThrowException()));
    }

/**
首先,通过redisClient.setWithExpireIfNotExist()方法尝试获取锁,如果获取不到则根据ifThrow参数决定是否抛出异常或返回null。如果获取到锁,则创建一个DistributedLock对象,设置其心跳时间、下一次心跳时间以及redisClient等属性,并将其存储在allLock中,最后返回DistributedLock对象。 
 具体步骤如下: 
1. 使用redisClient.setWithExpireIfNotExist()方法尝试获取锁,如果获取不到则根据ifThrow参数决定是否抛出异常或返回null。 
2. 如果获取到锁,则创建一个DistributedLock对象,设置其心跳时间、下一次心跳时间以及redisClient等属性。 
3. 将DistributedLock对象存储在allLock中。 
4. 返回DistributedLock对象。
*/
    public static DistributedLock lock(String redisKey, String owner, boolean ifThrow) {
        boolean ifGot = redisClient.setWithExpireIfNotExist(redisKey, owner, distributedLockConf.getSecondsToExpire());
        if (!ifGot) {
            if (ifThrow) {
                throw new ErrorCodeException(1, "failed to lock.");
            } else {
                return null;
            }
        } else {
            logger.debug("get lock: {}", redisKey);
            DistributedLock distributedLock = new DistributedLock(redisKey, owner, distributedLockConf.getSecondsToExpire(), distributedLockConf.getHeartbeatSeconds(), redisClient);
            distributedLock.nextHeartbeatTs = System.currentTimeMillis() + (long)(distributedLockConf.getHeartbeatSeconds() * 1000);
            allLock.put(distributedLock.sn, distributedLock);
            return distributedLock;
        }
    }

/**
这段代码是一个心跳函数,用于维护分布式锁的状态。首先,通过迭代器遍历所有的锁。
如果锁已经关闭,则从allLock中删除该锁。如果锁应该释放,则在releaseDoExecutor中执行释放操作,
并从allLock中删除该锁。
如果当前时间超过了锁的下一个心跳时间,则在heartbeatDoExecutor中执行心跳操作,
并更新锁的下一个心跳时间。
*/
    private void heartbeat() {
        Iterator<Map.Entry<String, DistributedLock>> iterator = allLock.entrySet().iterator();

        while(iterator.hasNext()) {
            Map.Entry<String, DistributedLock> entry = (Map.Entry)iterator.next();
            DistributedLock lock = (DistributedLock)entry.getValue();
            if (lock.ifHasClosed) {
                iterator.remove();
            } else if (lock.ifShouldRelease) {
                this.releaseDoExecutor.execute(() -> {
                    try {
                        redisClient.compareAndDelete(lock.redisKey, lock.owner);
                    } catch (Exception var2) {
                        logger.error(var2.toString(), var2);
                    }

                });
                iterator.remove();
            } else if (System.currentTimeMillis() >= lock.nextHeartbeatTs) {
                this.heartbeatDoExecutor.execute(() -> {
                    try {
                        redisClient.compareAndExpire(lock.redisKey, lock.owner, lock.secondsToExpire);
                        lock.nextHeartbeatTs = System.currentTimeMillis() + (long)lock.heartbeatSeconds * 1000L;
                    } catch (Exception var2) {
                        logger.error(var2.toString(), var2);
                    }

                });
            }
        }

    }
}
/**
这段代码是一个用于设置Redis中某个键值对的方法,如果该键不存在,则设置该键的过期时间为指定秒数。代码的具体步骤如下: 
 1. 初始化重试次数为2。 
 2. 进入一个无限循环。 
 3. 尝试从Redis连接池中获取一个Jedis连接。 
 4. 使用Jedis连接执行set操作,并设置nx和ex参数,表示只有在键不存在时才能设置该键的值,并设置该键的过期时间为指定秒数。 
 5. 如果set操作返回结果为"OK",则表示设置成功,返回true。 
 6. 如果set操作抛出了异常,则关闭Jedis连接,并将异常抛出。 
 7. 如果Jedis连接成功创建,则关闭该连接。 
 8. 如果set操作返回结果不为"OK",则表示设置失败,返回false。 
 9. 如果set操作抛出异常或者返回结果不为"OK",则记录错误日志,并将重试次数减1。 
 10. 如果重试次数小于等于0,则表示设置失败,返回false。
*/

public class RedisClient { 
    
    private final RedisFactory redisFactory;
    
	public boolean setWithExpireIfNotExist(String key, String value, int secondsToExpire) {
        int retryTimes = 2;
        while(true) {
            try {
                Jedis jedis = this.redisFactory.getJedis();
                boolean var7;
                try {
                    String result = jedis.set(key, value, SetParams.setParams().nx().ex(secondsToExpire));
                    var7 = "OK".equalsIgnoreCase(result);
                } catch (Throwable var9) {
                    if (jedis != null) {
                        try {
                            jedis.close();
                        } catch (Throwable var8) {
                            var9.addSuppressed(var8);
                        }
                    }
                    throw var9;
                }
                if (jedis != null) {
                    jedis.close();
                }
                return var7;
            } catch (Exception var10) {
                if (retryTimes-- <= 0) {
                    return false;
                }
            }
        }
    }

实体类:

public class DistributedLock implements Closeable {
    private static final Logger logger = LoggerFactory.getLogger(DistributedLock.class);
    public final String sn = UUID.randomUUID().toString();
    public final String redisKey;
    public final String owner;
    public final int secondsToExpire;
    public final int heartbeatSeconds;
    public final RedisClient redisClient;
    public long nextHeartbeatTs;
    public boolean ifShouldRelease = false;
    public boolean ifHasClosed = false;

    public DistributedLock(String redisKey, String owner, int secondsToExpire, int heartbeatSeconds, RedisClient redisClient) {
        this.redisKey = redisKey;
        this.owner = owner;
        this.secondsToExpire = secondsToExpire;
        this.heartbeatSeconds = heartbeatSeconds;
        this.redisClient = redisClient;
    }

    public void close() {
        try {
            this.ifHasClosed = this.redisClient.compareAndDelete(this.redisKey, this.owner);
        } catch (Exception var2) {
            logger.error(var2.toString(), var2);
        }

        logger.debug("release lock: {}", this.redisKey);
        this.ifShouldRelease = true;
    }
}

实际中应用

/**
采取增强try的方式进行
*/
try (DistributedLock distributedLock = DistributedLockBuilder.lock(redisKey, owner, ifThrow)) {
 	//通过判断distributedLock是否为空,依次来判断此线程是否被占用,
 	//为空表示未获得当前资源,跳过即可
 }

增强try备注:

在Java 7中引入的try-with-resources是增强的try结构。与传统的try-catch结构相比,
try-with-resources结构能在处理异常的过程中优雅地执行资源的释放,
避免了代码中显式地进行资源释放的繁琐和容易出错的操作。

try-with-resources结构可以自动关闭实现了Java 7中AutoCloseable接口或其父接口Closeable的资源,
例如文件流、数据库连接、网络连接等资源。
在try块中初始化资源对象,当try块结束时,无论是否发生异常,都会自动关闭资源。
如果try块和catch块都会抛出异常,那么先关闭资源,然后再抛出异常,保证了代码的一致性和可靠性。

try-with-resources的基本语法如下:

​```
try (resource) {
   // 代码块
} catch (exceptionType exception) {
   // 处理异常的代码
}
​```

其中resource是需要管理的资源,可以是一个或者多个以分号分隔的资源对象。
在try块执行完毕后会自动关闭这些资源。

需要注意的是,在使用try-with-resources结构时,资源对象必须在try块之前声明,且初始化必须在try块之内。
在声明资源对象时,也可以使用final关键字显式声明对象是不可变的,从而增加代码的可读性和安全性。

总之,try-with-resources是在Java 7中引入的增强型try语句块结构。
它可以更优雅地管理资源、避免代码重复、防止错误并发生代码问题。

更多推荐