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

如何为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)

问题分析

  1. 共享NNTP连接线程不安全:nntplib.NNTP实例不是线程安全的,多个线程共用同一个连接会导致命令执行混乱,出现重复获取同一消息的情况。
  2. 任务提交逻辑错误:仅调用executor.submit(download_nntp_file)但未传递消息ID参数,且函数内部逻辑混乱,没有正确处理传入的消息ID。
  3. 目录创建无并发保护:多个线程同时创建同一目录会抛出异常。
  4. 未遍历所有消息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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 11:41:05