多PM2副本下Bull队列写入同一文件的防损坏方案咨询
我用Redis实现了一个Bull队列,接收的数据格式如下:
{ "sessionID": "43i43ko4", "events": [{...}] }
数据送入队列后由consumer.js处理,流程是:
- 从负载中获取
sessionID,检查对应{sessionID}.csv文件是否存在,不存在则创建 - 将该
sessionID任务中的events追加至文件中
之前consumer.js只运行单PM2副本,能保证同sessionID的events按序写入,运行正常。但任务量增加后,单副本处理能力不足,而且events数据量大导致内存占用持续上升,于是我把consumer.js扩展为4个PM2副本,但现在担心同sessionID的4个任务并发执行,会引发同一文件的读写冲突,导致文件损坏。
核心代码如下:
async build(job) { let { events, sessionId,... } = job.data; let filename = constructSessionReplayFilePath(this.SESSION_REPLAYS_FOLDER_PATH, target, videoID); for(event of events){ await compressDecompress(filename, event); } }
compressDecompress和constructSessionReplayFilePath方法:
export function constructSessionReplayFilePath(path, target, videoID) { return path + target + '/' + videoID + '.csv'; } export async function compressDecompress(filename, event) { const serializedObject = JSON.stringify(event); const compressedData = zlib.deflateSync(serializedObject); const dataLength = Buffer.alloc(4); dataLength.writeInt32BE(compressedData.length, 0); await fs.promises.appendFile(filename, dataLength); await fs.promises.appendFile(filename, compressedData); }
请问如何避免文件损坏?目前想到的方案是拆分4个队列,每个队列对应一个副本。
方案1:基于sessionID哈希拆分队列(你的思路优化)
不用手动拆分队列,让生产者根据sessionID的哈希值将任务分配到对应队列:把sessionID做哈希后取模4,映射到4个队列(queue-0到queue-3),每个PM2副本只监听一个固定队列。这样同sessionID的任务只会进入同一个队列,由同一个副本处理,从根源避免并发写入冲突。
生产者端示例代码:
// 计算sessionID的哈希值,取模得到队列索引 function getQueueIndex(sessionID) { let hash = 0; for (let i = 0; i < sessionID.length; i++) { hash = (hash << 5) - hash + sessionID.charCodeAt(i); hash |= 0; // 转为32位整数 } return Math.abs(hash) % 4; // 4对应副本数 } // 发送任务到对应队列 const queueIndex = getQueueIndex(sessionID); const targetQueue = new Bull(`session-queue-${queueIndex}`, redisConfig); await targetQueue.add(jobData);
消费者端示例代码:
// 通过PM2环境变量传递队列索引,比如启动命令:pm2 start consumer.js --name consumer-0 --env QUEUE_INDEX=0 const queueIndex = process.env.QUEUE_INDEX || 0; const queue = new Bull(`session-queue-${queueIndex}`, redisConfig); queue.process(async (job) => { // 原build方法逻辑 });
方案2:使用文件锁控制并发写入
如果不想改动队列结构,可以在写入文件前加文件锁,确保同一时间只有一个进程能写入指定文件。可以用proper-lockfile这类库实现文件锁。
修改后的compressDecompress方法:
import lockfile from 'proper-lockfile'; import fs from 'fs/promises'; import zlib from 'zlib'; export async function compressDecompress(filename, event) { let release; try { // 获取文件锁,超时时间设为5秒(可调整) release = await lockfile.lock(filename, { timeout: 5000 }); const serializedObject = JSON.stringify(event); const compressedData = zlib.deflateSync(serializedObject); const dataLength = Buffer.alloc(4); dataLength.writeInt32BE(compressedData.length, 0); await fs.appendFile(filename, dataLength); await fs.appendFile(filename, compressedData); } catch (err) { console.error(`写入文件${filename}失败:`, err); throw err; } finally { if (release) { await release(); // 释放锁 } } }
注意:文件锁能解决冲突,但同sessionID任务量过大时会导致锁等待,影响处理效率,适合并发度不极端的场景。
方案3:利用BullMQ的分组特性
如果使用的是BullMQ(Bull的升级版本),可以给同sessionID的任务设置相同的分组ID,确保同组任务串行处理,不同组任务并行执行,兼顾处理效率和写入顺序。
生产者端:
await queue.add(jobData, { group: { id: sessionID } });
消费者端:
queue.process(async (job) => { // 处理逻辑 }, { concurrency: 4, group: { concurrency: 1 } });
方案4:拆分预处理与写入流程
把文件写入逻辑从consumer中拆分出来,单独用一个单副本进程负责写入:
- 多副本consumer只做
events的解析、压缩等预处理工作 - 将预处理后的
sessionID和event发送到专门的write-queue - 单副本consumer监听
write-queue,按顺序写入对应文件
这种方式既利用多副本处理计算密集的预处理,降低内存占用,又从根源避免了文件写入冲突。
内容的提问来源于stack exchange,提问作者Dren

