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

如何在Apache Ignite中使用SELECT ... FOR UPDATE实现原子选改避免重复处理

Apache Ignite 实现类似 SELECT FOR UPDATE 的分布式锁与原子操作

针对你多Spring Boot客户端节点重复处理缓存记录的问题,以下是几种适配场景的实现方案:

1. 悲观锁机制(最贴近 SELECT FOR UPDATE)

Ignite支持分布式悲观锁,能直接锁定指定缓存条目,确保同一时间只有一个客户端能修改该记录。

实现方式:

  • 通过IgniteCache.lock(K key)显式获取锁,操作完成后手动释放(或利用try-with-resources自动释放)
  • 在SQL查询中使用FOR UPDATE子句,直接锁定查询到的记录

代码示例(SQL方式):

// 获取Ignite缓存
IgniteCache<Long, Record> cache = ignite.cache("recordCache");

// 执行带FOR UPDATE的SQL查询,锁定TO_BE_PROCESSED状态的记录
SqlFieldsQuery query = new SqlFieldsQuery(
    "SELECT id, status FROM Record WHERE status = ? FOR UPDATE")
    .setArgs("TO_BE_PROCESSED");

try (QueryCursor<List<?>> cursor = cache.query(query)) {
    for (List<?> row : cursor) {
        Long recordId = (Long) row.get(0);
        // 原子更新状态为IN_PROGRESS
        cache.invoke(recordId, (entry, args) -> {
            Record record = entry.getValue();
            if ("TO_BE_PROCESSED".equals(record.getStatus())) {
                record.setStatus("IN_PROGRESS");
                entry.setValue(record);
                return true;
            }
            return false;
        });
        
        // 执行业务逻辑处理
        processRecord(recordId);
    }
}

代码示例(显式锁方式):

Long recordId = getNextToProcessRecordId(); // 先查询待处理记录ID

try (IgniteLock lock = cache.lock(recordId)) {
    // 再次检查状态,避免锁等待期间被其他节点修改
    Record record = cache.get(recordId);
    if (record != null && "TO_BE_PROCESSED".equals(record.getStatus())) {
        record.setStatus("IN_PROGRESS");
        cache.put(recordId, record);
        
        // 执行业务处理
        processRecord(record);
    }
}

2. 乐观锁机制(基于版本号)

通过版本字段实现乐观锁,适合并发冲突频率较低的场景,性能比悲观锁更高。

实现步骤:

  1. 给缓存实体添加版本字段(如@Version注解):
public class Record {
    private Long id;
    private String status;
    @Version
    private long version; // 版本字段,Ignite自动维护
    // getter、setter
}
  1. 配置缓存启用乐观锁:
CacheConfiguration<Long, Record> cacheCfg = new CacheConfiguration<>("recordCache");
cacheCfg.setWriteSynchronizationMode(CacheWriteSynchronizationMode.FULL_SYNC);
cacheCfg.setIndexedTypes(Long.class, Record.class);
  1. 原子更新操作:
Record record = cache.get(recordId);
if (record != null && "TO_BE_PROCESSED".equals(record.getStatus())) {
    record.setStatus("IN_PROGRESS");
    try {
        cache.put(record); // 版本不匹配时会抛出OptimisticLockingException
        processRecord(record);
    } catch (OptimisticLockingException e) {
        // 处理版本冲突,说明该记录已被其他节点修改
        log.warn("Record {} has been modified by another node", recordId);
    }
}

3. 分布式原子操作(Invoke方法)

利用Ignite的invoke方法将查询与更新逻辑封装在一个分布式原子操作中,由Ignite保证全局原子性,无需手动处理锁。

代码示例:

// 批量获取并锁定待处理记录(限制每次处理数量)
SqlFieldsQuery fetchQuery = new SqlFieldsQuery(
    "SELECT id FROM Record WHERE status = ? LIMIT 10")
    .setArgs("TO_BE_PROCESSED");

try (QueryCursor<List<?>> cursor = cache.query(fetchQuery)) {
    for (List<?> row : cursor) {
        Long recordId = (Long) row.get(0);
        
        Boolean updated = cache.invoke(recordId, (entry, args) -> {
            Record record = entry.getValue();
            if (record != null && "TO_BE_PROCESSED".equals(record.getStatus())) {
                record.setStatus("IN_PROGRESS");
                entry.setValue(record);
                return true;
            }
            return false;
        });
        
        if (updated) {
            processRecord(recordId);
        }
    }
}

场景适配建议

  • 如果并发冲突频繁,推荐使用悲观锁,能直接避免冲突
  • 若冲突较少,优先选择乐观锁或原子invoke操作,性能更优
  • 定时任务中建议限制每次获取的记录数量,避免长时间持有锁影响并发

内容的提问来源于stack exchange,提问作者Junaid Akhtar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 07:01:12