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

Node.js IMAP连接意外终止:定时爬取邮件更新DB时最后4封处理失败

问题诊断与修复方案

核心问题

  1. 异步操作不同步:邮件解析(simpleParser)是异步操作,但fetch的end事件会在所有邮件的流读取完成后立即触发,此时可能还有部分邮件的解析逻辑未执行完毕,导致visitedData数据不完整,后续数据库操作时可能引发异常,间接导致IMAP连接异常。
  2. IMAP连接闲置超时:数据库处理逻辑("some complex logic")耗时较长,此时IMAP连接处于闲置状态,容易被邮件服务器主动断开连接。

修复步骤

1. 确保邮件解析全部完成后再处理数据库

将邮件解析的异步操作收集为Promise,等待所有解析完成后再进入数据库处理环节,同时把数据库操作移到IMAP连接关闭之后执行:

export const readAllEmailReports = async () => {  
  try {
    let visitedData = {};
    const imap = new Imap(imapConfig);
    const yesterday = moment().subtract(5, "hour").toDate();

    // 用Promise包裹IMAP流程,统一控制异步逻辑
    await new Promise((resolve, reject) => {
      imap.once("ready", () => {
        imap.openBox("INBOX", false, (err) => {
          if (err) {
            console.error(err);
            reject(err);
            return;
          }

          imap.search(
            [
              "ALL",
              ["FROM", "test@test.com"],
              [
                "OR",
                ["SUBJECT", "Complete_Report"],
                ["SUBJECT", "Partial_Report"],
              ],
              ["SINCE", yesterday],
            ],
            (err, results) => {
              if (err) {
                console.error(err);
                reject(err);
                return;
              }

              if (results.length === 0) {
                console.log("No new messages to fetch.");
                imap.end();
                resolve();
                return;
              }

              const f = imap.fetch(results, { bodies: "" });
              // 收集所有邮件解析的Promise
              const parsePromises = [];

              f.on("message", (msg) => {
                const parsePromise = new Promise((msgResolve) => {
                  msg.on("body", (stream) => {
                    simpleParser(stream, (err, parsed) => {
                      if (err) {
                        console.error("Error parsing email:", err);
                        msgResolve();
                        return;
                      }

                      const text = parsed.text ? parsed.text : parsed.html;
                      const parsedText = extractVidInfo(text, parsed?.subject);

                      const key = `${parsedText.visitNumber}-${parsedText.reportType}`;
                      if (!visitedData[key]) {
                        visitedData[key] = {
                          visitNumber: parsedText?.visitNumber,
                          reportType: parsedText?.reportType,
                          subject: parsed?.subject,
                          attachments: [],
                        };

                        if (parsed.attachments && parsed.attachments.length > 0) {
                          for (let attachment of parsed.attachments) {
                            if (isSupportedAttachment(attachment))
                              visitedData[key].attachments.push(attachment);
                          }
                        }
                      } else {
                        console.log(
                          "not inserting",
                          parsedText.visitNumber,
                          "----",
                          parsedText.reportType
                        );
                      }
                      msgResolve();
                    });
                  });
                });
                parsePromises.push(parsePromise);
              });

              f.once("error", (ex) => {
                console.error(ex);
                reject(ex);
              });

              f.once("end", async () => {
                // 等待所有邮件解析完成
                await Promise.all(parsePromises);
                console.log("Done parsing all messages!");
                imap.end();
                resolve();
              });
            }
          );
        });
      });

      imap.once("error", (err) => {
        console.error(err);
        reject(err);
      });

      imap.once("end", () => {
        console.log("Connection ended");
      });

      imap.connect();
    });

    // 所有邮件解析完成、IMAP连接关闭后,再处理数据库
    const processingTasks = Object.entries(visitedData).map(async ([key, visitData]) => {
      // 执行你的复杂数据库逻辑
    });
    if (processingTasks.length > 0) {
      await Promise.all(processingTasks);
    }
    console.log("Done processing all data!");

  } catch (ex) {
    console.error("An error occurred", ex);  
  } 
};

2. 增加IMAP连接心跳(可选)

如果数据库操作确实耗时极长,可以在IMAP连接期间定时发送心跳包,防止服务器因连接闲置断开:

// 在imap.once("ready")之后添加心跳定时器
let heartbeatInterval = setInterval(() => {
  if (imap.state === 'authenticated' || imap.state === 'selected') {
    // 发送NOOP命令作为心跳,不执行实际操作仅保持连接
    imap.noop((err) => {
      if (err) console.error("Heartbeat failed:", err);
    });
  }
}, 30000); // 每30秒发送一次心跳

// 在imap.once("end")里清除定时器
imap.once("end", () => {
  clearInterval(heartbeatInterval);
  console.log("Connection ended");
});

3. 优化数据库操作

  • 采用批量更新/插入的方式,减少与数据库的交互次数,降低整体耗时。
  • 检查数据库逻辑中是否有重复查询或冗余操作,针对性优化。

内容的提问来源于stack exchange,提问作者Sushanth

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 20:04:55