Node.js中 Lambda处理SQS消息时的竞态条件
我在使用AWS Lambda + SQS处理Webhook事件时,遇到了一个竞态条件,即使有保护措施也会造成重复处理。我已经尝试在orderId上添加唯一索引,并将insertOne放在try/catch块中包裹。
code:
const processMessage = async (message) => {
const { orderId } = JSON.parse(message.body);
const existing = await db.collection("orders").findOne({ orderId });
if (existing) {
console.log("Already processed:", orderId);
return;
}
const details = await getOrderDetails(orderId);
await db.collection("orders").insertOne({
orderId,
details,
createdAt: new Date(),
});
console.log("Processed:", orderId);
};
Lambda处理程序:
export const handler = async (event) => {
await Promise.all(
event.Records.map((record) => processMessage(record))
);
};
解决方案
将分离的读写替换为一个原子操作,因为MongoDB保证原子性:
const result = await db.collection("orders").updateOne(
{ orderId },
{
$setOnInsert: {
orderId,
createdAt: new Date(),
},
},
{ upsert: true }
);
if (result.upsertedCount === 0) {
console.log("Already processed:", orderId);
return;
}
即使使用upsert,在数据库插入完成之前,仍会有多个Lambda同时调用 getOrderDetails() 的问题,因此要阻止外部API调用的重复:
const result = await db.collection("orders").findOneAndUpdate(
{ orderId },
{
$setOnInsert: {
orderId,
status: "processing",
createdAt: new Date(),
},
},
{ upsert: true, returnDocument: "after" }
);
if (result.lastErrorObject.updatedExisting) {
console.log("Already being processed or done:", orderId);
return;
}
// Only ONE Lambda reaches here
const details = await getOrderDetails(orderId);
await db.collection("orders").updateOne(
{ orderId },
{ $set: { status: "completed", details } }
)
并使用Redis来对重复进行短路处理
const cacheKey = `order:${orderId}`;
const cached = await redis.get(cacheKey);
if (cached) {
console.log("Cached result:", orderId);
return;
}
await redis.set(cacheKey, "processing", "EX", 60);
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。