You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Java多线程读取MySQL表如何实现记录独占读取的锁机制问题

集群环境下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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 10:33:11