u
This commit is contained in:
1 parent
9454e48bab
commit
a96f04630e
5 files changed
+233
-215
No files matched your search
@@ -1,82 +0,0 @@
|
||||
package com.email;
|
||||
|
||||
|
||||
import javax.mail.*;
|
||||
import javax.mail.internet.MimeBodyPart;
|
||||
import javax.mail.search.FlagTerm;
|
||||
import java.io.*;
|
||||
import java.util.*;
|
||||
import java.util.concurrent.*;
|
||||
|
||||
public class FolderHierarchyTest {
|
||||
private final ExecutorService executor = Executors.newFixedThreadPool(8);
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
new FolderHierarchyTest().run();
|
||||
}
|
||||
|
||||
public void run() throws Exception {
|
||||
Properties props = new Properties();
|
||||
props.put("mail.store.protocol", "imaps");
|
||||
|
||||
Session session = Session.getInstance(props, null);
|
||||
Store store = session.getStore();
|
||||
store.connect("imap.qq.com", "717406575@qq.com", "dhqesrwzgblrbcac");
|
||||
|
||||
Folder inbox = store.getFolder("INBOX");
|
||||
inbox.open(Folder.READ_ONLY);
|
||||
|
||||
if (!(inbox instanceof UIDFolder)) {
|
||||
throw new RuntimeException("This IMAP server does not support UIDFolder");
|
||||
}
|
||||
UIDFolder uidFolder = (UIDFolder) inbox;
|
||||
|
||||
Message[] messages = inbox.search(new FlagTerm(new Flags(Flags.Flag.SEEN), false));
|
||||
|
||||
int batchSize = 50;
|
||||
for (int i = 0; i < messages.length; i += batchSize) {
|
||||
int start = i;
|
||||
int end = Math.min(i + batchSize, messages.length);
|
||||
executor.submit(() -> processBatch(Arrays.copyOfRange(messages, start, end), uidFolder));
|
||||
}
|
||||
|
||||
executor.shutdown();
|
||||
executor.awaitTermination(30, TimeUnit.MINUTES);
|
||||
inbox.close(false);
|
||||
store.close();
|
||||
}
|
||||
|
||||
private void processBatch(Message[] batch, UIDFolder uidFolder) {
|
||||
for (Message msg : batch) {
|
||||
try {
|
||||
long uid = uidFolder.getUID(msg);
|
||||
String subject = msg.getSubject();
|
||||
Address[] froms = msg.getFrom();
|
||||
System.out.println("UID: " + uid + " Subject: " + subject);
|
||||
|
||||
Object content = msg.getContent();
|
||||
if (content instanceof Multipart) {
|
||||
Multipart mp = (Multipart) content;
|
||||
for (int i = 0; i < mp.getCount(); i++) {
|
||||
BodyPart bp = mp.getBodyPart(i);
|
||||
if (Part.ATTACHMENT.equalsIgnoreCase(bp.getDisposition())) {
|
||||
saveAttachment((MimeBodyPart) bp, uid);
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (Exception e) {
|
||||
System.err.println("Failed to process message: " + e.getMessage());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void saveAttachment(MimeBodyPart part, long uid) throws Exception {
|
||||
String fileName = part.getFileName();
|
||||
File dir = new File("attachments/" + uid);
|
||||
if (!dir.exists()) dir.mkdirs();
|
||||
|
||||
File file = new File(dir, fileName);
|
||||
part.saveFile(file);
|
||||
System.out.println("Saved attachment: " + file.getAbsolutePath());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,227 @@
|
||||
package com.email;
|
||||
|
||||
import javax.mail.*;
|
||||
import javax.mail.internet.MimeMessage;
|
||||
import javax.mail.internet.MimeUtility;
|
||||
import com.sun.mail.imap.IMAPFolder;
|
||||
import com.sun.mail.imap.IMAPStore;
|
||||
|
||||
import java.io.*;
|
||||
import java.util.*;
|
||||
import java.util.concurrent.*;
|
||||
|
||||
public class MultiFolderMailFetcher {
|
||||
|
||||
private static final String HOST = "imap.qq.com";
|
||||
private static final String USERNAME = "717406575@qq.com";
|
||||
private static final String PASSWORD = "dhqesrwzgblrbcac";
|
||||
|
||||
private static final int THREAD_COUNT = 20;
|
||||
private static final int BATCH_SIZE = 100;
|
||||
|
||||
private static final BlockingQueue<AttachmentTask> attachmentQueue = new LinkedBlockingQueue<>();
|
||||
|
||||
static class AttachmentTask {
|
||||
byte[] data;
|
||||
String filename;
|
||||
long uid;
|
||||
String folderName;
|
||||
|
||||
public AttachmentTask(byte[] data, String filename, long uid, String folderName) {
|
||||
this.data = data;
|
||||
this.filename = filename;
|
||||
this.uid = uid;
|
||||
this.folderName = folderName;
|
||||
}
|
||||
}
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
startAttachmentSavers(4);
|
||||
|
||||
long startTime = System.currentTimeMillis();
|
||||
|
||||
Properties props = new Properties();
|
||||
props.put("mail.store.protocol", "imap");
|
||||
props.put("mail.imap.ssl.enable", "true");
|
||||
props.put("mail.imap.fetchsize", "4194304");
|
||||
|
||||
Session session = Session.getInstance(props);
|
||||
IMAPStore store = (IMAPStore) session.getStore("imap");
|
||||
store.connect(HOST, USERNAME, PASSWORD);
|
||||
|
||||
Folder defaultFolder = store.getDefaultFolder();
|
||||
Folder[] folders = defaultFolder.list("*");
|
||||
|
||||
ExecutorService executor = Executors.newFixedThreadPool(THREAD_COUNT);
|
||||
List<Future<?>> futures = new ArrayList<>();
|
||||
|
||||
for (Folder f : folders) {
|
||||
if (!(f instanceof IMAPFolder)) continue;
|
||||
IMAPFolder folder = (IMAPFolder) f;
|
||||
|
||||
try {
|
||||
folder.open(Folder.READ_ONLY);
|
||||
int messageCount = folder.getMessageCount();
|
||||
if (messageCount == 0) {
|
||||
folder.close(false);
|
||||
continue;
|
||||
}
|
||||
|
||||
long startUID = folder.getUID(folder.getMessage(1));
|
||||
long endUID = folder.getUID(folder.getMessage(messageCount));
|
||||
folder.close(false);
|
||||
|
||||
System.out.printf("文件夹:%s,共 %d 封邮件,UID %d ~ %d%n",
|
||||
folder.getFullName(), messageCount, startUID, endUID);
|
||||
|
||||
List<long[]> segments = splitUIDRange(startUID, endUID, BATCH_SIZE);
|
||||
|
||||
for (long[] seg : segments) {
|
||||
long s = seg[0], e = seg[1];
|
||||
String folderName = folder.getFullName();
|
||||
|
||||
futures.add(executor.submit(() -> {
|
||||
try {
|
||||
fetchSegment(session, folderName, s, e);
|
||||
} catch (Exception ex) {
|
||||
ex.printStackTrace();
|
||||
}
|
||||
}));
|
||||
}
|
||||
|
||||
} catch (Exception ex) {
|
||||
System.err.println("无法处理文件夹: " + folder.getFullName() + " 错误: " + ex.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
// 等待所有线程结束
|
||||
for (Future<?> future : futures) {
|
||||
future.get();
|
||||
}
|
||||
|
||||
executor.shutdown();
|
||||
|
||||
long endTime = System.currentTimeMillis();
|
||||
System.out.println("总耗时:" + (endTime - startTime) + "ms");
|
||||
|
||||
store.close();
|
||||
}
|
||||
|
||||
private static void fetchSegment(Session session, String folderName, long startUID, long endUID) throws Exception {
|
||||
IMAPStore store = (IMAPStore) session.getStore("imap");
|
||||
store.connect(HOST, USERNAME, PASSWORD);
|
||||
IMAPFolder folder = (IMAPFolder) store.getFolder(folderName);
|
||||
folder.open(Folder.READ_ONLY);
|
||||
|
||||
try {
|
||||
Message[] messages = folder.getMessagesByUID(startUID, endUID);
|
||||
if (messages == null || messages.length == 0) return;
|
||||
|
||||
FetchProfile fp = new FetchProfile();
|
||||
fp.add(FetchProfile.Item.ENVELOPE);
|
||||
fp.add(FetchProfile.Item.FLAGS);
|
||||
fp.add(FetchProfile.Item.CONTENT_INFO);
|
||||
folder.fetch(messages, fp);
|
||||
|
||||
for (Message msg : messages) {
|
||||
if (msg == null) continue;
|
||||
MimeMessage mime = (MimeMessage) msg;
|
||||
long uid = folder.getUID(msg);
|
||||
Address[] froms = mime.getFrom();
|
||||
String from = (froms != null && froms.length > 0) ? froms[0].toString() : "未知";
|
||||
|
||||
System.out.println("线程 " + Thread.currentThread().getName()
|
||||
+ " - 文件夹: " + folderName
|
||||
+ " - UID: " + uid
|
||||
+ " - 发件人: " + from
|
||||
+ " - 主题: " + mime.getSubject());
|
||||
|
||||
Object content = mime.getContent();
|
||||
if (content instanceof String) {
|
||||
String body = ((String) content).trim();
|
||||
System.out.println("正文摘要: " + body.substring(0, Math.min(50, body.length())).replaceAll("[\r\n]+", " "));
|
||||
} else if (content instanceof Multipart) {
|
||||
parseMultipart((Multipart) content, uid, folderName);
|
||||
} else {
|
||||
System.out.println("未知正文类型: " + content.getClass().getName());
|
||||
}
|
||||
}
|
||||
|
||||
} finally {
|
||||
if (folder.isOpen()) folder.close(false);
|
||||
store.close();
|
||||
}
|
||||
}
|
||||
|
||||
private static void parseMultipart(Multipart multipart, long uid, String folderName) throws Exception {
|
||||
for (int i = 0; i < multipart.getCount(); i++) {
|
||||
BodyPart part = multipart.getBodyPart(i);
|
||||
String contentType = part.getContentType();
|
||||
|
||||
if (Part.ATTACHMENT.equalsIgnoreCase(part.getDisposition())
|
||||
|| (part.getFileName() != null && !part.getFileName().isEmpty())) {
|
||||
String filename = MimeUtility.decodeText(part.getFileName());
|
||||
InputStream is = part.getInputStream();
|
||||
byte[] data = inputStreamToByteArray(is);
|
||||
attachmentQueue.offer(new AttachmentTask(data, filename, uid, folderName));
|
||||
} else if (contentType.contains("text/plain") || contentType.contains("text/html")) {
|
||||
String text = (String) part.getContent();
|
||||
System.out.println("正文片段: " + text.substring(0, Math.min(50, text.length())).replaceAll("[\r\n]+", " "));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static byte[] inputStreamToByteArray(InputStream is) throws IOException {
|
||||
try (ByteArrayOutputStream baos = new ByteArrayOutputStream()) {
|
||||
byte[] buffer = new byte[4096];
|
||||
int len;
|
||||
while ((len = is.read(buffer)) != -1) {
|
||||
baos.write(buffer, 0, len);
|
||||
}
|
||||
return baos.toByteArray();
|
||||
}
|
||||
}
|
||||
|
||||
private static void startAttachmentSavers(int threadCount) {
|
||||
for (int i = 0; i < threadCount; i++) {
|
||||
Thread saver = new Thread(() -> {
|
||||
while (true) {
|
||||
try {
|
||||
AttachmentTask task = attachmentQueue.take();
|
||||
saveAttachmentSync(task.data, task.filename, task.uid, task.folderName);
|
||||
} catch (Exception e) {
|
||||
System.err.println("保存任务失败: " + e.getMessage());
|
||||
}
|
||||
}
|
||||
}, "Attachment-Saver-" + i);
|
||||
saver.setDaemon(true);
|
||||
saver.start();
|
||||
}
|
||||
}
|
||||
|
||||
private static void saveAttachmentSync(byte[] data, String filename, long uid, String folderName) {
|
||||
try {
|
||||
String safeFolder = folderName.replaceAll("[^a-zA-Z0-9_\\-]", "_");
|
||||
String dirPath = "./attachments/" + safeFolder + "/uid_" + uid;
|
||||
File dir = new File(dirPath);
|
||||
if (!dir.exists()) dir.mkdirs();
|
||||
|
||||
File file = new File(dir, filename);
|
||||
try (FileOutputStream fos = new FileOutputStream(file)) {
|
||||
fos.write(data);
|
||||
}
|
||||
System.out.println("保存附件: " + file.getAbsolutePath());
|
||||
} catch (Exception e) {
|
||||
System.err.println("保存附件失败: " + filename + " 错误: " + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
private static List<long[]> splitUIDRange(long startUID, long endUID, int batchSize) {
|
||||
List<long[]> segments = new ArrayList<>();
|
||||
for (long i = startUID; i <= endUID; i += batchSize) {
|
||||
long end = Math.min(i + batchSize - 1, endUID);
|
||||
segments.add(new long[]{i, end});
|
||||
}
|
||||
return segments;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
package com.email;
|
||||
|
||||
|
||||
public class Test {
|
||||
|
||||
}
|
||||
@@ -1,133 +0,0 @@
|
||||
package com.email;
|
||||
|
||||
import javax.mail.*;
|
||||
import javax.mail.internet.InternetAddress;
|
||||
import java.util.*;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
public class UltraMailFetcher {
|
||||
private static int DYNAMIC_PAGE_SIZE = 200;
|
||||
private final ExecutorService pipelineExecutor = Executors.newWorkStealingPool(8);
|
||||
private final Session session = createTurboSession();
|
||||
|
||||
private Session createTurboSession() {
|
||||
Properties props = new Properties();
|
||||
props.put("mail.store.protocol", "imap");
|
||||
props.put("mail.imap.host", "imap.qq.com");
|
||||
props.put("mail.imap.port", "993");
|
||||
props.put("mail.imap.ssl.enable", "true");
|
||||
|
||||
// Turbo优化参数
|
||||
props.put("mail.imap.connectionpoolsize", "6");
|
||||
props.put("mail.imap.connectionpooltimeout", "300000");
|
||||
props.put("mail.imap.ssl.sessioncache.size", "32");
|
||||
props.put("mail.imap.timeout", "15000");
|
||||
props.put("mail.imap.fetchsize", "1048576"); // 1MB缓冲区
|
||||
|
||||
return Session.getInstance(props);
|
||||
}
|
||||
|
||||
public void turboFetch() throws Exception {
|
||||
long totalStart = System.currentTimeMillis();
|
||||
|
||||
try (Store store = session.getStore();
|
||||
Folder inbox = store.getFolder("INBOX")) {
|
||||
store.connect("717406575@qq.com", "dhqesrwzgblrbcac");
|
||||
|
||||
inbox.open(Folder.READ_ONLY);
|
||||
|
||||
int total = inbox.getMessageCount();
|
||||
List<Message> allMessages = new ArrayList<>(total);
|
||||
|
||||
// 阶段1:智能分页加载
|
||||
CompletableFuture<Void> fetchFuture = CompletableFuture.runAsync(() -> {
|
||||
int remaining = total;
|
||||
int currentStart = 1;
|
||||
while (remaining > 0) {
|
||||
int currentPage = Math.min(DYNAMIC_PAGE_SIZE, remaining);
|
||||
int currentEnd = currentStart + currentPage - 1;
|
||||
|
||||
try {
|
||||
long pageStart = System.currentTimeMillis();
|
||||
Message[] batch = inbox.getMessages(currentStart, currentEnd);
|
||||
|
||||
FetchProfile fp = new FetchProfile();
|
||||
fp.add(FetchProfile.Item.ENVELOPE);
|
||||
fp.add(FetchProfile.Item.CONTENT_INFO);
|
||||
inbox.fetch(batch, fp);
|
||||
|
||||
synchronized (allMessages) {
|
||||
Collections.addAll(allMessages, batch);
|
||||
}
|
||||
|
||||
// 动态调整分页
|
||||
long loadTime = System.currentTimeMillis() - pageStart;
|
||||
DYNAMIC_PAGE_SIZE = adjustPageSize(DYNAMIC_PAGE_SIZE, loadTime);
|
||||
|
||||
remaining -= currentPage;
|
||||
currentStart = currentEnd + 1;
|
||||
} catch (MessagingException e) {
|
||||
System.err.println("Page load error: " + e.getMessage());
|
||||
}
|
||||
}
|
||||
}, pipelineExecutor);
|
||||
|
||||
// 阶段2:流水线处理
|
||||
CompletableFuture<Void> processFuture = fetchFuture.thenApplyAsync(v -> {
|
||||
processTurboMessages(allMessages);
|
||||
return null;
|
||||
}, pipelineExecutor);
|
||||
|
||||
processFuture.join();
|
||||
|
||||
System.out.printf("[Turbo] Total time: %dms%n",
|
||||
System.currentTimeMillis() - totalStart);
|
||||
}
|
||||
}
|
||||
|
||||
private int adjustPageSize(int currentSize, long loadTime) {
|
||||
if (loadTime < 800) return (int)(currentSize * 1.5);
|
||||
if (loadTime > 2000) return (int)(currentSize * 0.6);
|
||||
return currentSize;
|
||||
}
|
||||
|
||||
private void processTurboMessages(List<Message> messages) {
|
||||
final int cores = Runtime.getRuntime().availableProcessors();
|
||||
final int batchSize = Math.max(50, messages.size() / (cores * 2));
|
||||
|
||||
List<CompletableFuture<Void>> futures = new ArrayList<>();
|
||||
|
||||
for (int i = 0; i < messages.size(); i += batchSize) {
|
||||
int end = Math.min(i + batchSize, messages.size());
|
||||
List<Message> batch = messages.subList(i, end);
|
||||
|
||||
futures.add(CompletableFuture.runAsync(() -> {
|
||||
batch.parallelStream().forEach(msg -> {
|
||||
try {
|
||||
String from = extractAddress(msg.getFrom());
|
||||
String subject = msg.getSubject();
|
||||
// 实际处理逻辑...
|
||||
} catch (MessagingException e) {
|
||||
System.err.println("Process error: " + e.getMessage());
|
||||
}
|
||||
});
|
||||
}, pipelineExecutor));
|
||||
}
|
||||
|
||||
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
|
||||
}
|
||||
|
||||
private String extractAddress(Address[] addresses) {
|
||||
// ...地址提取逻辑...
|
||||
return "";
|
||||
}
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
UltraMailFetcher fetcher = new UltraMailFetcher();
|
||||
fetcher.turboFetch();
|
||||
fetcher.pipelineExecutor.shutdownNow();
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user