如何为NNTP归档程序实现多线程?(新手求助)
新闻组归档代码的多线程改造方案
问题描述
需要将单线程的新闻组文章归档代码改造为多线程版本,目标利用服务器支持的最大连接数(最多5线程)同时下载不同消息,但尝试使用ThreadPoolExecutor时出现重复处理同一消息ID的问题,无法实现预期的多线程并行下载。
原单线程代码
import nntplib import sys import datetime import os basetime = datetime.datetime.today() # daysback = int(sys.argv[1]) # date_list = [basetime - datetime.timedelta(days=x) for x in range(daysback)] # 建立NNTP连接,服务器限制最多5连接 s = nntplib.NNTP('free.xsusenet.com', user='USERNAME', password='PASSWORD') groups = [] resp, groups_list_tuple = s.list() def remove_non_ascii_2(string): return string.encode('ascii', errors='ignore').decode() for g_tuple in groups_list_tuple: # 解析新闻组信息 group = g_tuple[0] last = g_tuple[1] first = g_tuple[2] flag = g_tuple[3] resp, count, first, last, name = s.group(group) # 遍历当前组的所有消息ID并下载 for message_id in range(first, last): resp, number, mes_id = s.next() resp, info = s.article(mes_id) # 创建组目录 if not os.path.exists(f'.\\{group}'): os.mkdir(f'.\\{group}') print(f"Downloading: {message_id}") # 写入文件 with open(f'.\\{group}\\{message_id}', 'a', encoding="utf-8") as outfile: for line in info.lines: outfile.write(remove_non_ascii_2(str(line)) + '\n')
失败的多线程尝试代码
import nntplib import sys import datetime import os import concurrent.futures basetime = datetime.datetime.today() # daysback = int(sys.argv[1]) # date_list = [basetime - datetime.timedelta(days=x) for x in range(daysback)] s = nntplib.NNTP('free.xsusenet.com', user='USERNAME', password='PASSWORD') groups = [] resp, groups_list_tuple = s.list() def remove_non_ascii_2(string): return string.encode('ascii', errors='ignore').decode() def download_nntp_file(mess_id): resp, count, first, last, name = s.group(group) message_id = range(first, last) resp, number, mes_id = s.next() resp, info = s.article(mes_id) if os.path.exists('.\\' + group): pass else: os.mkdir('.\\' + group) print(f"Downloading: {mess_id}") outfile = open('.\\' + group + '\\' + str(mess_id), 'a', encoding="utf-8") for line in info.lines: outfile.write(remove_non_ascii_2(str(line)) + '\n') outfile.close() for g_tuple in groups_list_tuple: group = g_tuple[0] last = g_tuple[1] first = g_tuple[2] flag = g_tuple[3] with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor: futures = executor.submit(download_nntp_file)
问题分析
- 共享NNTP连接线程不安全:
nntplib.NNTP实例不是线程安全的,多个线程共用同一个连接会导致命令执行混乱,出现重复获取同一消息的情况。 - 任务提交逻辑错误:仅调用
executor.submit(download_nntp_file)但未传递消息ID参数,且函数内部逻辑混乱,没有正确处理传入的消息ID。 - 目录创建无并发保护:多个线程同时创建同一目录会抛出异常。
- 未遍历所有消息ID:没有将所有消息ID作为独立任务提交给线程池,仅执行了单次下载操作。
修正后的多线程代码
import nntplib import os import concurrent.futures # 配置信息 NNTP_SERVER = 'free.xsusenet.com' USERNAME = 'USERNAME' PASSWORD = 'PASSWORD' MAX_WORKERS = 5 # 匹配服务器限制的最大连接数 def remove_non_ascii_2(string): return string.encode('ascii', errors='ignore').decode() def download_message(group_name, message_id): # 每个线程创建独立的NNTP连接,避免线程安全问题 try: with nntplib.NNTP(NNTP_SERVER, user=USERNAME, password=PASSWORD) as s: # 切换到目标新闻组 s.group(group_name) # 获取指定ID的文章 resp, info = s.article(str(message_id)) # 创建组目录,exist_ok=True避免重复创建报错 os.makedirs(f'.\\{group_name}', exist_ok=True) # 写入文件 file_path = f'.\\{group_name}\\{message_id}' with open(file_path, 'w', encoding='utf-8') as outfile: for line in info.lines: outfile.write(remove_non_ascii_2(str(line)) + '\n') print(f"Completed: {group_name} - {message_id}") except Exception as e: print(f"Failed to download {group_name}-{message_id}: {str(e)}") def main(): # 先获取所有新闻组列表 with nntplib.NNTP(NNTP_SERVER, user=USERNAME, password=PASSWORD) as s: resp, groups_list_tuple = s.list() # 遍历每个新闻组,收集所有需要下载的任务 tasks = [] for g_tuple in groups_list_tuple: group_name = g_tuple[0] first_msg_id = int(g_tuple[2]) last_msg_id = int(g_tuple[1]) # 生成当前组的所有消息ID任务 for msg_id in range(first_msg_id, last_msg_id + 1): tasks.append( (group_name, msg_id) ) # 使用线程池执行所有任务 with concurrent.futures.ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: # 用map批量提交任务,自动解包参数 executor.map(lambda p: download_message(*p), tasks) if __name__ == "__main__": main()
关键改进点
- 每个线程独立连接:避免共享连接导致的线程安全问题,确保每个线程的NNTP操作互不干扰。
- 批量任务提交:收集所有消息ID作为任务,通过
executor.map批量提交,确保每个任务处理不同的消息。 - 安全创建目录:使用
os.makedirs并设置exist_ok=True,避免并发创建目录的异常。 - 异常处理:增加异常捕获,避免单个任务失败导致整个程序崩溃。
内容的提问来源于stack exchange,提问作者ANK
相关产品推荐
相关产品推荐

