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

十万条邮件列表并发发送问题:程序无报错提前退出

问题分析与解决方案

核心问题

程序无报错退出的根本原因是主线程提前终止:

  • 初始调用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 12:24:56