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

使用asyncio实现多文件并发写入的设计模式及实现方法

实现方案说明

你这个场景适用生产者-消费者设计模式,核心是把耗时的记录生成逻辑和IO密集的文件写入逻辑解耦,两者可以并行执行,不会因为文件IO阻塞记录生成流程。
首先要注意Python标准库的原生文件操作是同步阻塞的,asyncio生态下一般用第三方库aiofiles实现异步文件写入,如果你不想引入第三方依赖,也可以用asyncio提供的线程池包装同步写逻辑实现。

可运行实现代码(基于aiofiles)

首先先修复你原单线程代码的小问题:generate_records里的循环变量是i,你代码里误用了index变量,下面的实现已经修正:

import asyncio
import aiofiles
import datetime
import json
from asyncio import Queue

# 生产者协程:生成记录,写入队列
async def produce(queue: Queue):
    for i in range(100):
        # 模拟生成记录的耗时操作,可替换为实际业务逻辑
        await asyncio.sleep(0.01)
        record = {
            'index': i,
            'message': f"Message {i}",
            'timestamp': datetime.datetime.now().isoformat()
        }
        # 把记录和序号丢进队列
        await queue.put((i, record))
    # 生产完成,往队列丢结束标记,数量和消费者数量一致
    for _ in range(3):
        await queue.put(None)

# 消费者协程:从队列取数据,异步写文件
async def consume(queue: Queue):
    while True:
        item = await queue.get()
        if item is None:
            # 收到结束标记,退出循环
            queue.task_done()
            break
        idx, record = item
        # 异步写文件
        async with aiofiles.open(f"record_{idx:04}.json", 'w', encoding='utf-8') as f:
            await f.write(json.dumps(record, ensure_ascii=False, indent=2))
        queue.task_done()

async def main():
    # 可以设置队列最大长度,避免生成太快内存占用过高
    queue = Queue(maxsize=10)
    # 启动生产者
    producer_task = asyncio.create_task(produce(queue))
    # 启动3个消费者并发写文件,数量可以根据实际IO情况调整
    consumer_tasks = [asyncio.create_task(consume(queue)) for _ in range(3)]
    # 等生产者生产完成
    await producer_task
    # 等队列所有数据消费完成
    await queue.join()
    # 等所有消费者退出
    await asyncio.gather(*consumer_tasks)

if __name__ == "__main__":
    asyncio.run(main())

无第三方依赖的实现方案(用线程池包装同步写)

如果你不想安装aiofiles,可以用loop.run_in_executor把同步的写文件逻辑丢到线程池里跑,也能达到并发IO的效果,核心逻辑和上面一致,只需要修改消费者的写文件部分:

# 同步写文件函数
def write_to_file(idx, record):
    with open(f"record_{idx:04}.json", 'w', encoding='utf-8') as f:
        json.dump(record, f, ensure_ascii=False, indent=2)

# 修改后的消费者协程
async def consume(queue: Queue):
    loop = asyncio.get_running_loop()
    while True:
        item = await queue.get()
        if item is None:
            queue.task_done()
            break
        idx, record = item
        # 把同步写操作丢到默认线程池执行
        await loop.run_in_executor(None, write_to_file, idx, record)
        queue.task_done()

核心说明

  • 生产者消费者模式下,记录生成和文件写入完全并行,不会互相阻塞
  • 消费者数量可以根据你的磁盘IO性能调整,普通机械盘3-5个足够,SSD可以开更多
  • 队列设最大长度可以避免生成速度远快于写入速度时内存占用过高

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 09:57:01