【问题标题】:Read stream with settimeout maximum value reached error读取具有 settimeout 最大值的流已达到错误
【发布时间】:2020-05-17 02:35:00
【问题描述】:

我正在尝试读取一些大型 CSV 文件并处理这些数据,因此处理中存在速率限制,因此我想在每个请求之间添加 1mnt 延迟。我尝试了设置超时,但最后,知道设置超时有限制并得到以下错误。我不确定任何其他方式来处理这种情况,CSV 文件有超过 1M 的记录。我在这里做错什么了吗?

错误

超时持续时间设置为 1。(节点:41) TimeoutOverflowWarning: 2241362000 不适合 32 位有符号整数。

示例代码:

   const Queue = require('bull');
const domainQueue = new Queue(config.api.crawlerQ, {
  redis: connectRedis(),
});
let ctr = 0;
function processCSV (name, fileName, options)  {
  return new Promise((resolve, reject) => {
    console.log('process csv started', new Date());
    let filePath = config.api.basePath + fileName;
    stream = fs.createReadStream(filePath)
        .on('error', (error) => {
          // handle error
          console.log('error processing csv');
          reject(error);
        })
        .pipe(csv())
        .on('data', async (row) => {
          ctr++
          increment(row, ctr)
        })
        .on('end', () => {
          console.log('stream processCSV end', fileName, new Date());
          resolve(filePath);
        })
  });

}

async function increment(raw, counter) {
  setTimeout(async function(){
    console.log('say i am inside a function', counter, new Date());
    domainQueue.add(data, options); // Add jobs to queue - Here i Need a delay say 1mnt, if i
    // add jobs without delay it will hit ratelimit 
  }, 60000 * counter);

}

function queueWorkerProcess(value) { // Process jobs in queue and save in text file 
  console.log('value', value, new Date());
  return new Promise(resolve => {
    resolve();
  });

}

【问题讨论】:

  • 我真的不知道你想在这里做什么。您在每个文件的每一行都调用increment()。显然,60000 * counter 的数字太大了。计时器似乎实际上并没有做任何有用的事情,所以我不知道你想用它来完成什么。
  • 因此,此代码处理一个本地 CSV 文件。你想用它解决什么问题?您在哪里遇到速率限制问题?我在这段代码中看不到它。我们需要查看实际存在问题的真实代码,并且需要对实际问题进行完整描述。
  • 是的,excel包含域名,我处理该文件的每一行并进行外部api调用,所以我希望每个请求之间有延迟,所以我曾经反击。见:borgs.cybrilla.com/tils/… 问题见:stackoverflow.com/questions/3468607/…
  • 如果您展示您的实际代码以及您实际尝试做什么以及实际问题是什么,我可能会在大约 5 分钟内为您提供帮助。从字面上看,如果你展示了整个问题,我可以在几分钟内解决你的问题。但是,我拒绝写试图猜测你真正想要做什么的答案。到目前为止,您显示的代码实际存在的唯一问题是您试图创建一个太大的整数。删除此处没有实际用途的setTimeout(),此代码没有问题。向我们展示您真正的问题和真正的代码。
  • 我编辑了我的问题,我正在使用公牛队列作业,我希望在将作业添加到队列时延迟

标签: javascript node.js loops settimeout delay


【解决方案1】:

这是一个总体思路。您需要跟踪正在处理的项目数量,以限制使用的内存量并控制存储结果的任何资源的负载。

当您达到飞行中的数量限制时,您会暂停流。当您回到限制以下时,您将恢复流。您在.add() 上增加一个计数器并在completed 消息上减少一个计数器以跟踪事情。您可以在此处暂停或恢复直播。

仅供参考,仅在某处插入 setTimeout() 对您没有帮助。为了控制内存使用,一旦处理中的项目过多,您必须暂停流中的数据流。然后,当项目恢复到阈值以下时,您可以恢复流。

这是一个大概的样子:

const Queue = require('bull');
const domainQueue = new Queue(config.api.crawlerQ, {
    redis: connectRedis(),
});

// counter that keeps track of how many items in the queue
let queueCntr = 0;

// you tune this constant up or down to manage memory usage or tweak performance
// this is what keeps you from having too many requests going at once
const queueMax = 20;

function processCSV(name, fileName, options) {
    return new Promise((resolve, reject) => {
        let paused = false;

        console.log('process csv started', new Date());
        const filePath = config.api.basePath + fileName;

        const stream = fs.createReadStream(filePath)
            .on('error', (error) => {
                // handle error
                console.log('error processing csv');
                domainQueue.off('completed', completed);
                reject(error);
            }).pipe(csv())
            .on('data', async (row) => {
                increment(row, ctr);
                if (queueCntr)
            })
            .on('end', () => {
                console.log('stream processCSV end', fileName, new Date());
                domainQueue.off('completed', completed);
                resolve(filePath);
            });

        function completed() {
            --queueCntr;
            // see if queue got small enough we now resume the stream
            if (paused && queueCntr < queueMax) {
                stream.resume();
                paused = false;
            }
        }

        domainQueue.on('completed', completed);

        function increment(raw, counter) {
            ++queueCntr;
            domainQueue.add(data, options);
            if (!paused && queueCntr > queueMax) {
                stream.pause();
                paused = true;
            }
        }
    });
}

而且,如果您使用不同的文件多次调用 processCSV(),您应该对它们进行排序,以便在第一个完成之前不要调用第二个,在第二个完成之前不要调用第三个已完成,依此类推...您没有显示该代码,因此我们无法对此提出具体建议。

【讨论】:

  • 谢谢@jfriend00,您的评论对我很有帮助。我采用暂停/恢复的想法来控制问题。我尝试了多种方法来延迟,但我一直坚持下去。感谢您建议暂停直播。
  • 嗨,它对我有用,但我现在有一个问题,因为我正在使用上面的单个流异步处理多个 CSV 文件,在暂停/恢复之后,一些流没有结束,即那些没有给出' end' 事件,似乎那些在中间被击中或丢失。我是否需要使用动态流来处理每个 csv?
  • async function readDirectory() { try { fs.readdir(config.api.basePath + '/csvdir/', (err, files) => { if (err) throw err; for (const文件文件) { console.log('file', file); processCSV(file) } }); } 捕捉(错误){ } } 读取目录();
猜你喜欢
  • 1970-01-01
  • 2020-06-07
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-03-05
  • 2017-03-13
  • 2020-02-25
  • 2023-03-22
相关资源
最近更新 更多