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

多PM2副本下Bull队列写入同一文件的防损坏方案咨询

问题描述

我用Redis实现了一个Bull队列,接收的数据格式如下:

{
  "sessionID": "43i43ko4",
  "events": [{...}]
}

数据送入队列后由consumer.js处理,流程是:

  1. 从负载中获取sessionID,检查对应{sessionID}.csv文件是否存在,不存在则创建
  2. 将该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中拆分出来,单独用一个单副本进程负责写入:

  1. 多副本consumer只做events的解析、压缩等预处理工作
  2. 将预处理后的sessionID和event发送到专门的write-queue
  3. 单副本consumer监听write-queue,按顺序写入对应文件

这种方式既利用多副本处理计算密集的预处理,降低内存占用,又从根源避免了文件写入冲突。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 01:30:59