多JVM间同步实现方案咨询:K8s多Pod下线程安全问题
解决K8s多Pod环境下特定用户的分布式同步更新问题
你遇到的核心问题非常典型:synchronized(this)是JVM本地锁,只能在单个JVM进程内保证线程安全,但K8s的多Pod是完全独立的JVM实例,彼此之间没有锁共享,所以同一用户的请求被分发到不同Pod时,就会出现并发更新的问题。
要解决这个问题,必须引入分布式锁,而且锁的粒度要绑定到特定用户(比如用用户ID作为锁的key),这样既能保证同一用户的更新操作互斥,又不会影响其他用户的并发执行。下面给你两个可行的方案,附代码示例和关键注意事项:
方案一:基于数据库的分布式锁(适合小流量场景)
如果你的系统流量不大,不需要额外引入中间件,用现有数据库就能实现分布式锁。核心思路是利用数据库的唯一约束和行级锁来保证同一用户的锁只能被一个实例获取。
步骤1:创建锁表
先在数据库中创建一张锁表,用来存储用户的锁状态:
CREATE TABLE user_distributed_lock ( user_id VARCHAR(64) NOT NULL PRIMARY KEY COMMENT '用户ID,作为锁的唯一标识', lock_time BIGINT NOT NULL COMMENT '锁获取时间戳(毫秒)', expire_time BIGINT NOT NULL COMMENT '锁过期时间戳(毫秒)' ) COMMENT '用户分布式锁表';
步骤2:实现锁工具类
用Spring的JdbcTemplate封装锁的获取和释放逻辑:
@Component public class DbDistributedLock { @Autowired private JdbcTemplate jdbcTemplate; // 尝试获取锁,返回是否成功 public boolean tryLock(String userId, long expireMillis) { try { // 用INSERT ON DUPLICATE KEY UPDATE实现:如果锁不存在则插入,存在则更新过期时间 String sql = "INSERT INTO user_distributed_lock (user_id, lock_time, expire_time) " + "VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE expire_time = ?"; long now = System.currentTimeMillis(); long expireTime = now + expireMillis; int affectedRows = jdbcTemplate.update(sql, userId, now, expireTime, expireTime); return affectedRows > 0; } catch (Exception e) { // 捕获主键冲突等异常,返回获取锁失败 return false; } } // 释放锁 public void releaseLock(String userId) { String sql = "DELETE FROM user_distributed_lock WHERE user_id = ?"; jdbcTemplate.update(sql, userId); } // 定时清理过期锁(防止实例崩溃导致锁一直占用) @Scheduled(fixedRate = 300000) // 每5分钟执行一次 public void cleanExpiredLocks() { String sql = "DELETE FROM user_distributed_lock WHERE expire_time < ?"; jdbcTemplate.update(sql, System.currentTimeMillis()); } }
步骤3:改造业务代码
在业务逻辑中加入分布式锁,并且必须重新查询最新的person对象(多实例环境下本地对象可能不是最新的):
@Autowired private DbDistributedLock dbDistributedLock; @Autowired private PersonRepository personRepository; public void updateUserAmount(String userId, Person person, Sal sal, int refCount) { String lockKey = userId; long lockExpireMillis = 300000; // 锁超时5分钟,根据业务调整 try { // 最多重试3次,避免瞬时并发冲突 boolean locked = false; int retryCount = 0; while (!locked && retryCount < 3) { locked = dbDistributedLock.tryLock(lockKey, lockExpireMillis); if (!locked) { Thread.sleep(100); // 重试间隔100ms retryCount++; } } if (!locked) { throw new RuntimeException("当前操作正在处理中,请稍后重试"); } // 关键:重新从数据库查询最新的person对象,避免本地缓存的旧数据导致不一致 Person latestPerson = personRepository.findById(person.getId()) .orElseThrow(() -> new RuntimeException("用户信息不存在")); int calcAmount = latestPerson.getSum() + sal.getMonthlySal(); calcAmount = calcAmount + refCount; int status = customQuery.updateUnit(calcAmount); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("操作被中断", e); } finally { // 无论业务成功还是失败,都要释放锁 dbDistributedLock.releaseLock(lockKey); } }
方案二:基于Redis的分布式锁(适合高流量场景)
如果你的系统并发量较高,推荐用Redis实现分布式锁,性能更好。这里推荐用Redisson框架,它已经封装了锁的自动续期、可重入、公平锁等细节,不用自己造轮子。
步骤1:引入Redisson依赖
在Maven中加入Redisson的Spring Boot starter:
<dependency> <groupId>org.redisson</groupId> <artifactId>redisson-spring-boot-starter</artifactId> <version>3.23.3</version> <!-- 使用最新稳定版 --> </dependency>
步骤2:配置Redis连接
在application.yml中配置Redis的连接信息:
spring: redis: host: your-redis-host port: 6379 password: your-redis-password # 如果有密码的话 database: 0
步骤3:改造业务代码
用Redisson的RLock实现分布式锁:
@Autowired private RedissonClient redissonClient; @Autowired private PersonRepository personRepository; public void updateUserAmount(String userId, Person person, Sal sal, int refCount) { // 锁的key要加上前缀,避免和其他业务的锁冲突 RLock lock = redissonClient.getLock("user_update_lock:" + userId); try { // 尝试获取锁:最多等待3秒,锁自动过期5分钟 // Redisson的看门狗机制会自动续期锁,避免业务执行时间过长导致锁失效 boolean locked = lock.tryLock(3, 5, TimeUnit.MINUTES); if (!locked) { throw new RuntimeException("当前操作正在进行,请稍后重试"); } // 同样要重新查询最新的person对象 Person latestPerson = personRepository.findById(person.getId()) .orElseThrow(() -> new RuntimeException("用户信息不存在")); int calcAmount = latestPerson.getSum() + sal.getMonthlySal(); calcAmount = calcAmount + refCount; int status = customQuery.updateUnit(calcAmount); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("操作被中断", e); } finally { // 只有当前持有锁的线程才能释放锁,避免误释放其他线程的锁 if (lock.isHeldByCurrentThread()) { lock.unlock(); } } }
关键注意事项
- 锁粒度必须绑定用户:绝对不能用全局锁(比如固定的一个key),否则会导致所有用户的操作都互斥,严重影响系统性能。
- 必须重新查询最新数据:多实例环境下,本地的
person对象可能已经被其他Pod修改,直接使用会导致数据不一致,一定要从数据库重新查询。 - 锁的超时和释放要可靠:用
finally块保证锁被释放,或者用自动过期机制(比如Redisson的看门狗)避免死锁。 - 重试机制:获取锁失败时可以重试几次,或者返回友好提示给用户,提升用户体验。
内容的提问来源于stack exchange,提问作者djyo02
相关产品推荐
相关产品推荐

