如何在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. 乐观锁机制(基于版本号)
通过版本字段实现乐观锁,适合并发冲突频率较低的场景,性能比悲观锁更高。
实现步骤:
- 给缓存实体添加版本字段(如
@Version注解):
public class Record { private Long id; private String status; @Version private long version; // 版本字段,Ignite自动维护 // getter、setter }
- 配置缓存启用乐观锁:
CacheConfiguration<Long, Record> cacheCfg = new CacheConfiguration<>("recordCache"); cacheCfg.setWriteSynchronizationMode(CacheWriteSynchronizationMode.FULL_SYNC); cacheCfg.setIndexedTypes(Long.class, Record.class);
- 原子更新操作:
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
相关产品推荐
相关产品推荐

