【问题标题】:AWS Lambda function that executes 5000+ promises to AWS SQS is extremely unreliable对 AWS SQS 执行 5000 多个承诺的 AWS Lambda 函数极其不可靠
【发布时间】:2018-07-17 00:42:55
【问题描述】:

我正在编写一个 Node AWS Lambda 函数,该函数从我的数据库中查询大约 5,000 个项目,并通过消息将它们发送到 AWS SQS 队列中。

我的本​​地环境包括我使用本地 AWS SAM 运行我的 lambda,并使用 GoAWS 模拟 AWS SQS。

我的 Lambda 的一个示例框架是:

async run() {
  try {
    const accounts = await this.getAccountsFromDB();
    const results = await this.writeAccountsIntoQueue(accounts);
    return 'I\'ve written: ' + results + ' messages into SQS';
  } catch (e) {
    console.log('Caught error running job: ');
    console.log(e);
    return e;
  }
}

我的 getAccountsFromDB() 函数没有性能问题,它几乎可以立即运行,返回包含 5,000 个帐户的数组。

我的writeAccountsIntoQueue 函数如下所示:

async writeAccountsIntoQueue(accounts) {
  // Extract the sqsClient and queueUrl from the class 
  const { sqsClient, queueUrl } = this;
  try {
    // Create array of functions to concurrenctly call later
    let promises = accounts.map(acc => async () => await sqsClient.sendMessage({
        QueueUrl: queueUrl,
        MessageBody: JSON.stringify(acc),
        DelaySeconds: 10,
      })
    );

    // Invoke the functions concurrently, using helper function `eachLimit`
    let writtenMessages = await eachLimit(promises, 3);
    return writtenMessages;
  } catch (e) {
    console.log('Error writing accounts into queue');
    console.log(e);
    return e;
  }
}

我的助手eachLimit 看起来像:

async function eachLimit (funcs, limit) {
  let rest = funcs.slice(limit);
  await Promise.all(
    funcs.slice(0, limit).map(
      async (func) => {
        await func();
        while (rest.length) {
          await rest.shift()();
        }
      }
    )
  );
}

据我所知,它应该将并发执行限制为limit

此外,我还包装了 AWS 开发工具包 SQS 客户端以返回一个带有 sendMessage 函数的对象,如下所示:

sendMessage(params) {
  const { client } = this;
  return new Promise((resolve, reject) => {
    client.sendMessage(params, (err, data) => {
      if (err) {
        console.log('Error sending message');
        console.log(err);
        return reject(err);
      }
      return resolve(data);
    });
  });
}

所以没什么特别的,只是承诺回调。

我已将我的 lambda 设置为 300 秒后超时,并且 lambda 总是超时,如果不是,它会突然结束并错过一些应该继续进行的最终日志记录,这让我觉得它甚至可能默默地在某个地方出错。当我检查 SQS 队列时,我丢失了大约 1,000 个条目。

【问题讨论】:

    标签: node.js amazon-web-services async-await aws-lambda


    【解决方案1】:

    我可以在您的代码中看到几个问题,

    第一:

        let promises = accounts.map(acc => async () => await sqsClient.sendMessage({
            QueueUrl: queueUrl,
            MessageBody: JSON.stringify(acc),
            DelaySeconds: 10,
          })
        );
    

    你在滥用async / await。永远记住await 将等到你的承诺得到解决后再继续下一个承诺,在这种情况下,每当你映射数组promises 并调用每个函数项时,它都会在继续之前等待该函数包装的承诺,这很糟糕。由于您只对收回承诺感兴趣,因此您可以简单地这样做:

    const promises = accounts.map(acc => () => sqsClient.sendMessage({
           QueueUrl: queueUrl,
           MessageBody: JSON.stringify(acc),
           DelaySeconds: 10,
        })
    );
    

    现在,对于第二部分,您的 eachLimit 实现看起来错误且非常冗长,我在 es6-promise-pool 的帮助下对其进行了重构,以便为您处理并发限制:

    const PromisePool = require('es6-promise-pool')
    
    function eachLimit(promiseFuncs, limit) {    
        const promiseProducer = function () {
            while(promiseFuncs.length) {
                const promiseFunc = promiseFuncs.shift();
                return promiseFunc();
            }
    
            return null;
        }
    
        const pool = new PromisePool(promiseProducer, limit)
        const poolPromise = pool.start();
        return poolPromise;
    }
    

    最后,但非常重要的是,看看SQS Limits,SQS FIFO 每秒发送高达 300 次。由于您正在处理 5k 个项目,您可能会将并发限制提高到 5k / (300 + 50) ,大约为 15。50 可以是任何正数,只是为了稍微远离限制。 此外,考虑使用SendMessageBatch,您可以拥有更高的吞吐量并达到 3k 发送/秒。

    编辑

    正如我上面建议的那样,使用sendMessageBatch 的吞吐量要好得多,所以我重构了映射您的承诺的代码以支持sendMessageBatch

    function chunkArray(myArray, chunk_size){
        var index = 0;
        var arrayLength = myArray.length;
        var tempArray = [];
    
        for (index = 0; index < arrayLength; index += chunk_size) {
            myChunk = myArray.slice(index, index+chunk_size);
            tempArray.push(myChunk);
        }
    
        return tempArray;
    }
    
    const groupedAccounts = chunkArray(accounts, 10);
    
    const promiseFuncs = groupedAccounts.map(accountsGroup => {
        const messages = accountsGroup.map((acc,i) => {
            return {
                Id: `pos_${i}`,
                MessageBody: JSON.stringify(acc),
                DelaySeconds: 10
            }
        });
    
        return () => sqsClient.sendMessageBatch({
            Entries: messages,
            QueueUrl: queueUrl
         })
    });
    

    然后你可以像往常一样拨打eachLimit

    const result = await eachLimit(promiseFuncs, 3);
    

    现在的不同之处在于每个处理过的 Promise 都会发送一批大小为 n(上例中为 10)的消息。

    【讨论】:

    • 我相信一些引擎现在已经足够聪明,可以内联一些只返回另一个 Promise 的 Promise 链。我可能错了 - 需要在以后的 chrome/nodejs 版本上进行基准测试。
    • 非常感谢您的回复和见解。我已经进行了更改,再次运行后,看起来 lambda 执行了一些时间,但随后突然下降。检查队列后,看起来有 5266 个条目进入其中,但应该写入了 5501。关于可能导致这种情况的任何想法?
    • 很难说。你有任何错误吗?你设置了什么并发限制?此外,您应该尝试使用真正的 SQS,并将那里的结果与 GoAWS 进行比较。
    • 没有错误或输出,我什至将.then(d =&gt; {console.log('success')}) 绑定到sqsClient.sendMessage(...) 以及一个计数器,我看到它发生了大约 5266 次,挂起,然后退出,没有更多输出或任何错误。我将并发设置为 10,当我回到家时,我可以在普通 SQS 上试试这个
    • “你在滥用 async / await”应该是“你在 misusing async / await”。
    猜你喜欢
    • 1970-01-01
    • 2019-01-28
    • 2019-11-16
    • 1970-01-01
    • 2016-10-16
    • 2020-02-15
    • 2016-04-10
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多