Node.js中如何阻塞所有工作线程直至首个工作线程完成?
解决Node.js多Worker线程中fetchDbData方法的同步问题
核心思路
Node.js的Worker线程内存相互隔离,无法直接共享内存锁,因此需要通过主线程中转、文件锁或者数据库原生锁机制来实现跨线程的排他性访问,模拟Java中synchronized的效果,确保同一时间只有一个线程执行fetchDbData。
具体方案
方案一:主线程消息传递实现锁
让主线程维护锁状态,所有Worker线程调用fetchDbData前先向主线程申请锁,获得授权后再执行操作,完成后释放锁。
- 主线程代码:
const { Worker } = require('worker_threads'); let isLocked = false; const lockQueue = []; // 处理锁申请 function handleLockRequest(worker) { if (!isLocked) { isLocked = true; worker.postMessage({ type: 'GRANT_LOCK' }); } else { lockQueue.push(worker); } } // 处理锁释放 function handleLockRelease() { isLocked = false; if (lockQueue.length > 0) { const nextWorker = lockQueue.shift(); isLocked = true; nextWorker.postMessage({ type: 'GRANT_LOCK' }); } } // 创建Worker实例 for (let i = 0; i < 2; i++) { const worker = new Worker('./worker.js'); worker.on('message', (msg) => { if (msg.type === 'REQUEST_LOCK') { handleLockRequest(worker); } else if (msg.type === 'RELEASE_LOCK') { handleLockRelease(); } }); }
- Worker线程代码(worker.js):
const { parentPort } = require('worker_threads'); const fetchDbData = require('./your-db-module').fetchDbData; async function getLockedDbData() { // 向主线程申请锁 parentPort.postMessage({ type: 'REQUEST_LOCK' }); // 等待锁授权 await new Promise(resolve => { parentPort.once('message', (msg) => { if (msg.type === 'GRANT_LOCK') { resolve(); } }); }); try { // 执行数据库操作 const data = await fetchDbData(); return data; } finally { // 操作完成后释放锁 parentPort.postMessage({ type: 'RELEASE_LOCK' }); } } // 调用带锁的数据库查询方法 getLockedDbData().then(data => console.log('Worker获取到数据:', data));
方案二:文件锁实现跨线程/进程同步
如果需要跨多个Node.js进程同步,或者Worker线程数量较多,可以用文件锁(依赖fs-ext库),通过文件的排他锁确保同一时间只有一个线程执行目标方法。
- 安装依赖:
npm install fs-ext
- 封装带锁的fetchDbData:
const fs = require('fs'); const fsext = require('fs-ext'); const lockFile = './db-lock.lock'; const fetchDbData = require('./your-db-module').fetchDbData; async function lockedFetchDbData() { // 创建锁文件(不存在则自动创建) await fs.promises.open(lockFile, 'w').then(fd => fd.close()); // 获取排他锁 await new Promise((resolve, reject) => { fsext.flock(fs.openSync(lockFile, 'r'), 'ex', (err) => { if (err) reject(err); else resolve(); }); }); try { // 执行原数据库查询逻辑 return await fetchDbData(); } finally { // 释放锁 fsext.flock(fs.openSync(lockFile, 'r'), 'un', () => {}); } }
方案三:利用数据库原生锁机制
直接借助数据库的锁能力,比如MySQL的SELECT ... FOR UPDATE行级锁,从根源上保证同一时间只有一个事务能获取目标数据,避免应用层锁的潜在问题。
示例MySQL查询语句:
SELECT * FROM your_table WHERE status = 'available' LIMIT 1 FOR UPDATE;
执行该语句后,其他线程的同语句会阻塞,直到当前事务提交或回滚,确保每次获取的是未被锁定的数据集。
方案选择建议
- 仅Worker线程间同步:优先选方案一,轻量无额外依赖。
- 跨进程同步:选方案二,支持多进程场景。
- 追求可靠性:选方案三,数据库原生锁更稳定,避免应用层锁的BUG。
内容的提问来源于stack exchange,提问作者Rajiv Sharma
相关产品推荐
相关产品推荐

