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

Python AsyncIO、多线程与多进程写入大量小文件性能疑问

批量写入小文件的性能优化困惑与解决方案

我正尝试加速一段写入数千个小文件的代码,最初考虑使用AsyncIO(因文件写入为阻塞I/O操作),同时也计划测试multithreading(多线程)和multiprocessing(多进程)方案。

经过数次测试后,我对结果深感困惑:

  • 同步与异步算法的性能几乎无差异;
  • 性能会因首次调用的方法不同而产生巨大波动;
  • 即便限制了线程或进程数量,multithreading和multiprocessing的运行速度依然极慢。

请问我的基准测试是否存在问题?或者针对这类场景有更优的解决方案?

测试代码

import asyncio
import multiprocessing
import os
import random
import shutil
import threading
import time
from multiprocessing import Process, cpu_count


def compute_time(f):
    def wrapped_func(*args, **kwargs):
        start = time.time()
        result = f(*args, **kwargs)
        print(time.time() - start)
        return result

    return wrapped_func


def async_compute_time(f):
    async def wrapped_func(*args, **kwargs):
        start = time.time()
        result = await f(*args, **kwargs)
        print(time.time() - start)
        return result

    return wrapped_func


def clean_dir():
    try:
        shutil.rmtree("toto")
    except:
        pass
    os.mkdir("toto")


# =========================SYNC========================
def write_to_file(filename, content):
    with open(filename, "w") as f:
        f.write(content)

@compute_time
def write_sync(n):
    content = str(random.randint(0, 100))
    for i in range(n): 
        filename = f"toto/file_{i}.txt"
        write_to_file(filename, content)


# =========================ASYNC
import aiofiles


async def async_write_to_file(filename, content):
    async with aiofiles.open(filename, "w") as f:
        await f.write(content)


@async_compute_time
async def async_write_sync(n):
    MAX_CORO = 1600 #found this number to be otpimum "by hand"
    content = str(random.randint(0, 1000))

    for sub in range(0, n, MAX_CORO):
        task = []
        for i in range(sub, min(sub + MAX_CORO, n)):
            filename = f"toto/file_{i}.txt"
            task.append(async_write_to_file(filename, content))
        await asyncio.gather(*task)


# =====================Thread
@compute_time
def write_sync_multi_thread(n):
    MAX_THREAD = 8
    content = str(random.randint(0, 100))
    threads = []
    for sub in range(0, n, MAX_THREAD):
        for i in range(sub, min(sub + MAX_THREAD, n)):
            filename = f"toto/file_{i}.txt"
            t = threading.Thread(target=write_to_file, args=(filename, content))
            threads.append(t)
            t.start()
        for t in threads:
            t.join()


# Multiprocessing
@compute_time
def write_sync_multi_processing(n):
    NUMBER_OF_WORKER = 5
    content = str(random.randint(0, 100))
    processes = []
    for sub in range(0, n, 5):
        for i in range(sub, min(sub + NUMBER_OF_WORKER, n)):
            filename = f"toto/file_{i}.txt"
            p = multiprocessing.Process(target=write_to_file, args=(filename, content))
            processes.append(p)
            p.start()
        for p in processes:
            p.join()


if __name__ == "__main__":
    n = 10000
    # clean_dir()
    # write_sync_multi_thread(n)
    clean_dir()
    asyncio.run(async_write_sync(n))
    clean_dir()
    write_sync(n)

    # clean_dir()
    # write_sync_multi_processing(n)

问题分析与解决方案

一、基准测试的问题

  1. 异步并发过载
    设置MAX_CORO = 1600会瞬间创建大量协程写入请求,超出磁盘的并发处理能力。磁盘的寻道和调度资源有限,过多并发会引发队列拥堵,抵消异步的优势,导致和同步性能无差异。

  2. 多线程/多进程实现缺陷

    • 多线程代码中,threads列表是全局累积的,随着循环推进,实际运行的线程数会远超MAX_THREAD,引发线程调度开销爆炸。
    • 多进程的创建、销毁开销远大于线程,且文件I/O是磁盘瓶颈,多进程无法利用CPU并行优势,反而因进程间资源竞争拖慢速度。
  3. 冷启动影响测试结果
    首次调用方法时,操作系统文件缓存、磁盘调度器处于冷状态,后续调用受益于缓存,导致性能波动。

二、优化方案

  1. 限制异步协程并发数
    用asyncio.Semaphore将并发数控制在磁盘可承受范围(如10-50,根据硬件调整),避免I/O队列过载:

    async def async_write_sync(n):
        MAX_CONCURRENT = 20
        semaphore = asyncio.Semaphore(MAX_CONCURRENT)
        content = str(random.randint(0, 1000))
    
        async def bounded_write(filename):
            async with semaphore:
                await async_write_to_file(filename, content)
    
        tasks = [bounded_write(f"toto/file_{i}.txt") for i in range(n)]
        await asyncio.gather(*tasks)
    
  2. 修复多线程实现
    每个子块结束后清空线程列表,确保每次仅运行MAX_THREAD个线程:

    @compute_time
    def write_sync_multi_thread(n):
        MAX_THREAD = 8
        content = str(random.randint(0, 100))
        for sub in range(0, n, MAX_THREAD):
            threads = []
            for i in range(sub, min(sub + MAX_THREAD, n)):
                filename = f"toto/file_{i}.txt"
                t = threading.Thread(target=write_to_file, args=(filename, content))
                threads.append(t)
                t.start()
            for t in threads:
                t.join()
    
  3. 放弃多进程方案
    多进程适合CPU密集型任务,文件I/O密集场景下,进程开销远大于收益,优先选择异步或多线程即可。

  4. 测试前预热
    正式测试前先运行所有方法,让操作系统缓存预热,消除冷启动偏差:

    if __name__ == "__main__":
        n = 10000
        # 预热
        clean_dir()
        asyncio.run(async_write_sync(n))
        clean_dir()
        write_sync(n)
        clean_dir()
        write_sync_multi_thread(n)
        clean_dir()
    
        # 正式测试
        print("=== 正式测试 ===")
        clean_dir()
        asyncio.run(async_write_sync(n))
        clean_dir()
        write_sync(n)
        clean_dir()
        write_sync_multi_thread(n)
    
  5. 额外优化技巧

    • 分组存储文件:将文件按子目录分组(如每1000个文件一个子目录),减少根目录文件数量,提升文件系统访问效率。
    • 使用SSD或内存文件系统:SSD的随机写入性能远优于HDD,测试时可使用tmpfs排除硬件瓶颈。
    • 减少元数据操作:场景允许时复用文件句柄,避免频繁open/close操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 06:27:35