Zac*_*ach 3 amazon-web-services node.js async-await aws-lambda
我正在编写一个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;
}
}
Run Code Online (Sandbox Code Playgroud)
我的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;
}
}
Run Code Online (Sandbox Code Playgroud)
我的助手,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()();
}
}
)
);
}
Run Code Online (Sandbox Code Playgroud)
据我所知,它应该限制并发执行limit.
另外,我已经包装了AWS SDK 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);
});
});
}
Run Code Online (Sandbox Code Playgroud)
所以没有什么花哨的,只是宣传回调.
我已经将我的lambda设置为在300秒后超时,并且lambda总是超时,如果它没有突然结束并且错过了一些应该继续的最终日志记录,这使得我甚至可能在某处出错,默默地.当我检查SQS队列时,我缺少大约1,000个条目.
我可以在你的代码中看到几个问题,
第一:
let promises = accounts.map(acc => async () => await sqsClient.sendMessage({
QueueUrl: queueUrl,
MessageBody: JSON.stringify(acc),
DelaySeconds: 10,
})
);
Run Code Online (Sandbox Code Playgroud)
你在滥用async / await.总是要记住,await等到你的承诺得到解决,然后再继续下一个,在这种情况下,无论何时映射数组promises并调用每个函数项,它都会在继续之前等待该函数包含的承诺,这很糟糕.既然你只对获得承诺感兴趣,你可以简单地这样做:
const promises = accounts.map(acc => () => sqsClient.sendMessage({
QueueUrl: queueUrl,
MessageBody: JSON.stringify(acc),
DelaySeconds: 10,
})
);
Run Code Online (Sandbox Code Playgroud)
现在,对于第二部分,您的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;
}
Run Code Online (Sandbox Code Playgroud)
最后,但非常重要的是,看看SQS限制,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
})
});
Run Code Online (Sandbox Code Playgroud)
然后你可以eachLimit照常打电话:
const result = await eachLimit(promiseFuncs, 3);
Run Code Online (Sandbox Code Playgroud)
现在的区别是每个处理的承诺都会发送一批大小为n的消息(上例中为10).
| 归档时间: |
|
| 查看次数: |
890 次 |
| 最近记录: |