十万条邮件列表并发发送问题:程序无报错提前退出
问题分析与解决方案
核心问题
程序无报错退出的根本原因是主线程提前终止:
- 初始调用
concurrentExecute时,主线程进入with ThreadPoolExecutor()块,提交50个任务后,会等待这50个任务全部完成才会退出该块。 - 当这50个任务全部执行完毕(无论成功还是重试完成),主线程就会结束,整个进程随之终止——哪怕子线程中还在调用
concurrentExecute处理后续邮件,也会被强制终止。
除此之外,代码还有两个隐患:
- 递归调用
sendEmail重试,若API持续失败,会触发栈溢出,最终导致程序崩溃。 - 每次创建新的
ThreadPoolExecutor,不仅有性能开销,还可能导致实际并发数超过50的限制。
修正后的代码
改用全局固定大小的线程池,通过循环处理重试,主线程等待所有任务完成:
import concurrent.futures import requests # 全局线程池,固定最大并发数50 MAX_WORKERS = 50 executor = concurrent.futures.ThreadPoolExecutor(max_workers=MAX_WORKERS) def send_email(email, send_mail_api): max_retries = 3 # 设置最大重试次数,避免无限循环 retry_count = 0 while retry_count < max_retries: try: res = requests.post(send_mail_api, data={"email": email}) if res.status_code == 200: return True print(f"发送邮件{email}失败,状态码{res.status_code},重试中...") except Exception as e: print(f"发送邮件{email}出现异常:{str(e)},重试中...") retry_count += 1 print(f"邮件{email}重试{max_retries}次后仍失败") return False def main(all_emails, send_mail_api): # 提交所有任务到线程池 futures = [executor.submit(send_email, email, send_mail_api) for email in all_emails] # 主线程等待所有任务完成 concurrent.futures.wait(futures) # 关闭线程池 executor.shutdown() # 调用示例 # all_emails = [...] # 你的10万条邮件列表 # send_mail_api = "https://your-send-mail-api.com" # main(all_emails, send_mail_api)
优化说明
- 固定线程池:全局只创建一个线程池,严格控制最大并发数为50,避免频繁创建销毁线程的开销。
- 循环重试:用while循环替代递归重试,避免栈溢出风险,同时设置最大重试次数,防止无限循环。
- 主线程等待:通过
concurrent.futures.wait()让主线程等待所有邮件任务处理完毕再退出,确保所有邮件都被处理。
如果想要更高效的“任务完成一个就补一个”的模式(避免一次性提交10万个任务占用内存),可以用以下方式:
import concurrent.futures import requests from itertools import islice MAX_WORKERS = 50 executor = concurrent.futures.ThreadPoolExecutor(max_workers=MAX_WORKERS) def send_email(email, send_mail_api): max_retries = 3 retry_count = 0 while retry_count < max_retries: try: res = requests.post(send_mail_api, data={"email": email}) if res.status_code == 200: return True print(f"发送邮件{email}失败,状态码{res.status_code},重试中...") except Exception as e: print(f"发送邮件{email}出现异常:{str(e)},重试中...") retry_count += 1 print(f"邮件{email}重试{max_retries}次后仍失败") return False def main(all_emails, send_mail_api): email_iter = iter(all_emails) # 先提交初始的MAX_WORKERS个任务 futures = {executor.submit(send_email, email, send_mail_api): email for email in islice(email_iter, MAX_WORKERS)} while futures: # 等待任意一个任务完成 done, _ = concurrent.futures.wait(futures, return_when=concurrent.futures.FIRST_COMPLETED) for future in done: # 移除已完成的任务 del futures[future] # 获取下一个邮件,提交新任务 try: next_email = next(email_iter) new_future = executor.submit(send_email, next_email, send_mail_api) futures[new_future] = next_email except StopIteration: # 邮件列表已遍历完 pass executor.shutdown() # 调用示例 # all_emails = [...] # send_mail_api = "https://your-send-mail-api.com" # main(all_emails, send_mail_api)
这个版本会在每个任务完成后立即补充新任务,保持并发数稳定在50,同时不会一次性创建10万个future对象,内存占用更低。
内容的提问来源于stack exchange,提问作者News Entertainment
相关产品推荐
相关产品推荐

