邮件队列重复发送邮件及执行异常问题排查求助
邮件队列实现问题排查
问题现象
我正在实现一个邮件队列,目标是让接口快速响应,由队列管理器在后台处理邮件发送任务,但目前存在两个问题:
- 当
this.running设为true时,定时器仍会持续执行 - 部分邮件会重复发送,而非仅发送一次
我意识到initialize方法存在逻辑问题,试过用Redis,也了解bull库可作为解决方案,但希望尽量避免不必要的依赖。理想情况下想实现一个可复用基类,适配短信队列等不同类型队列,但当前优先解决邮件队列的问题。
后来我发现漏加了这段代码:
const emails = [...this.emailQueue]; for (let i = 0; i < emails.length; i++) { const [emailUuid, { attempts, emailLogEvent, ...emailConfig }] = emails[i]; + processedEmails[emailUuid] = false; emailClient(emailConfig) .then((result) => {
实现代码
import emailClient from 'services/emailClient'; import randomUuid from 'services/randomUuid' export interface EmailClientConfig { to: string; } const MAX_ATTEMPTS_FOR_EMAILS = 5; class EmailQueueManager { private running = false; private emailQueue = new Map(); private handleRemoveSentEmailsFromQueue(sentEmailUuids: string[]) { // delete each sent out email from queue sentEmailUuids.forEach((sentEmailUuid) => { this.emailQueue.delete(sentEmailUuid); }); // set state to "ready to process" this.running = false; } private async handleProcessEmailQueue() { // determines if queue is currently being executed this.running = true; const processedEmails: Record<string, boolean> = {}; // list of successfully sent emails const sentEmailUuids: string[] = []; // do nothing if currently being executed, so as // to not send out duplicate emails if (!this.emailQueue.size) { this.running = false; return; } const emails = [...this.emailQueue]; for (let i = 0; i < emails.length; i++) { const [emailUuid, { attempts, emailLogEvent, ...emailConfig }] = emails[i]; emailClient(emailConfig) .then((result) => { sentEmailUuids.push(emailUuid); // maybe do something else }) .catch((err) => { if (attempts > MAX_ATTEMPTS_FOR_EMAILS) { // remove from queue permanently if // attempt count exceeds what is allowed this.emailQueue.delete(emailUuid); } else { // send back to queue and attempt to send email out again this.emailQueue.set(emailUuid, { attempts: attempts + 1, emailLogEvent, ...emailConfig, }); } }) .finally(() => { // flag email as "processed", regardless of whether // or not is was successfully sent processedEmails[emailUuid] = true; }); } // check every 250ms if all emails in queue // were processed (not necessarily successfully sent out) const checkProcessedEmailsInterval = setInterval(() => { if (Object.values(processedEmails).every(Boolean)) { this.handleRemoveSentEmailsFromQueue(sentEmailUuids); clearInterval(checkProcessedEmailsInterval); } }, 250); } public addEmailToQueue({ emailConfig, emailLogEvent, }: { emailConfig: EmailClientConfig; emailLogEvent: string; }) { // add email to queue this.emailQueue.set(randomUuid(), { attempts: 0, emailLogEvent, ...emailConfig }); // and immediately process the queue this.handleProcessEmailQueue(); } public initialize() { // check every minute if there are // emails in queue that need to be processed setInterval(() => { if (!this.running) { this.handleProcessEmailQueue(); } }, 1000 * 60); } } const emailQueueManager = new EmailQueueManager(); // instantiate the interval emailQueueManager.initialize(); export default emailQueueManager;
接口使用示例
import emailQueueManager from 'src/managers/queues/emailQueue'; . . . async function someEndpoint(req, res) { // do stuff emailQueueManager.addEmailToQueue({ emailConfig: { to: 'test@test.com', }, emailLogEvent: 'sent_email', }); // do more stuff res.end() } . . .
问题分析与修复方案
1. running状态设置时机错误导致定时器重复触发
handleProcessEmailQueue开头直接设置this.running = true,但如果此时队列为空,会立刻把this.running设回false,这短暂的时间差内定时器可能触发,导致重复调用处理逻辑。
修复:调整状态设置顺序,先判断队列是否为空,再标记运行状态:
private async handleProcessEmailQueue() { // 先判断队列是否为空,直接返回 if (!this.emailQueue.size) { return; } // 确认有任务需要处理,再标记为运行中 this.running = true; // ... 后续代码 }
2. 邮件重复发送的核心原因
- 漏加的
processedEmails[emailUuid] = false是关键:当前代码中processedEmails为空对象,只有在finally阶段才标记为true,如果在处理这批邮件时,新的handleProcessEmailQueue调用被触发(比如addEmailToQueue调用),会重新读取队列中正在处理的邮件,导致重复发送。 - 异步任务并发执行时,队列中的邮件未被及时移除,新的处理逻辑会重复读取。
修复:
- 提前标记待处理邮件,处理前先从队列移除,避免重复读取;
- 用
Promise.all替代轮询定时器,更可靠地等待所有异步任务完成。
修改后的handleProcessEmailQueue示例:
private async handleProcessEmailQueue() { // 正在运行或队列为空时直接返回 if (this.running || !this.emailQueue.size) { return; } this.running = true; const emails = [...this.emailQueue]; // 先移除当前要处理的邮件,避免重复读取 emails.forEach(([uuid]) => this.emailQueue.delete(uuid)); const sentEmailUuids: string[] = []; const processPromises = emails.map(async ([emailUuid, { attempts, emailLogEvent, ...emailConfig }]) => { try { await emailClient(emailConfig); sentEmailUuids.push(emailUuid); } catch (err) { // 未超过重试次数则重新加入队列 if (attempts < MAX_ATTEMPTS_FOR_EMAILS) { this.emailQueue.set(emailUuid, { attempts: attempts + 1, emailLogEvent, ...emailConfig, }); } } }); // 等待所有邮件处理完成 await Promise.all(processPromises); // 标记为可处理状态 this.running = false; }
3. addEmailToQueue的冗余调用优化
每次添加邮件都调用handleProcessEmailQueue,如果当前正在运行会被拦截,但新邮件会留在队列中,等待定时器或下一次触发。可以在handleProcessEmailQueue结束后检查队列是否为空,不为空则自动继续处理,避免等待定时器。
内容的提问来源于stack exchange,提问作者Mike K
相关产品推荐
相关产品推荐

