Java多线程读取MySQL表如何实现记录独占读取的锁机制问题
我计划将应用迁移至集群环境,因此需要处理MySQL锁机制问题。现有一张包含id和flag字段(标记记录是否已更新)的测试表,通过3个线程模拟集群中的独立进程读取5条记录,期望每个线程仅读取未被其他线程读取的记录。
我尝试禁用autoCommit并使用带LOCK的SELECT语句,以下是线程读取的示例代码:
public Long read(Connection c) throws Exception { c.setAutoCommit(false); c.setTransactionIsolation(Connection.TRANSACTION_READ_COMMITTED); PreparedStatement sth = c.prepareStatement("SELECT * FROM test WHERE flag!='Y' LIMIT 1 FOR UPDATE SKIP LOCKED") ; ResultSet rst = sth.executeQuery() ; Long id = null ; if (rst.next()) { id = rst.getLong("id"); sth = c.prepareStatement("UPDATE test SET FLAG='Y' WHERE id=" + id) ; sth.executeUpdate(); System.out.println(getName() + ": Commited id " + id); } c.commit(); sth.close(); return id ; }
测试运行3个线程后输出如下:
LockThread#2: Commited id 1 LockThread#2: Commited id 2 LockThread#3: Commited id 1 LockThread#2: Commited id 3 LockThread#1: Commited id 1 LockThread#1: Commited id 4 LockThread#3: Commited id 4 LockThread#2: Commited id 4 LockThread#3: Commited id 5 LockThread#2: Commited id 5 LockThread#1: Commited id 5
实际运行后发现同一记录被多个线程重复读取,以下是几种可行的解决方法:
可行解决方案
1. 合并查询与更新为原子操作
将SELECT和UPDATE合并成一条SQL语句,消除事务中两个操作之间的时间窗口,确保锁定与更新的原子性:
UPDATE test SET flag='Y' WHERE id = (SELECT id FROM test WHERE flag!='Y' LIMIT 1 FOR UPDATE SKIP LOCKED)
也可以用JOIN方式实现,避免子查询可能的性能问题:
UPDATE test t1 JOIN (SELECT id FROM test WHERE flag!='Y' LIMIT 1 FOR UPDATE SKIP LOCKED) t2 ON t1.id = t2.id SET t1.flag='Y'
对应的Java代码可以直接执行这条UPDATE语句,通过getUpdateCount()判断是否有记录被更新,再按需获取目标id。
2. 去掉SKIP LOCKED,使用阻塞式悲观锁
如果业务允许线程等待锁释放而非跳过,可以移除SKIP LOCKED关键字,这样当一个线程锁定符合条件的行后,其他线程会阻塞直到锁释放,从根源避免重复读取:
PreparedStatement sth = c.prepareStatement("SELECT * FROM test WHERE flag!='Y' LIMIT 1 FOR UPDATE") ;
这种方式适合对并发吞吐量要求不高、但必须保证数据不重复的场景,缺点是高并发下可能出现线程排队等待的情况。
3. 乐观锁+版本号机制
放弃悲观锁,给表新增一个version整数字段,通过版本号判断记录是否被修改,实现无锁化的并发控制:
public Long read(Connection c) throws Exception { c.setAutoCommit(false); c.setTransactionIsolation(Connection.TRANSACTION_READ_COMMITTED); while (true) { // 读取未处理的记录及其版本号 PreparedStatement selectStmt = c.prepareStatement("SELECT id, version FROM test WHERE flag!='Y' LIMIT 1"); ResultSet rst = selectStmt.executeQuery(); Long id = null; Integer version = null; if (rst.next()) { id = rst.getLong("id"); version = rst.getInt("version"); } else { c.commit(); selectStmt.close(); return null; // 没有未处理记录 } selectStmt.close(); // 尝试更新,只有版本号匹配时才生效 PreparedStatement updateStmt = c.prepareStatement( "UPDATE test SET flag='Y', version=version+1 WHERE id=? AND version=?" ); updateStmt.setLong(1, id); updateStmt.setInt(2, version); int affectedRows = updateStmt.executeUpdate(); updateStmt.close(); if (affectedRows > 0) { System.out.println(getName() + ": Commited id " + id); c.commit(); return id; } // 更新失败,说明记录已被其他线程处理,回滚后重试 c.rollback(); } }
这种方式适合高并发场景,避免锁竞争,但需要处理重试逻辑。
4. 引入分布式锁(跨节点集群场景)
如果是跨多个应用节点的集群,除了数据库锁,还可以用Redis等实现分布式锁:在读取记录前,先尝试获取对应id的分布式锁(或全局的任务分配锁),获取成功后再执行数据库的读取和更新操作。需要注意设置合理的锁过期时间,避免因节点故障导致死锁。
现有代码的潜在问题
你的现有代码中存在两个可能导致重复读取的隐患:
- 重复赋值
PreparedStatement变量,导致初始的SELECT语句对应的Statement未被正确关闭,可能造成连接资源泄漏,间接影响锁的及时释放; - SELECT与UPDATE是两个独立操作,虽然有行锁,但在某些极端场景下(如事务提交延迟、数据库锁机制的边缘情况),可能出现其他线程读取到未更新flag的记录。
内容的提问来源于stack exchange,提问作者WonderMouse

