如何实现多线程(单邮箱对应单线程)SMTP队列发件及线程负载均衡
解决方案:均衡Office365邮箱邮件发送负载
问题根源
现有方案中,队列空时各线程独立休眠5秒,醒来后随机抢占队列中的邮件,导致不同线程(对应不同邮箱)处理的邮件数量出现差异。同时手动轮询队列的方式效率较低,还容易引发竞争问题。
优化方案一:使用BlockingCollection替代ConcurrentQueue
BlockingCollection封装了ConcurrentQueue,自带阻塞等待功能,能在队列空时自动挂起线程,有新邮件时唤醒线程,避免手动轮询和不同步休眠的问题。
修改后的核心代码
using Message = KeyValuePair<long, Email>; // 用BlockingCollection包装ConcurrentQueue,实现阻塞取元素 static BlockingCollection<Message> emails_queue = new BlockingCollection<Message>(new ConcurrentQueue<Message>()); static Thread GetDBEmailsThread;
// 每个线程绑定固定邮箱,参数传入对应账号 static void SendEmails(string emailAccount) { while (!emails_queue.IsCompleted) { try { // 队列空时自动阻塞,直到有新邮件或队列标记为完成 Message msg = emails_queue.Take(); Email email = msg.Value; Console.WriteLine($"{DateTime.Now} {Thread.CurrentThread.Name} - 正在处理: {msg.Key} - {email.Subject}"); // 发送邮件到指定邮箱 SendEmail(email, emailAccount); Console.WriteLine($"{DateTime.Now} {Thread.CurrentThread.Name} - 邮件 {msg.Key} 已发送."); } catch (InvalidOperationException) { // 队列已完成,退出循环 break; } catch (Exception e) { Console.WriteLine($"{DateTime.Now} {Thread.CurrentThread.Name} - 异常. \r\n{e.Message}"); } finally { // 控制单邮箱发送速率:每3秒1封,60秒最多20封,符合Office365限制 Thread.Sleep(3000); } } }
方案优势
- 自动阻塞等待:无需手动判断队列长度和休眠,线程在队列空时统一挂起,有新邮件时同步唤醒,避免了线程间的休眠不同步问题,保证各邮箱处理机会均等。
- 严格速率控制:每个线程绑定固定邮箱,通过
Thread.Sleep(3000)严格控制单邮箱发送频率,满足Office365的60秒最多30封的限制(此处设置为20封更保守)。 - 线程安全:基于
ConcurrentQueue实现,天然支持多线程并发操作,无需额外的锁控制。
优化方案二:使用Channel实现异步发送(推荐.NET Core/.NET 5+)
对于现代.NET应用,使用Channel(异步优先的队列)替代传统线程和阻塞集合,能获得更高的线程利用率和更简洁的代码。
核心代码示例
using Message = KeyValuePair<long, Email>; // 创建有界通道,避免内存溢出;满队列时等待写入 static Channel<Message> emails_channel = Channel.CreateBounded<Message>(new BoundedChannelOptions(1000) { FullMode = BoundedChannelFullMode.Wait, AllowSynchronousContinuations = false }); // 异步拉取数据库邮件到通道 static async Task FetchEmailsFromDbAsync() { while (true) { // 从数据库拉取待发送邮件 var emails = await GetEmailsFromDbAsync(); if (emails.Count == 0) { // 无邮件时异步休眠,不占用线程资源 await Task.Delay(5000); continue; } foreach (var email in emails) { await emails_channel.Writer.WriteAsync(email); } } } // 异步发送邮件任务,每个任务绑定固定邮箱 static async Task SendEmailsAsync(string emailAccount) { // 异步遍历通道中的邮件,空通道时自动等待 await foreach (var msg in emails_channel.Reader.ReadAllAsync()) { try { Email email = msg.Value; Console.WriteLine($"{DateTime.Now} 任务ID:{Task.CurrentId} - 正在处理: {msg.Key} - {email.Subject}"); // 使用异步发送邮件方法,避免阻塞线程 await SendEmailAsync(email, emailAccount); Console.WriteLine($"{DateTime.Now} 任务ID:{Task.CurrentId} - 邮件 {msg.Key} 已发送."); } catch (Exception e) { Console.WriteLine($"{DateTime.Now} 任务ID:{Task.CurrentId} - 异常. \r\n{e.Message}"); } finally { // 异步休眠,不占用线程池资源 await Task.Delay(3000); } } } // 启动程序入口 static void Main() { // 启动拉取邮件任务 _ = FetchEmailsFromDbAsync(); // 启动10个发送任务,每个对应一个邮箱 for (int i = 0; i < 10; i++) { var emailAccount = GetEmailAccount(i); // 获取第i个邮箱账号 _ = SendEmailsAsync(emailAccount); } Console.ReadLine(); }
方案优势
- 异步高效:
Task.Delay替代Thread.Sleep,不会阻塞线程池线程,资源利用率更高。 - 天然支持异步等待:
ReadAllAsync()自动处理空通道的情况,异步等待新邮件,无需轮询。 - 可扩展性强:通道支持配置边界、超时等参数,更容易应对高并发场景。
关键注意事项
- 每个线程/任务必须绑定固定邮箱,不能混用,否则会突破单个邮箱的发送限制。
- 数据库拉取线程的速度要匹配发送速度,避免队列/通道溢出或长期空队列。
- 可以添加失败重试机制,针对Office365的限流错误(如429)进行延迟重试,避免邮件丢失。
内容的提问来源于stack exchange,提问作者Bilal Halayqa
相关产品推荐
相关产品推荐

