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

按需停止多线程(如按键输入):串口设备数据采集问题

多串口设备数据采集的线程控制与集成方案

问题背景

我有两个需同时运行的串口设备:一个用于计数水样中的颗粒,另一个用于测量水样流速。需求是持续采集数据直至水样耗尽,需支持按需停止(按下q+回车终止采集)。目前用线程实现双设备并行,但存在以下问题:

  • 线程无法正确响应停止指令,启动后未等待用户输入就直接执行完毕
  • 后续需要集成流速采集函数get_FlowData,同时要对采集数据做格式化处理并存储到SQL数据库

现有代码问题分析

  1. 主线程逻辑错误:启动线程后立刻调用stop_abakus_event.set(),直接终止采集线程,完全跳过用户输入等待环节
  2. 串口对象管理混乱:ser变量未在主线程初始化就传递给线程,且线程内部又重新创建Serial实例,导致资源冗余
  3. 停止逻辑不严谨:用户输入任意内容都会触发停止,未判断是否为指定的q指令
  4. 数据共享风险:使用全局变量output传递采集数据,多线程环境下可能出现数据竞争
  5. 依赖缺失:数据处理和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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 11:57:09