目录
在讲分布式锁之前,我们先来看一个业务案例场景,需求描述如下:
我们需要开发一个接口,该接口接收ASR语音识别模块的识别结果信息,将一通通话的对话信息按照顺序保存起来。
接口入参的JSON格式如下:
- {
- "params":{
- "extensionNumber":"1001",
- "txt":"您好,请问需要办理信用卡吗?"
- "type":"0",
- "sessionId":"00030001",
- "callid":"3000000001",
- "callingType":"1",
- "startTime":"2020/05/11 16:30:30",
- "tel":"13839930146",
- "provId":"0311",
- "dataTime":"2020/05/11 16:30:40"
- }
- }
其中我们需要重点关注一下几个字段:
callid是通话记录id,同一个callid会对应多个sessionid。
callingType是通话状态,1是通话开始,2是通话中,3是通话结束
dataTime是该消息被组装发送的时间,这个字段可以作为会话的排序字段
这里有几个细节点:
(1)一通通话大多数情况下都是多轮对话
(2)如何保证对话的有序性
(3)高并发场景处理
首先因为话务量非常高,所以该接口一定会面临非常高的并发请求,其次,该接口同一时间要处理的会话信息一定是多个,再次,可能因为网络原因,出现先发后置或者后发先至的情况,也就是对话信息是无序的,所以接口内部逻辑,我们肯定不能直接写入数据库,或者通过数据库进行数据排序再更新,因为高并发的缘故,数据库大概率会被打满,然后接口就拒绝服务了。
针对以上问题,我们进行如下设计:
(1)使用redis做数据缓存,每次请求过来,我们根据callid到redis中获取对话列表,然后将本次入参信息更新到对话列表,然后再回写redis
(2)引入“水印”概念,设置水印留存时长为WaterTime,当一通对话的callingType=3的时候,记录当前时间为NowTime,后续如果再来该通话的会话记录,判断dataTime是是否小于NowTime+WaterTIme,如果小于则重复第一步操作,否则就丢弃消息。
(3)新增一个独立的线程定时轮训redis,根据callid查询对象,判断对话列表对象是否静止(当前时间>NowTime+WaterTime*2),如果数据静止就将数据取出根据dataTIme进行排序,然后入库到mysql,同时从redis中删除数据。
(4)将应用横向扩展,部署多个节点,增加节点的并发处理能力
上面的方案基本可以解决前面提到的三个问题,高并发问题、无序问题、多轮对话合并。但是这里面有一个潜在的隐藏问题,那就是多节点分布式部署的时候,一定会出现数据的问题。
比如,我们部署了三个节点,因为网络问题缘故,callid为100001,但是sessionId不同的请求,同一时间分别请求到了A/B/C三个节点,按照上面的逻辑,有可能有会丢失A/B的数据,请求同时到到达接口,内部逻辑先根据callId取出对象,然后将本次信息附上去,在存入redis。完全存在B/C节点上处理的时候,A还没存入的情况。这样数据就丢了。
这时候我们就需要引入分布式锁来解决这个问题。
分布式锁,顾名思义,就是解决分布式问题时候的锁。如果是单体应用,我们可以使用java的synchronize关键字。 如果是分布式的话,synchronize就不行了。
根据上面的安利,可以大概了解到,分布式锁要解决的问题,就是在分布式部署环境下,不同进程的不同线程在对相同资源进行请求的时候,需要考虑加锁。
另外:分布式锁和分布式事务是两码事,大多数时候,数据库事务是分布式锁实现的一种方式,比如我们的应用涉及到对mysql数据表的更新或者写入,因为mysql自身带有的事物,所以就可以避免因为高并发分布式操作的时候出现的问题,事物锁一定可以保证操作的唯一性。
分布式事务是另外一个概念,主要是微服务架构中,服务链式调用时候需要考虑的问题。这里不做过多记录。
参考Redis实现分布式锁这篇文档,写的非常详细,博文写的也很好,介绍了几种分布式锁的实现方案。
简单记录下我在开发过程中使用的到分布式锁。
- public class JedisHelper {
- private static ThreadLocal
-
- /**
- * 尝试获取分布式锁
- *
- * @param jedis Redis客户端
- * @param lockKey 锁
- * @param expireTime 超期时间
- * @return 是否成功获取分布式锁
- */
- public static boolean tryGetDistributedLock(Jedis jedis, String lockKey, int expireTime) {
- String requestId = buildRequestId();
- Object[] local = new Object[]{jedis, lockKey, requestId};
-
- final String LOCK_SUCCESS = "OK";
- final String SET_IF_NOT_EXIST = "NX";
- final String SET_WITH_EXPIRE_TIME = "PX";
- long maxMillis = expireTime + 5000;
- long begin = System.currentTimeMillis();
- while (true) {
- long now = System.currentTimeMillis();
- if ((now - begin) > maxMillis) {
- return false;
- }
-
- String result = jedis.set(lockKey, requestId, SET_IF_NOT_EXIST, SET_WITH_EXPIRE_TIME, expireTime);
- if (LOCK_SUCCESS.equals(result)) {
- threadLocal.set(local);
- return true;
- }
-
- sleep(10);
- }
- }
-
- /**
- * 释放分布式锁
- *
- * @param maxCloseMillis 关闭操作最大时长
- */
- public static void releaseDistributedLock(long maxCloseMillis) {
- Object[] local = threadLocal.get();
- if (local == null) {
- return;
- }
-
- threadLocal.remove();
- Jedis jedis = (Jedis) local[0];
- String lockKey = (String) local[1];
- String requestId = (String) local[2];
-
- try {
- long begin = System.currentTimeMillis();
- while (true) {
- long now = System.currentTimeMillis();
- if ((now - begin) > maxCloseMillis) {
- //关闭不成功时,放弃关闭,让其自动过期
- return;
- }
-
- String script = "if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('del', KEYS[1]) else return 0 end";
- Object result = jedis.eval(script, Collections.singletonList(lockKey), Collections.singletonList(requestId));
- final Long RELEASE_SUCCESS = 1L;
- if (RELEASE_SUCCESS.equals(result)) {
- return;
- }
-
- sleep(10);
- }
- } finally {
- if (jedis != null) {
- jedis.close();
- }
- }
- }
-
- /**
- * 构建一个唯一的requestId
- */
- private static String buildRequestId() {
- return new Date().getTime() + "-" + UUID.randomUUID().toString();
- }
-
- /**
- * 休眠指定时间
- */
- private static void sleep(long millis) {
- try {
- Thread.sleep(millis);
- } catch (Exception e) {
- }
- }
- try{
- //获取分布式锁,锁等待时间 单位毫秒,最多等待这么长时间,会自动销毁 (1分钟),作用:处理一通通话的时间,最多1分钟,超过一分钟会自动释放锁
- boolean flag = JedisHelper.tryGetDistributedLock(jedis,lockKey,60000);
- if(!flag){
- throw new RuntimeException("获取分布式锁失败!");
- }
- //下面进行业务逻辑处理即可
-
- }finally {
- //使用完一定要是释放掉锁
- JedisHelper.releaseDistributedLock(5000);
- }