如何修改Python 3.12代码使write_to_file任务并行执行?
让
write_to_file调用并行执行的具体修改方案 核心问题分析
当前代码用await write_to_file(...)会串行等待每个写入任务完成,导致效率低下;直接去掉await仅返回协程对象,不会被事件循环调度执行,因此需要通过异步任务机制启动协程。
具体修改步骤
维护任务跟踪列表
在read_source_file函数内新增列表,用于保存所有未完成的写入任务,避免任务被垃圾回收提前终止:pending_tasks = []用
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()]调整
write_to_file函数
移除函数内的linesBuffer.clear()(因为操作的是副本,原buffer在read_source_file中已经处理):# 删除这一行 # linesBuffer.clear()确保所有任务完成
如果程序有正常退出逻辑,在退出前等待所有未完成任务执行完毕: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
相关产品推荐
相关产品推荐

