从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))
关键修改点
- 字节偏移追踪进度:通过
ftp.size()获取文件最新大小,对比已读取字节数判断是否有新增内容 - REST命令定位读取:每次读取前用
REST命令告诉FTP服务器从指定字节位置开始传输,避免重复下载整个文件 - 独立获取文件头:初始化时单独下载一次文件获取表头,避免每次追踪重复处理
- 异常容错处理:添加异常捕获,避免网络波动导致程序崩溃
内容的提问来源于stack exchange,提问作者ewm
相关产品推荐
相关产品推荐

