Node.js中 Lambda处理SQS消息时的竞态条件

后端开发 2026-07-10

我在使用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导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。

相关文章