You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

邮件队列重复发送邮件及执行异常问题排查求助

邮件队列实现问题排查

问题现象

我正在实现一个邮件队列,目标是让接口快速响应,由队列管理器在后台处理邮件发送任务,但目前存在两个问题:

  • 当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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.05 11:05:27