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

Python读取大文件避免RAM耗尽 ThreadPoolExecutor并发处理问题

核心问题

你当前内存占用过高的直接原因:

  • 列表推导式[line.strip() for line in stream]会一次性把全文件500万行全部加载到内存生成列表,加上Python字符串对象的额外内存开销,大文件下很容易占满RAM。
  • 即便把列表换成生成器传给pool.map,ThreadPoolExecutor默认会预取批量任务存入内部待调度队列,几百万个Future对象同样会占用大量内存。
  • 现有代码存在逻辑错误:在提交任何爬取任务前就调用queue.join(),程序会直接阻塞在这一步无法继续执行。
实现方案

核心逻辑是惰性逐行读文件+严格控制在途任务数量,保证内存中同时存在的URL、待处理任务、待写入结果永远维持在固定小数量级,内存占用不会随文件行数增长而升高。
实现要点:

  • 直接迭代打开的文件对象:Python的文件对象本身就是惰性迭代器,每次迭代仅读取一行到内存,处理完立即回收,不会加载全量文件。
  • 用信号量限制最大在途任务数:信号量计数和线程池最大线程数保持一致,确保线程池内部待调度队列不会堆积过量任务。
  • 修正线程等待逻辑:所有爬取任务执行完成后,再等待写入队列清空,最后通知写入线程退出,避免数据丢失。

修正后完整代码

import requests
from concurrent.futures import ThreadPoolExecutor
from bs4 import BeautifulSoup
import warnings
from threading import Thread, Semaphore, Lock
from queue import Queue

# 关闭无关警告
requests.packages.urllib3.disable_warnings(requests.packages.urllib3.exceptions.InsecureRequestWarning)
warnings.filterwarnings("ignore", category=UserWarning, module='bs4')

# 可调参数
MAX_WORKERS = 50
INPUT_FILE = "1_1.txt"
OUTPUT_FILE = "output.txt"
REQUEST_TIMEOUT = 40

# 全局共享组件
task_queue = Queue()
sem = Semaphore(MAX_WORKERS)
count_lock = Lock()
success_count = 0
fail_count = 0

def get_url(url):
    global success_count, fail_count
    try:
        resp = requests.get(url, verify=False, timeout=REQUEST_TIMEOUT)
        soup = BeautifulSoup(resp.text, 'html.parser')
        title = str(soup.title.get_text().splitlines(False))[:10000]
        with count_lock:
            success_count += 1
        task_queue.put(f'{url} - {title} \n')
    except Exception:
        with count_lock:
            fail_count += 1
        task_queue.put(f'FAILED : {url} \n')
    finally:
        # 任务完成释放信号量,允许提交新任务
        sem.release()

def file_writer(filepath, q):
    with open(filepath, 'a', encoding="utf-8") as f:
        while True:
            line = q.get()
            if line is None:
                q.task_done()
                break
            f.write(line)
            # 每100条刷一次盘,平衡性能和落盘实时性,可自行调整阈值
            if success_count % 100 == 0:
                f.flush()
            q.task_done()

if __name__ == "__main__":
    # 启动独立写入线程
    writer_thread = Thread(target=file_writer, args=(OUTPUT_FILE, task_queue), daemon=True)
    writer_thread.start()

    # 逐行读文件提交任务,全程不加载全量文件到内存
    with open(INPUT_FILE, "r", encoding="utf-8") as f, ThreadPoolExecutor(max_workers=MAX_WORKERS) as pool:
        for line in f:
            url = line.strip()
            if not url:
                continue
            # 信号量满时自动阻塞,避免任务队列堆积
            sem.acquire()
            pool.submit(get_url, url)
    
    # 等待所有爬取结果写入完成
    task_queue.join()
    # 发送终止信号,等待写入线程退出
    task_queue.put(None)
    writer_thread.join()

    print(f"处理完成,成功请求:{success_count},失败请求:{fail_count}")

效果说明

  • 内存占用稳定:运行时内存中最多同时存在50个正在处理的请求、少量待写入结果,哪怕文件行数上亿,内存占用也会稳定在几十MB级别,不会随文件大小上涨。
  • 顺序偏差可控:同一时间最多50个请求并发处理,结果顺序偏差不会超过50行,完全符合几十行偏差可接受的要求。
  • 无数据丢失:所有任务完成、队列内容全部落盘后才会终止程序,避免中途退出导致结果丢失。

可选优化点

  • 可以初始化全局requests.Session()对象复用TCP连接,请求速度会提升30%以上。
  • 如果不需要严格实时落盘,可以把file.flush()的触发阈值调高(比如每1000行刷一次),能明显降低磁盘IO开销。
  • 可以给请求加随机UA、代理轮换逻辑,减少被目标站点封禁的概率。

内容的提问来源于stack exchange,提问作者jdm0

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 03:06:12