实时更新二进制文件处理咨询:ADC多通道数据滤波与拆分
实时处理ADC动态二进制文件的方案分析与实现
Great question—your initial approach is totally feasible, but we can tweak it to be more efficient and reliable for real-time processing. Let’s break this down step by step:
你的初始方案是否可行?
你的思路完全可行,但固定time.sleep(1)的方式有优化空间:固定睡眠时间可能导致两种问题——要么新数据已经产生了很多却还在等待,要么数据还没写入就重复读取空内容。更好的方式是基于文件大小变化触发读取,而不是固定等待。
优化版初始方案实现
我们可以修改现有代码,实现增量式读取、处理,同时避免无效等待:
核心逻辑步骤
- 记录文件初始大小,作为下次读取的起始位置
- 循环监听文件大小,当检测到新增数据时,读取新增的字节块
- 对新增数据进行滤波、拆分写入
- 刷新实时图形
- 加入短时间的低延迟等待(比如0.1秒),避免CPU空转
代码修改示例(结合你的现有代码)
import os import time import array from PyQt5 import QtWidgets # 初始化参数(从你的UI控件获取) channel_number = int(self.lineEdit_3.text()) folder_base = str(self.lineEdit_2.text()) file_prefix = str(self.lineEdit.text()) coef_1 = float(self.lineEdit_9.text()) coef_2 = float(self.lineEdit_10.text()) num_spe = int(self.lineEdit_11.text()) # 选择目标文件 file_path, ext = QtWidgets.QFileDialog.getOpenFileName() if not file_path: exit() # 创建所有通道的文件夹(提前一次性创建,避免重复判断) base_output_path = os.path.join(os.getcwd(), "Data File", folder_base) for ch in range(1, channel_number + 1): ch_folder = os.path.join(base_output_path, f"{folder_base}{ch}") os.makedirs(ch_folder, exist_ok=True) # 用exist_ok避免重复创建报错 # 初始化文件读取位置和滤波状态(如果是IIR滤波,需要保留每个通道的历史值) current_pos = 0 # 示例:如果是IIR滤波,存储每个通道的上一次滤波结果 filter_states = [0.0 for _ in range(channel_number)] # 实时处理循环 try: while True: # 获取当前文件大小 stat_info = os.stat(file_path) file_size = stat_info.st_size # 计算新增数据的字节数 new_bytes = file_size - current_pos if new_bytes <= 0: time.sleep(0.1) # 无新数据,短暂等待 continue # 确保新增数据是2字节的整数倍(因为用的是array.array("h"),每个元素2字节) new_bytes = new_bytes - (new_bytes % 2) if new_bytes <= 0: time.sleep(0.1) continue # 读取新增的二进制数据 with open(file_path, 'rb') as fb: fb.seek(current_pos) bin_f = array.array("h") # 计算新增的元素个数:字节数 / 2 num_elements = new_bytes // 2 bin_f.fromfile(fb, num_elements) # 更新读取位置 current_pos += new_bytes # 拆分通道数据并滤波 # 假设ADC数据是按通道顺序循环写入的:ch1, ch2, ..., chN, ch1, ch2... for ch_idx in range(channel_number): # 提取当前通道的数据:从ch_idx开始,步长为channel_number ch_data = bin_f[ch_idx::channel_number] # 滤波处理(这里用你的滤波逻辑,示例用简单的IIR) filtered_data = [] for val in ch_data: filtered_val = coef_1 * val + coef_2 * filter_states[ch_idx] filtered_data.append(filtered_val) filter_states[ch_idx] = filtered_val # 写入对应通道的文件(采用追加模式,避免覆盖) ch_file_path = os.path.join(base_output_path, f"{folder_base}{ch_idx+1}", f"{file_prefix}.txt") with open(ch_file_path, 'a') as f: # 可以批量写入,比逐行写高效 f.write('\n'.join(map(str, filtered_data)) + '\n') # 刷新实时图形(这里调用你的绘图函数,比如更新Matplotlib画布) self.update_plot(filtered_data) # 替换成你的实际绘图逻辑 except KeyboardInterrupt: print("实时处理已停止")
更优方案:多线程+内存映射
如果追求更低延迟和更好的UI响应,推荐采用多线程+内存映射的方案:
核心优势
- 多线程:将文件读取/处理逻辑放在子线程,主线程专注于UI和实时绘图,避免UI卡顿
- 内存映射(
mmap):直接将文件映射到内存,无需每次打开/读取,大幅提升读取效率 - 批量写入:攒够一定量的数据再写入文件,减少磁盘IO次数,降低延迟
关键实现要点
- 子线程处理数据:用
QThread(因为你用的是PyQt)创建数据处理线程,避免阻塞主线程 - 内存映射读取:用
mmap.mmap打开文件,实时获取新增数据 - 线程间通信:用
pyqtSignal将处理后的数据传递给主线程更新图形 - 缓存写入:设置一个缓存阈值(比如1000条数据),达到阈值再写入文件
简化示例(多线程部分)
from PyQt5.QtCore import QThread, pyqtSignal import mmap class ADCProcessorThread(QThread): data_processed = pyqtSignal(list) # 传递处理后的数据给主线程 def __init__(self, file_path, channel_number, coefs): super().__init__() self.file_path = file_path self.channel_number = channel_number self.coef_1, self.coef_2 = coefs self.running = True self.filter_states = [0.0]*channel_number self.cache = [[] for _ in range(channel_number)] self.cache_threshold = 1000 # 缓存阈值 def run(self): with open(self.file_path, 'rb') as f: with mmap.mmap(f.fileno(), length=0, access=mmap.ACCESS_READ) as mm: current_pos = 0 while self.running: # 获取当前文件大小(内存映射的size会自动更新) file_size = mm.size() new_bytes = file_size - current_pos if new_bytes <= 0: time.sleep(0.05) continue new_bytes -= new_bytes % 2 if new_bytes <=0: continue # 从内存映射读取数据 data = mm[current_pos:current_pos+new_bytes] bin_f = array.array("h") bin_f.frombytes(data) current_pos += new_bytes # 处理每个通道 for ch_idx in range(self.channel_number): ch_data = bin_f[ch_idx::self.channel_number] # 滤波 filtered = [] for val in ch_data: filt_val = self.coef_1 * val + self.coef_2 * self.filter_states[ch_idx] filtered.append(filt_val) self.filter_states[ch_idx] = filt_val # 加入缓存 self.cache[ch_idx].extend(filtered) # 达到阈值则写入 ch_file = os.path.join(base_output_path, f"{folder_base}{ch_idx+1}", f"{file_prefix}.txt") with open(ch_file, 'a') as f: f.write('\n'.join(map(str, self.cache[ch_idx])) + '\n') self.cache[ch_idx] = [] # 发送数据给主线程绘图 self.data_processed.emit(filtered) def stop(self): self.running = False self.wait() # 在主线程中启动线程 processor = ADCProcessorThread(file_path, channel_number, (coef_1, coef_2)) processor.data_processed.connect(self.update_plot) processor.start() # 关闭程序时停止线程 # 在closeEvent中调用processor.stop()
关键注意事项
- 数据格式一致性:确保ADC写入的二进制文件是按固定通道顺序循环的(比如ch1, ch2,...ch32, ch1...),否则拆分逻辑会出错
- 滤波状态保留:如果使用需要历史数据的滤波算法(如IIR),必须保留每个通道的滤波状态,不能每次处理新块就重置
- 文件访问冲突:确保ADC程序以追加模式写入文件,你的程序以只读模式打开,避免读写冲突
- 延迟优化:测试不同的缓存阈值和等待时间,找到平衡CPU占用和延迟的最优值
内容的提问来源于stack exchange,提问作者Ziad Lucka
相关产品推荐
相关产品推荐

