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

如何修改Python 3.12代码使write_to_file任务并行执行?

让write_to_file调用并行执行的具体修改方案

核心问题分析

当前代码用await write_to_file(...)会串行等待每个写入任务完成,导致效率低下;直接去掉await仅返回协程对象,不会被事件循环调度执行,因此需要通过异步任务机制启动协程。

具体修改步骤

  1. 维护任务跟踪列表
    在read_source_file函数内新增列表,用于保存所有未完成的写入任务,避免任务被垃圾回收提前终止:

    pending_tasks = []
    
  2. 用asyncio.create_task启动并行任务
    将原有的await write_to_file(linesBuffer)替换为创建异步任务的代码,必须传入linesBuffer的副本(否则原buffer被清空会导致写入空内容):

    # 替换原await语句
    task = asyncio.create_task(write_to_file(linesBuffer.copy()))
    pending_tasks.append(task)
    # 定期清理已完成的任务,避免内存占用
    pending_tasks = [t for t in pending_tasks if not t.done()]
    
  3. 调整write_to_file函数
    移除函数内的linesBuffer.clear()(因为操作的是副本,原buffer在read_source_file中已经处理):

    # 删除这一行
    # linesBuffer.clear()
    
  4. 确保所有任务完成
    如果程序有正常退出逻辑,在退出前等待所有未完成任务执行完毕:

    await asyncio.gather(*pending_tasks)
    

修改后的完整代码

import os
import platform
import asyncio

numLines = 10

def get_source_file_path():
    if platform.system() == 'Windows':
        return 'C:\\path\\to\\sourceFile.txt'
    else:
        return '/path/to/sourceFile.txt'

async def write_to_file(linesBuffer):
    print("inside Writing to file...")
    with open('newFile.txt', 'a') as new_destination_file:
        for line in linesBuffer:
            new_destination_file.write(line)
    # 获取newFile.txt所在目录并打印
    directory_name = os.path.dirname(os.path.abspath('newFile.txt'))
    print("directory_name: ", directory_name)
    # 每1秒打印一次,持续2秒
    for i in range(2):
        print("HI HO, HI HO.  IT'S OFF TO WORK WE GO...")
        await asyncio.sleep(1)
    print("inside done Writing to file...")

async def read_source_file():
    source_file_path = get_source_file_path()
    linesBuffer = []
    counter = 0
    pending_tasks = []  # 新增:跟踪未完成的写入任务
    print("Reading source file...")
    print("source_file_path: ", source_file_path)
    # 获取源文件大小
    file_size = os.path.getsize(source_file_path)
    print("file_size: ", file_size)
    with open(source_file_path, 'r') as source_file:
        source_file.seek(0, os.SEEK_END)
        while True:
            line = source_file.readline()
            new_file_size = os.path.getsize(source_file_path)
            if new_file_size < file_size:
                print("The file has been truncated.")
                source_file.seek(0, os.SEEK_SET)
                file_size = new_file_size
                linesBuffer.clear()
                counter = 0
                print("new_file_size: ", new_file_size)
            if len(line) > 0:
              new_line = str(counter) + " line: " + line
              print(new_line)
              linesBuffer.append(new_line)
              print("len(linesBuffer): ", len(linesBuffer))
              if len(linesBuffer) >= numLines:
                print("Writing to file...")
                # 替换原await语句,创建并行任务
                task = asyncio.create_task(write_to_file(linesBuffer.copy()))
                pending_tasks.append(task)
                # 清理已完成的任务
                pending_tasks = [t for t in pending_tasks if not t.done()]
                print("启动写入任务,继续读取...")
                linesBuffer.clear()
              counter += 1
              print("counter: ", counter)
            if not line:
                await asyncio.sleep(0.1)
                continue
            # 检测是否为文件最后一行,是的话写入
            if source_file.tell() == file_size:
                print("LAST LINE IN FILE FOUND.  Writing to file...")
                # 同样用create_task启动并行任务
                task = asyncio.create_task(write_to_file(linesBuffer.copy()))
                pending_tasks.append(task)
                pending_tasks = [t for t in pending_tasks if not t.done()]
                print("启动写入任务,继续读取...")
                linesBuffer.clear()
                counter = 0
        # 若循环有退出条件,此处添加等待所有任务完成的代码
        # await asyncio.gather(*pending_tasks)

async def main():
    await read_source_file()

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

关键修改说明

  • asyncio.create_task():将协程包装成异步任务,提交给事件循环调度,无需等待任务完成即可继续执行主逻辑,实现并行效果。
  • 传入buffer副本:原代码中linesBuffer会被立即清空,直接传原对象会导致写入任务读取空内容,因此必须用copy()创建副本。
  • 任务列表维护:保存任务对象防止被垃圾回收,定期清理已完成任务避免内存占用。
  • 等待任务完成:确保程序退出前所有写入操作都执行完毕,不会丢失数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 07:54:53