修改Python 3.12的myReader.py,确保单次停止指令即可触发终止
让myReader.py响应单次停止指令立即终止的修改方案
问题背景
当前myReader.py仅在部分场景下能检测到"Stop, damnit!"行并停止,临时方案是让myWriter.py重复写入该指令,但耗时极不稳定(曾测试15-30分钟)。问题根源是异步API调用耗时波动导致任务积压,程序需等待所有异步任务完成后才会检测停止指令。目标是实现myReader.py响应单次写入的停止指令即可立即终止。
原代码参考
myWriter.py
import time #Repeat 900 times to test output. Sleep for 1 second between each. for i in range(900): writeToFile("Stop, damnit!") time.sleep(1)
myReader.py
import os import platform import asyncio import aiofiles BATCH_SIZE = 10 def get_source_file_path(): if platform.system() == 'Windows': return 'C:\\path\\to\\sourceFile.txt' else: return '/path/to/sourceFile.txt' async def send_to_api(linesBuffer): success = runAPI(linesBuffer) return success async def read_source_file(): source_file_path = get_source_file_path() counter = 0 print("Reading source file...") print("source_file_path: ", source_file_path) #Detect the size of the file located at source_file_path and store it in the variable file_size. file_size = os.path.getsize(source_file_path) print("file_size: ", file_size) taskCountList = [] background_tasks = set() async with aiofiles.open(source_file_path, 'r') as source_file: await source_file.seek(0, os.SEEK_END) linesBuffer = [] while True: # Always make sure that file_size is the current size: line = await source_file.readline() new_file_size = os.path.getsize(source_file_path) if new_file_size < file_size: print("The file has been truncated.") print("old file_size: ", file_size) print("new_file_size: ", new_file_size) await source_file.seek(0, os.SEEK_SET) file_size = new_file_size # Allocate a new list instead of clearing the current one linesBuffer = [] counter = 0 continue line = await source_file.readline() if line: new_line = str(counter) + " line: " + line print(new_line) linesBuffer.append(new_line) print("len(linesBuffer): ", len(linesBuffer)) if len(linesBuffer) == BATCH_SIZE: print("sending to api...") task = asyncio.create_task(send_to_api(linesBuffer)) background_tasks.add(task) task.add_done_callback(background_tasks.discard) pendingTasks = len(background_tasks) taskCountList.append(pendingTasks) print("") print("pendingTasks: ", pendingTasks) print("") # Do not clear the buffer; allocate a new one: linesBuffer = [] counter += 1 print("counter: ", counter) #detect whether or not the present line is the last line in the file. # If it is the last line in the file, then write whatever batch # we have even if it is not complete. if "Stop, damnit!" in line: #Print the next line 30 times to simulate a large file. for i in range(30): print("LAST LINE IN FILE FOUND.") #sleep for 1 second to simulate a large file. await asyncio.sleep(1) #Omitting other stuff for brevity. break else: await asyncio.sleep(0.1) async def main(): await read_source_file() if __name__ == '__main__': asyncio.run(main())
具体修改点
1. 修复重复读取行的BUG
原代码连续调用两次readline(),导致每轮循环跳过一行,可能直接漏掉停止指令。
原代码片段:
line = await source_file.readline() new_file_size = os.path.getsize(source_file_path) # ... 中间截断检查代码 ... line = await source_file.readline()
修改后:
new_file_size = os.path.getsize(source_file_path) if new_file_size < file_size: print("The file has been truncated.") print("old file_size: ", file_size) print("new_file_size: ", new_file_size) await source_file.seek(0, os.SEEK_SET) file_size = new_file_size linesBuffer = [] counter = 0 continue # 仅保留一次readline调用 line = await source_file.readline()
2. 优先检测停止指令,不等待批次满
原代码将停止检测放在批次逻辑之后,若当前buffer未达到BATCH_SIZE,会继续等待新行,无法及时响应停止指令。调整顺序为读取行后立即检测停止指令,优先处理终止逻辑。
原代码片段:
linesBuffer.append(new_line) print("len(linesBuffer): ", len(linesBuffer)) if len(linesBuffer) == BATCH_SIZE: # ... 批次处理代码 ... if "Stop, damnit!" in line: # ... 停止处理代码 ...
修改后:
# 先检测停止指令,优先处理终止逻辑 if "Stop, damnit!" in line: print("LAST LINE IN FILE FOUND.") # 可选:处理当前buffer中剩余的未提交行 if linesBuffer: task = asyncio.create_task(send_to_api(linesBuffer)) background_tasks.add(task) # 可选:等待所有已提交的异步API任务完成(若不需要等待可直接跳过) if background_tasks: await asyncio.gather(*background_tasks) break # 再处理批次收集逻辑 linesBuffer.append(new_line) print("len(linesBuffer): ", len(linesBuffer)) if len(linesBuffer) == BATCH_SIZE: # ... 原批次处理代码保持不变 ...
3. 移除冗余的模拟延迟代码
原代码中检测到停止指令后有30次循环打印和1秒睡眠,这会强制延迟程序终止,直接移除这部分无意义的代码。
原代码片段:
if "Stop, damnit!" in line: #Print the next line 30 times to simulate a large file. for i in range(30): print("LAST LINE IN FILE FOUND.") #sleep for 1 second to simulate a large file. await asyncio.sleep(1) break
修改后:
if "Stop, damnit!" in line: print("LAST LINE IN FILE FOUND.") # ... 剩余任务处理逻辑 ... break
修改后的完整myReader.py代码
import os import platform import asyncio import aiofiles BATCH_SIZE = 10 def get_source_file_path(): if platform.system() == 'Windows': return 'C:\\path\\to\\sourceFile.txt' else: return '/path/to/sourceFile.txt' async def send_to_api(linesBuffer): success = runAPI(linesBuffer) return success async def read_source_file(): source_file_path = get_source_file_path() counter = 0 print("Reading source file...") print("source_file_path: ", source_file_path) file_size = os.path.getsize(source_file_path) print("file_size: ", file_size) taskCountList = [] background_tasks = set() async with aiofiles.open(source_file_path, 'r') as source_file: await source_file.seek(0, os.SEEK_END) linesBuffer = [] while True: # 先检查文件大小变化 new_file_size = os.path.getsize(source_file_path) if new_file_size < file_size: print("The file has been truncated.") print("old file_size: ", file_size) print("new_file_size: ", new_file_size) await source_file.seek(0, os.SEEK_SET) file_size = new_file_size linesBuffer = [] counter = 0 continue # 读取单行 line = await source_file.readline() if line: new_line = str(counter) + " line: " + line print(new_line) # 优先检测停止指令 if "Stop, damnit!" in line: print("LAST LINE IN FILE FOUND.") # 处理剩余未提交的行 if linesBuffer: task = asyncio.create_task(send_to_api(linesBuffer)) background_tasks.add(task) # 等待所有异步任务完成(可选,根据业务需求调整) if background_tasks: await asyncio.gather(*background_tasks) break # 批次收集逻辑 linesBuffer.append(new_line) print("len(linesBuffer): ", len(linesBuffer)) if len(linesBuffer) == BATCH_SIZE: print("sending to api...") task = asyncio.create_task(send_to_api(linesBuffer)) background_tasks.add(task) task.add_done_callback(background_tasks.discard) pendingTasks = len(background_tasks) taskCountList.append(pendingTasks) print("") print("pendingTasks: ", pendingTasks) print("") linesBuffer = [] counter += 1 print("counter: ", counter) else: await asyncio.sleep(0.1) async def main(): await read_source_file() if __name__ == '__main__': asyncio.run(main())
内容的提问来源于stack exchange,提问作者CodeMed
相关产品推荐
相关产品推荐

