【发布时间】: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