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

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 });
        }
    });
});

// 启动代码同前

关键说明

  1. 按key锁而非全局锁:只阻塞同一key的操作,不同key的请求可以并行处理,避免性能损耗。
  2. 锁的释放时机:必须在finally块中释放锁,防止异常导致锁永久持有。
  3. 修复原代码错误:修正了balace、tipo的拼写错误,避免逻辑异常。
  4. 边界处理:添加了cache[key]不存在时的默认值,防止解构报错。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 06:00:14