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

从FTP服务器实时追踪CSV文件时无法读取新行的问题求助

实时追踪FTP服务器上CSV文件时无法读取新行的问题解决

问题现象

程序可读取CSV文件所有现有行,但后续仅能读取空行,打印输出示例:

b'+1.000000000E+01,+1.98000E+01,+3.15500E-04,-4.23700E-03,-3.20300E-03,-1.13750E-03,-3.26900E-03,+1.47450E-03,-3.39800E-03,-1.05500E-04,-4.24500E-03,0\r\n'
b'+1.100000000E+01,+1.98000E+01,-2.63000E-04,-4.24300E-03,-3.23350E-03,-4.79000E-04,-3.26000E-03,+1.43450E-03,-3.49800E-03,-6.47500E-04,-4.11750E-03,0\r\n'
b'+1.200000000E+01,+1.98000E+01,-9.49500E-04,-3.80900E-03,-3.22450E-03,+1.14500E-04,-3.26800E-03,+1.26850E-03,-3.63450E-03,-1.32050E-03,-3.59800E-03,0\r\n'
...
b''
b''
b''
b''

原实现代码

from ftplib import FTP
import os
import pandas as pd
import time
from io import StringIO
from io import BytesIO

ftp = FTP('192.168.10.100')
ftp.login()

def follow(thefile):
    global ftp
    with BytesIO() as flo:
        ftp.retrbinary('RETR ' + thefile, flo.write)
        flo.seek(0)
        while True:
            line = flo.readline()
            print(line)
            if not line or not line.endswith(b'\n'):
                time.sleep(1)
                continue
            yield line.decode('utf-8')

if __name__ == "__main__":
    global df
    # set the file we want to log the current line to
    log_file = "./current_line"

    # check if the last line processed has been saved
    if os.path.exists(log_file):
        with open(log_file, 'r') as ifile:
            # get the last line processed
            start_line = int(ifile.read())
    else:
        # set the last line processed to be the first data row (not the header).  If there is no header then set to 0
        start_line = 1

    # set the file we are reading
    myfile = '/sdcard/HIOKI/LR8450/DATA/24-04-12/tmpWvData240412_172748.CSV'

    # remove this line if you don't need the header
    # flo.seek(0)
    # header = pd.read_csv(flo, nrows=0, on_bad_lines='skip', engine="python", encoding='unicode_escape')

    # initialize the list to store the lines in
    lines = []

    # loop through each line in the file
    for nline, line in enumerate(follow(myfile)):
        # if we have already processed this file
        if nline < start_line:
            continue
        # append to the lines list
        lines.append(line)

        # check if the we have hit the number of lines we want to handle
        if len(lines) == 10:
            # read the csv from the lines we have processed
            df = pd.read_csv(StringIO(''.join(lines)), header=None, on_bad_lines='skip', engine="python", encoding='unicode_escape')
            # update the header.  Delete this row if there is no header
            df.columns = header.columns

            # do something with df
            print(df)

            # reset the lines list
            lines = []  

            # open the log file and note the line we have processed up to
            with open(log_file, 'w') as lfile:
                lfile.write(str(nline))  # only write the processed lines when we have actually done something

问题根源

原代码的follow函数仅在初始化时通过ftp.retrbinary将整个文件下载到本地BytesIO对象中,之后一直从这个静态内存对象读取数据,完全没有与FTP服务器交互获取新增内容,因此读完现有行后只能持续返回空行。

修复方案

利用FTP的REST命令支持字节范围请求的特性,跟踪已读取的字节数,每次循环时从上次读取的位置开始获取文件新增部分,实现真正的实时追踪:

修复后的代码

from ftplib import FTP
import os
import pandas as pd
import time
from io import StringIO
from io import BytesIO

ftp = FTP('192.168.10.100')
ftp.login()

def follow(ftp_path):
    global ftp
    # 记录已读取的字节数,初始为0
    bytes_read = 0
    while True:
        try:
            # 获取当前文件的大小
            file_size = ftp.size(ftp_path)
            if file_size > bytes_read:
                # 定位到上次读取的位置
                ftp.resp = '200 OK'  # 重置响应状态,避免REST命令报错
                ftp.sendcmd(f'REST {bytes_read}')
                # 读取新增的字节内容
                with BytesIO() as flo:
                    ftp.retrbinary('RETR ' + ftp_path, flo.write)
                    flo.seek(0)
                    new_content = flo.read()
                # 更新已读取的字节数
                bytes_read = file_size
                # 按行分割并返回有效行
                lines = new_content.splitlines(keepends=True)
                for line in lines:
                    if line.endswith(b'\n') or line.endswith(b'\r\n'):
                        yield line.decode('utf-8')
            else:
                # 没有新增内容,等待1秒后重试
                time.sleep(1)
        except Exception as e:
            print(f"读取异常: {e}")
            time.sleep(1)

if __name__ == "__main__":
    # 记录已处理的最后一行索引
    log_file = "./current_line"
    start_line = 1
    if os.path.exists(log_file):
        with open(log_file, 'r') as ifile:
            start_line = int(ifile.read().strip())
    
    myfile = '/sdcard/HIOKI/LR8450/DATA/24-04-12/tmpWvData240412_172748.CSV'
    
    # 先获取文件头(如果需要)
    header = None
    with BytesIO() as flo:
        ftp.retrbinary('RETR ' + myfile, flo.write)
        flo.seek(0)
        header = pd.read_csv(flo, nrows=0, on_bad_lines='skip', engine="python", encoding='unicode_escape')
    
    lines = []
    current_line_idx = 0
    
    for line in follow(myfile):
        current_line_idx += 1
        # 跳过已处理的行
        if current_line_idx < start_line:
            continue
        
        lines.append(line)
        
        # 每积累10行处理一次
        if len(lines) == 10:
            df = pd.read_csv(StringIO(''.join(lines)), header=None, on_bad_lines='skip', engine="python", encoding='unicode_escape')
            df.columns = header.columns
            
            # 这里添加你的数据处理逻辑
            print(df)
            
            lines = []
            # 更新已处理的最后一行
            with open(log_file, 'w') as lfile:
                lfile.write(str(current_line_idx))

关键修改点

  1. 字节偏移追踪进度:通过ftp.size()获取文件最新大小,对比已读取字节数判断是否有新增内容
  2. REST命令定位读取:每次读取前用REST命令告诉FTP服务器从指定字节位置开始传输,避免重复下载整个文件
  3. 独立获取文件头:初始化时单独下载一次文件获取表头,避免每次追踪重复处理
  4. 异常容错处理:添加异常捕获,避免网络波动导致程序崩溃

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 00:08:11