按需停止多线程(如按键输入):串口设备数据采集问题
多串口设备数据采集的线程控制与集成方案
问题背景
我有两个需同时运行的串口设备:一个用于计数水样中的颗粒,另一个用于测量水样流速。需求是持续采集数据直至水样耗尽,需支持按需停止(按下q+回车终止采集)。目前用线程实现双设备并行,但存在以下问题:
- 线程无法正确响应停止指令,启动后未等待用户输入就直接执行完毕
- 后续需要集成流速采集函数
get_FlowData,同时要对采集数据做格式化处理并存储到SQL数据库
现有代码问题分析
- 主线程逻辑错误:启动线程后立刻调用
stop_abakus_event.set(),直接终止采集线程,完全跳过用户输入等待环节 - 串口对象管理混乱:
ser变量未在主线程初始化就传递给线程,且线程内部又重新创建Serial实例,导致资源冗余 - 停止逻辑不严谨:用户输入任意内容都会触发停止,未判断是否为指定的
q指令 - 数据共享风险:使用全局变量
output传递采集数据,多线程环境下可能出现数据竞争 - 依赖缺失:数据处理和SQL存储所需的模块未导入
修正后的完整代码
import threading import serial import re import pandas as pd from datetime import datetime from sqlalchemy import create_engine from queue import Queue # 线程安全队列,用于传递采集到的数据 data_queue = Queue() def get_AbakusData(stop_event): # 初始化串口 ser = serial.Serial('COM5', 38400, timeout=1) try: # 设备初始化指令 ser.write(b'C0\r\n') ser.write(b'C1\r\n') ser.write(b'C5\r\n') ser.write(b'C12\r\n') ser.readline() # 读取初始化响应 regex = r'\b0\d{7}\b' run_id = datetime.now().strftime("%Y%m%d_%H%M%S") # 生成唯一运行ID while not stop_event.is_set(): ser.write(b'C12\r\n') output = ser.readline().decode('utf-8').strip() if not output: continue # 跳过空响应 # 数据处理逻辑 matches = re.findall(regex, output) output_str = ' '.join([match.replace(' ', '')[:8] for match in matches]) if not output_str: continue # 转换为DataFrame series = pd.Series(output_str) split_data = series.str.split() df = pd.DataFrame({ 'bins': split_data.str[::2], 'counts': split_data.str[1::2] }) df = pd.concat([df[col].explode().reset_index(drop=True) for col in df], axis=1) # 添加时间戳和运行ID now = datetime.now() df['time'] = now.strftime("%H:%M:%S") df['date'] = now.strftime("%Y-%m-%d") df['run_id'] = run_id # 转换数据类型 df['bins'] = df['bins'].astype(int) df['counts'] = df['counts'].astype(int) # 将数据放入队列,供存储线程处理 data_queue.put(('abakus', df)) finally: ser.close() # 确保串口关闭 def get_FlowData(stop_event): # 流速设备采集逻辑,根据实际参数调整 ser = serial.Serial('COM6', 9600, timeout=1) try: # 设备初始化(根据实际指令修改) ser.write(b'INIT\r\n') ser.readline() run_id = datetime.now().strftime("%Y%m%d_%H%M%S") while not stop_event.is_set(): ser.write(b'QUERY_FLOW\r\n') # 根据实际查询指令修改 output = ser.readline().decode('utf-8').strip() if not output: continue # 流速数据处理逻辑(根据实际数据格式调整) try: flow_value = float(output) now = datetime.now() df = pd.DataFrame({ 'flow': [flow_value], 'time': [now.strftime("%H:%M:%S")], 'date': [now.strftime("%Y-%m-%d")], 'run_id': [run_id] }) data_queue.put(('flow', df)) except ValueError: print(f"无效的流速数据: {output}") finally: ser.close() def data_storage_worker(db_conn_str): # 数据存储线程,从队列中取数据并存入SQL engine = create_engine(db_conn_str) conn = engine.connect() try: while True: data_type, df = data_queue.get() if data_type == 'stop': # 停止信号 break # 写入数据库 df.to_sql(f'{data_type}_data', conn, if_exists='append', index=False) conn.commit() data_queue.task_done() finally: conn.close() def get_UserInput(stop_event): print("按下 'q' + 回车键停止所有采集线程...") while True: user_input = input().strip().lower() if user_input == 'q': stop_event.set() break print("停止指令已发送,正在等待线程结束...") if __name__ == "__main__": # 数据库连接字符串(根据实际数据库调整) DB_CONN_STR = 'mysql+pymysql://username:password@localhost/db_name' # 创建停止事件 stop_event = threading.Event() # 启动线程 t_abakus = threading.Thread(target=get_AbakusData, args=(stop_event,)) t_flow = threading.Thread(target=get_FlowData, args=(stop_event,)) t_storage = threading.Thread(target=data_storage_worker, args=(DB_CONN_STR,)) t_user_input = threading.Thread(target=get_UserInput, args=(stop_event,)) t_abakus.start() t_flow.start() t_storage.start() t_user_input.start() # 等待用户输入线程结束 t_user_input.join() # 等待采集线程结束 t_abakus.join() t_flow.join() # 发送停止信号给存储线程 data_queue.put(('stop', None)) t_storage.join() print("所有采集线程已停止,数据存储完成。")
关键优化说明
- 线程控制修复:主线程不再提前设置停止事件,等待用户输入线程触发停止信号后,再终止采集和存储线程
- 串口资源管理:每个采集线程内部独立初始化和关闭串口,避免跨线程资源冲突
- 线程安全数据传递:使用
Queue实现采集线程与存储线程的解耦,避免全局变量的数据竞争问题 - 严谨停止逻辑:仅当用户输入
q时才触发停止指令 - 模块化设计:将采集、数据处理、存储分离,便于后续功能扩展和维护
- 异常防护:添加
try/finally确保串口始终关闭,处理无效数据的异常情况
内容的提问来源于stack exchange,提问作者abby hudak
相关产品推荐
相关产品推荐

