This commit is contained in:
oneao committed 2025-07-11 17:26:48 +08:00
1 parent 7837b87959
commit 3026145786
7 files changed
+334 -106

No files matched your search

@@ -3,6 +3,7 @@ package com.email;
import com.email.core.fetch.EmailFetcher;
import com.email.core.fetch.ImapJavaMailEmailFetcher;
import com.email.domain.entity.UserAccount;
import lombok.extern.slf4j.Slf4j;
import javax.mail.*;
import javax.mail.Flags;
@@ -12,6 +13,7 @@ import javax.net.ssl.SSLSocketFactory;
import java.io.*;
import java.util.Properties;
@Slf4j
public class Test {
public static void main(String[] args) throws Exception {
UserAccount userAccount = new UserAccount();
@@ -24,22 +26,14 @@ public class Test {
emailFetcher.fetchAll(
userAccount,
// progressConsumer: 打印进度
progress -> System.out.printf("📩 [%s] 已拉取 %d / %d 封邮件%n",
progress.getFolderName(),
progress.getCurrent(),
progress.getTotal()
),
progress -> {
},
// folderConsumer: 开始处理某个文件夹
folder -> {
},
// fetchDataConsumer: 每封邮件拉取完成后处理
emailData -> {
}
);
}
@@ -191,7 +191,6 @@ public class EmailImapParser {
/** 从 Part 上提取并设置 contentType / charset / encoding */
private void populateMeta(Part part, EmailSummary emailSummary) throws Exception {
String ct = part.getContentType();
// 抽 charset
Matcher m = Pattern.compile("charset=\\s*\"?([^;\"\\s]+)", Pattern.CASE_INSENSITIVE)
@@ -5,11 +5,12 @@ import com.email.core.fetch.model.EmailFetchProgress;
import com.email.domain.entity.EmailFolder;
import com.email.domain.entity.UserAccount;
import java.util.List;
import java.util.function.Consumer;
public interface EmailFetcher {
void fetchAll(UserAccount userAccount,
Consumer<EmailFetchProgress> emailFetchProgressConsumer,
Consumer<EmailFolder> emailFolderConsumer,
Consumer<EmailFetchData> emailFetchDataConsumer);
Consumer<List<EmailFetchData>> emailFetchDataConsumer);
}
@@ -1,142 +1,361 @@
package com.email.core.fetch;
import com.baomidou.mybatisplus.core.toolkit.IdWorker;
import com.email.constants.ResponseStatusConstants;
import com.email.core.fetch.model.EmailFetchData;
import com.email.core.fetch.model.EmailFetchProgress;
import com.email.domain.entity.EmailFolder;
import com.email.domain.entity.UserAccount;
import com.email.core.fetch.model.FolderWrapper;
import com.email.domain.entity.*;
import com.email.enums.ResponseEnum;
import com.email.enums.email.EmailBodyTypeEnum;
import com.email.enums.email.EmailFolderTypeEnum;
import com.email.exception.auth.AuthException;
import com.sun.javaws.progress.Progress;
import com.sun.mail.imap.IMAPFolder;
import com.sun.mail.imap.IMAPStore;
import lombok.extern.slf4j.Slf4j;
import org.springframework.util.IdGenerator;
import org.springframework.util.ObjectUtils;
import javax.mail.*;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
import java.util.Properties;
import javax.mail.internet.MailDateFormat;
import javax.mail.internet.MimeMessage;
import java.text.ParseException;
import java.util.*;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Consumer;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
@Slf4j
public class ImapJavaMailEmailFetcher implements EmailFetcher {
private final ExecutorService folderExecutor = Executors.newFixedThreadPool(10);
private final ExecutorService folderExecutor = Executors.newFixedThreadPool(5);
private final ExecutorService emailExecutor = Executors.newFixedThreadPool(20);
private static final int BATCH_SIZE = 100;
@Override
public void fetchAll(UserAccount userAccount, Consumer<EmailFetchProgress> progressConsumer, Consumer<EmailFolder> folderConsumer, Consumer<EmailFetchData> fetchDataConsumer) {
public void fetchAll(UserAccount userAccount,
Consumer<EmailFetchProgress> progressConsumer,
Consumer<EmailFolder> folderConsumer,
Consumer<List<EmailFetchData>> fetchDataConsumer) {
Store store = null;
try {
store = connect(userAccount); // connect 内部已处理异常
store = connect(userAccount);
List<EmailFolder> folders = listAllFolders(store.getDefaultFolder(), 0L, userAccount.getId());
store.close(); // 提前关闭共享 store
List<Folder> allFolders = listAllFolders(store.getDefaultFolder());
CountDownLatch folderLatch = new CountDownLatch(allFolders.size());
for (Folder folder : allFolders) {
folderExecutor.submit(() -> {
try {
processFolder(folder, progressConsumer, folderConsumer, fetchDataConsumer);
} catch (Exception e) {
log.error("处理文件夹 [{}] 失败:{}", folder.getFullName(), e.getMessage(), e);
} finally {
folderLatch.countDown();
}
});
}
folderLatch.await(); // 等待所有文件夹处理完
} catch (Exception e) {
processAllFolders(userAccount, folders, progressConsumer, folderConsumer, fetchDataConsumer);
} catch (Exception e) {
throw new RuntimeException("抓取邮件失败:" + e.getMessage(), e);
} finally {
if (store != null && store.isConnected()) {
shutdownExecutors();
}
}
private void processAllFolders(UserAccount userAccount,
List<EmailFolder> folders,
Consumer<EmailFetchProgress> progressConsumer,
Consumer<EmailFolder> folderConsumer,
Consumer<List<EmailFetchData>> fetchDataConsumer) throws InterruptedException {
CountDownLatch folderLatch = new CountDownLatch(folders.size());
for (EmailFolder emailFolder : folders) {
folderExecutor.submit(() -> {
Store threadStore = null;
try {
store.close();
} catch (MessagingException closeEx) {
log.warn("关闭邮箱连接时异常:{}", closeEx.getMessage());
threadStore = connect(userAccount);
processSingleFolder(threadStore, emailFolder, progressConsumer, folderConsumer, fetchDataConsumer);
} catch (Exception e) {
log.error("处理文件夹 [{}] 异常:{}", emailFolder.getFullName(), e.getMessage(), e);
} finally {
if (threadStore != null && threadStore.isConnected()) {
try {
threadStore.close();
} catch (MessagingException e) {
log.warn("线程关闭邮箱连接异常", e);
}
}
folderLatch.countDown();
}
});
}
folderLatch.await();
}
private void processSingleFolder(Store store,
EmailFolder emailFolder,
Consumer<EmailFetchProgress> progressConsumer,
Consumer<EmailFolder> folderConsumer,
Consumer<List<EmailFetchData>> fetchDataConsumer) {
IMAPFolder imapFolder = null;
try {
imapFolder = (IMAPFolder) store.getFolder(emailFolder.getFullName());
// 文件夹 通知
folderConsumer.accept(emailFolder);
if ((imapFolder.getType() & Folder.HOLDS_MESSAGES) != 0) {
imapFolder.open(Folder.READ_ONLY);
int totalMessages = imapFolder.getMessageCount();
if (totalMessages > 0) {
processEmailsInFolder(imapFolder,emailFolder);
}
}
} catch (MessagingException e) {
log.error("打开或处理文件夹 [{}] 失败: {}", emailFolder.getFullName(), e.getMessage(), e);
} finally {
if (imapFolder != null && imapFolder.isOpen()) {
try {
imapFolder.close(false);
} catch (MessagingException e) {
log.warn("关闭文件夹异常", e);
}
}
}
}
// 递归获取所有文件夹(包括子文件夹)
private List<Folder> listAllFolders(Folder folder) throws MessagingException {
List<Folder> allFolders = new ArrayList<>();
if ((folder.getType() & Folder.HOLDS_MESSAGES) != 0) {
allFolders.add(folder);
private void processEmailsInFolder(IMAPFolder imapFolder,EmailFolder emailFolder) throws MessagingException {
if (!imapFolder.isOpen()) {
imapFolder.open(Folder.READ_ONLY);
}
if ((folder.getType() & Folder.HOLDS_FOLDERS) != 0) {
Folder[] subFolders = folder.list();
for (Folder sub : subFolders) {
allFolders.addAll(listAllFolders(sub));
}
}
return allFolders;
}
private void processFolder(Folder folder,
Consumer<EmailFetchProgress> progressConsumer,
Consumer<EmailFolder> folderConsumer,
Consumer<EmailFetchData> fetchDataConsumer) throws Exception {
if (!(folder instanceof IMAPFolder)) return;
folder.open(Folder.READ_ONLY);
int totalMessages = folder.getMessageCount();
if (totalMessages == 0) {
folder.close(false);
Message[] messages = imapFolder.getMessages();
if (messages == null || messages.length == 0) {
log.info("文件夹 [{}] 没有邮件,跳过", imapFolder.getFullName());
return;
}
EmailFolder emailFolder = new EmailFolder();
List<Message> messageList = Arrays.asList(messages);
int totalMessages = messageList.size();
log.info("文件夹 [{}] 邮件总数: {}", imapFolder.getFullName(), totalMessages);
folderConsumer.accept(emailFolder);
int batchCount = (totalMessages + BATCH_SIZE - 1) / BATCH_SIZE;
CountDownLatch latch = new CountDownLatch(batchCount);
FetchProfile fetchProfile = new FetchProfile();
fetchProfile.add(FetchProfile.Item.ENVELOPE);
fetchProfile.add(FetchProfile.Item.CONTENT_INFO);
for (int i = 0; i < totalMessages; i += BATCH_SIZE) {
int end = Math.min(i + BATCH_SIZE, totalMessages);
List<Message> batchMessages = messageList.subList(i, end);
Message[] messages = folder.getMessages();
folder.fetch(messages, fetchProfile);
emailExecutor.submit(() -> {
try {
for (Message message : batchMessages) {
for (int i = 0; i < messages.length; i++) {
Message message = messages[i];
try {
progressConsumer.accept(new EmailFetchProgress(folder.getFullName(), i + 1, totalMessages));
} catch (Exception e) {
log.warn("处理邮件失败 folder={} index={} error={}", folder.getFullName(), i, e.getMessage(), e);
}
}
} catch (Exception e) {
log.error("批量任务异常", e);
} finally {
latch.countDown();
}
});
}
folder.close(false);
try {
latch.await();
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
try {
imapFolder.close(false);
} catch (MessagingException e) {
log.warn("关闭文件夹异常", e);
}
}
private Store connect(UserAccount account) {
private EmailFetchData processSingleEmail(Message message,EmailFolder emailFolder){
EmailFetchData fetchData = new EmailFetchData();
return fetchData;
}
private EmailSummary processSingleEmailSummary(Folder folder,Message message,Long folderId) throws Exception {
EmailSummary emailSummary = new EmailSummary();
// subject
String subject = message.getSubject();
// messageId
String[] messageIdHeader = message.getHeader("Message-ID");
String messageId = (messageIdHeader != null && messageIdHeader.length > 0) ? messageIdHeader[0] : null;
// uid
UIDFolder uidFolder = (UIDFolder) folder;
Long uid = uidFolder.getUID(message);
// sentDate
Date sentDate = resolveMailSentDate(message);
// isSeen
Integer isSeen = message.isSet(Flags.Flag.SEEN) ? 1 : 0;
// emailId
long emailId = IdWorker.getId();
// 解析内容
Part part = message;
// bodyCharset
String contentType = part.getContentType();
Matcher matcher = Pattern.compile("charset=\\s*\"?([^;\"\\s]+)", Pattern.CASE_INSENSITIVE)
.matcher(contentType);
String bodyCharset = matcher.find() ? matcher.group(1).trim().toLowerCase()
: "UTF-8";
// bodyType
String bodyType = EmailBodyTypeEnum.fromMimeType(contentType);
// bodyEncoding
String[] encs = part.getHeader("Content-Transfer-Encoding");
String bodyEncoding = encs != null && encs.length > 0
? encs[0].trim().toLowerCase()
: "";
// bodyContent
String bodyContent = resolveMailContent(part);
return emailSummary;
}
private List<EmailAddress> processSingleEmailAddress(Message message) throws MessagingException {
List<EmailAddress> addressList = new ArrayList<>();
// 发件人
Address[] froms = message.getFrom();
// 回复地址
Address[] replyTos = message.getReplyTo();
// TO
Address[] tos = message.getRecipients(Message.RecipientType.TO);
// CC
Address[] ccs = message.getRecipients(Message.RecipientType.CC);
// BCC
Address[] bccs = message.getRecipients(Message.RecipientType.BCC);
return addressList;
}
private List<EmailFile> processSingleEmailFile(){
List<EmailFile> fileList = new ArrayList<>();
return fileList;
}
private String resolveMailContent(Part part) throws Exception {
if (part.isMimeType("text/html")) {
return (String) part.getContent();
} else if (part.isMimeType("text/plain")) {
return (String) part.getContent();
} else if (part.isMimeType("multipart/*")) {
Multipart multipart = (Multipart) part.getContent();
String html = null;
String plain = null;
for (int i = 0; i < multipart.getCount(); i++) {
BodyPart bodyPart = multipart.getBodyPart(i);
// ⛔ 跳过附件部分
if (Part.ATTACHMENT.equalsIgnoreCase(bodyPart.getDisposition())) {
continue;
}
String content = resolveMailContent(bodyPart);
if (bodyPart.isMimeType("text/html") && html == null) {
html = content;
} else if (bodyPart.isMimeType("text/plain") && plain == null) {
plain = content;
}
}
return html != null ? html : plain;
}
return "";
}
/**
* 解析邮件的 发送/接受 时间
* @param message 消息
* @return 发送/接受 时间
*/
private Date resolveMailSentDate(Message message) {
try {
if (message.getSentDate() != null) {
return message.getSentDate();
}
if (message.getReceivedDate() != null) {
return message.getReceivedDate();
}
String[] dateHeaders = message.getHeader("Date");
if (dateHeaders != null && dateHeaders.length > 0) {
try {
return new MailDateFormat().parse(dateHeaders[0]);
} catch (ParseException ignored) {
// 无法解析 Header 日期
}
}
} catch (Exception ignored) {
// 捕获所有异常以保证邮件拉取不中断
}
// 所有方式都失败,使用当前时间兜底
return new Date();
}
private List<EmailFolder> listAllFolders(Folder folder, Long parentId, Long userId) throws MessagingException {
List<EmailFolder> result = new ArrayList<>();
Folder[] folders = folder.list();
for (Folder f : folders) {
long id = IdWorker.getId();
EmailFolder emailFolder = new EmailFolder();
emailFolder.setId(id);
emailFolder.setUserId(userId);
emailFolder.setName(f.getName());
emailFolder.setFullName(f.getFullName());
emailFolder.setParentId(parentId);
emailFolder.setType(f.getType());
emailFolder.setCreateTime(new Date());
emailFolder.setUpdateTime(new Date());
result.add(emailFolder);
if ((f.getType() & Folder.HOLDS_FOLDERS) != 0) {
result.addAll(listAllFolders(f, id, userId));
}
}
return result;
}
private IMAPStore connect(UserAccount account) {
Properties props = new Properties();
props.put("mail.store.protocol", "imaps");
props.put("mail.imap.ssl.enable", "true");
props.put("mail.imap.connectiontimeout", "10000");
props.put("mail.imap.timeout", "10000");
// 在连接参数中强制指定ID信息
props.put("mail.imap.id.mechanism", "OAUTH2");
props.put("mail.imap.id.name", "name");
props.put("mail.imap.id.version", "1.0.0");
props.put("mail.imap.fetchsize", "4194304");
HashMap<String, String> imapId = new HashMap<>();
imapId.put("name", "name");
imapId.put("version", "1.0.0");
imapId.put("vendor", "vendor");
imapId.put("support-email", account.getEmail());
try {
Session session = Session.getInstance(props);
Store store = session.getStore("imaps");
store.connect(
account.getImapHost(),
Integer.parseInt(account.getImapPort()),
account.getEmail(),
account.getPassword()
);
IMAPStore store = (IMAPStore) session.getStore("imap");
store.connect(account.getImapHost(), Integer.parseInt(account.getImapPort()), account.getEmail(), account.getPassword());
store.id(imapId);
return store;
} catch (AuthenticationFailedException e) {
log.warn("邮箱认证失败:账号={},原因={}", account.getEmail(), e.getMessage(), e);
@@ -146,13 +365,7 @@ public class ImapJavaMailEmailFetcher implements EmailFetcher {
throw new AuthException(ResponseEnum.AUTH_EMAIL_PROTOCOL_UNSUPPORTED);
} catch (MessagingException e) {
log.warn("邮箱服务器连接失败:账号={},Host={},Port={},SSL={},原因={}",
account.getEmail(),
account.getImapHost(),
account.getImapPort(),
account.getImapSsl(),
e.getMessage(),
e
);
account.getEmail(), account.getImapHost(), account.getImapPort(), account.getImapSsl(), e.getMessage(), e);
throw new AuthException(ResponseEnum.AUTH_EMAIL_CONNECTION_FAILED);
} catch (Exception e) {
log.error("邮箱连接发生未知错误:账号={},原因={}", account.getEmail(), e.getMessage(), e);
@@ -160,4 +373,8 @@ public class ImapJavaMailEmailFetcher implements EmailFetcher {
}
}
private void shutdownExecutors() {
folderExecutor.shutdownNow();
emailExecutor.shutdownNow();
}
}
@@ -0,0 +1,16 @@
package com.email.core.fetch.model;
import com.email.domain.entity.EmailFolder;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import javax.mail.Folder;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class FolderWrapper {
private EmailFolder emailFolder;
private Folder folder;
}
@@ -44,6 +44,7 @@ public class EmailFileUtils {
InputStream inputStream = new ByteArrayInputStream(
bodyContent.getBytes(java.nio.charset.Charset.forName(charset))
);
return doUpload(email, "body.html", inputStream, BODY_PREFIX, "html");
}
@@ -89,9 +89,9 @@ const handleClickMenuItem = (key: string) => {
const handleClickSubMenuItem = (key: string) => {
const index = openkeys.value.indexOf(key);
if (index === -1) {
openkeys.value.push(key); // 如果不包含就添加
openkeys.value.push(key);
} else {
openkeys.value.splice(index, 1); // 如果包含就移除
openkeys.value.splice(index, 1);
}
};
</script>