Socket.io缓存并发更新竞态问题求助:查询结果不符预期
Socket.io缓存竞态条件解决方案
问题根源
你的代码中,并发的transaction请求会同时读取cache[data.key]的当前值,各自修改后再写回缓存,导致后续请求覆盖前面的修改结果——这是典型的读取-修改-写入竞态问题。之前尝试的mutex/lock无效,大概率是因为没有针对缓存的每个key单独加锁,而是用了全局锁或者锁的范围/时机不正确。
最优解决方案:按key串行化操作
针对每个缓存key维护独立的异步锁队列,确保同一key的所有操作(包括get和transaction)串行执行,不同key的操作互不阻塞,既解决竞态问题又保证性能。
手动实现锁队列(无额外依赖)
const { Server } = require("socket.io"); const io = new Server({}); var cache = {}; // 每个key对应一个锁Promise链,确保同一key的操作串行 const keyLocks = {}; // 获取key的锁,返回释放锁的函数 async function acquireLock(key) { let currentLock = keyLocks[key] || Promise.resolve(); let resolveLock; const newLock = new Promise(resolve => resolveLock = resolve); keyLocks[key] = currentLock.then(() => newLock); await currentLock; return resolveLock; } io.on('connection', (socket) => { socket.on('clear', () => { cache = {}; Object.keys(keyLocks).forEach(key => delete keyLocks[key]); }); socket.on('addToQueue', async (data, cb) => { const { key } = data; const releaseLock = await acquireLock(key); try { if (data.type === `get`) { await new Promise(resolve => setTimeout(resolve, 20)); cb({ value: cache[key] }); } if (data.type === `transaction`) { // 处理key不存在的默认值 let { balance, last_transactions } = cache[key] || { balance: { total: 0, limit: 0 }, last_transactions: [] }; const { value, type, description } = data.value; // 修复原代码拼写错误:balace→balance,tipo→type if (type === 'd') { if (balance.limit < ((balance.total - value) * -1)) { return cb({ value: { error: { type: type, error: `no-limit` }, ...cache[key] } }); } balance.total -= value; } if (type === 'c') { balance.total += value; } let latest_transactions = last_transactions || []; latest_transactions.push({ value, type, description, created_at: new Date().toISOString(), }); latest_transactions.sort((a, b) => new Date(b.created_at) - new Date(a.created_at)); latest_transactions = latest_transactions.slice(0, 10); cache[key] = { balance, last_transactions: latest_transactions }; cb({ value: cache[key] }); } } catch (err) { cb({ error: err.message }); } finally { // 必须释放锁,避免阻塞后续操作 releaseLock(); // 清理已完成的锁,防止内存泄漏 if (keyLocks[key] === keyLocks[key].then(() => {})) { delete keyLocks[key]; } } }); }); (async () => { const PORT = 6456; io.listen(PORT); console.log(`Server listening at *:${PORT}`); })();
用第三方库简化实现(async-mutex)
如果不想手动实现锁队列,可以使用async-mutex库,代码更简洁:
npm install async-mutex
const { Server } = require("socket.io"); const { Mutex } = require('async-mutex'); const io = new Server({}); var cache = {}; const keyMutexes = new Map(); // 每个key对应一个Mutex实例 function getKeyMutex(key) { if (!keyMutexes.has(key)) { keyMutexes.set(key, new Mutex()); } return keyMutexes.get(key); } io.on('connection', (socket) => { socket.on('clear', () => { cache = {}; keyMutexes.clear(); }); socket.on('addToQueue', async (data, cb) => { const { key } = data; const mutex = getKeyMutex(key); try { await mutex.runExclusive(async () => { if (data.type === `get`) { await new Promise(resolve => setTimeout(resolve, 20)); cb({ value: cache[key] }); } if (data.type === `transaction`) { // 同上面的transaction逻辑,省略重复代码 } }); } catch (err) { cb({ error: err.message }); } }); }); // 启动代码同前
关键说明
- 按key锁而非全局锁:只阻塞同一key的操作,不同key的请求可以并行处理,避免性能损耗。
- 锁的释放时机:必须在
finally块中释放锁,防止异常导致锁永久持有。 - 修复原代码错误:修正了
balace、tipo的拼写错误,避免逻辑异常。 - 边界处理:添加了
cache[key]不存在时的默认值,防止解构报错。
内容的提问来源于stack exchange,提问作者Fount
相关产品推荐
相关产品推荐

