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

修改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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 05:35:00